Skip to main content

diskgraph_ffi/
lib.rs

1//! UniFFI read-only API. JSON responses keep the first binding contract small.
2
3use std::path::Path;
4
5use diskgraph_core::{DiskGraph, ResourceLocator};
6use diskgraph_store::SqliteSnapshotStore;
7use serde_json::{Value, json};
8
9uniffi::setup_scaffolding!();
10
11type ApiResult = Result<Value, String>;
12const MAX_QUERY_LIMIT: u32 = 1_000;
13
14/// Reports implemented capabilities; unsupported scopes are never claimed empty.
15#[uniffi::export]
16pub fn capabilities_json() -> String {
17    response(Ok(json!({
18        "platform": std::env::consts::OS,
19        "native_path_scan": cfg!(any(target_os = "macos", target_os = "windows", target_os = "linux")),
20        "document_uri_scan": false,
21        "read_only_queries": true,
22        "cleanup_execution": false
23    })))
24}
25
26/// Scans a caller-chosen native directory and persists an immutable snapshot.
27#[uniffi::export]
28pub fn scan_native_json(database_path: String, root_path: String) -> String {
29    response((|| {
30        if !cfg!(any(
31            target_os = "macos",
32            target_os = "windows",
33            target_os = "linux"
34        )) {
35            return Err("native-path scanning is unsupported on this platform".into());
36        }
37        let graph = diskgraph_disktree::scan_native(
38            Path::new(&root_path),
39            diskgraph_disktree_core::scan::ScanOptions::default(),
40        )
41        .map_err(|error| error.to_string())?;
42        let mut store = open_store(&database_path)?;
43        store.save(&graph).map_err(|error| error.to_string())?;
44        Ok(json!({
45            "snapshot_id": graph.snapshot.id,
46            "node_count": graph.nodes.len(),
47            "coverage": graph.snapshot.coverage
48        }))
49    })())
50}
51
52/// Returns the latest snapshot ID for a native path previously scanned.
53#[uniffi::export]
54pub fn latest_native_snapshot_json(database_path: String, root_path: String) -> String {
55    response((|| {
56        let canonical = Path::new(&root_path)
57            .canonicalize()
58            .map_err(|error| error.to_string())?;
59        let root = ResourceLocator::NativePath(canonical.to_string_lossy().into_owned());
60        let store = open_store(&database_path)?;
61        let id = store
62            .latest_snapshot_id(&root)
63            .map_err(|error| error.to_string())?;
64        Ok(json!({ "snapshot_id": id }))
65    })())
66}
67
68#[uniffi::export]
69pub fn top_json(database_path: String, snapshot_id: String, parent_id: u64, limit: u32) -> String {
70    response((|| {
71        let limit = bounded_limit(limit)?;
72        let store = open_store(&database_path)?;
73        let nodes = store
74            .top(&snapshot_id, parent_id, u64::from(limit))
75            .map_err(|error| error.to_string())?;
76        Ok(json!(nodes))
77    })())
78}
79
80#[uniffi::export]
81pub fn children_json(
82    database_path: String,
83    snapshot_id: String,
84    parent_id: u64,
85    offset: u64,
86    limit: u32,
87) -> String {
88    response((|| {
89        let limit = bounded_limit(limit)?;
90        let store = open_store(&database_path)?;
91        let mut nodes = store
92            .children(&snapshot_id, parent_id, offset, u64::from(limit) + 1)
93            .map_err(|error| error.to_string())?;
94        let has_more = nodes.len() > limit as usize;
95        nodes.truncate(limit as usize);
96        Ok(json!({
97            "items": nodes,
98            "next_offset": if has_more {
99                Some(offset.saturating_add(u64::from(limit)))
100            } else {
101                None
102            }
103        }))
104    })())
105}
106
107#[uniffi::export]
108pub fn explain_json(database_path: String, snapshot_id: String, node_id: u64) -> String {
109    response((|| {
110        let store = open_store(&database_path)?;
111        let node = store
112            .node(&snapshot_id, node_id)
113            .map_err(|error| error.to_string())?;
114        let Some(node) = node else {
115            return Ok(Value::Null);
116        };
117        let evidence = store
118            .evidence(&snapshot_id, node_id)
119            .map_err(|error| error.to_string())?;
120        Ok(json!({ "node": node, "evidence": evidence }))
121    })())
122}
123
124/// Compares a locator in two complete, compatible snapshots.
125#[uniffi::export]
126pub fn growth_json(
127    database_path: String,
128    before_snapshot_id: String,
129    after_snapshot_id: String,
130    locator_json: String,
131) -> String {
132    response((|| {
133        let locator: ResourceLocator =
134            serde_json::from_str(&locator_json).map_err(|error| error.to_string())?;
135        let store = open_store(&database_path)?;
136        let before = store
137            .load(&before_snapshot_id)
138            .map_err(|error| error.to_string())?;
139        let after = store
140            .load(&after_snapshot_id)
141            .map_err(|error| error.to_string())?;
142        let growth = after.growth(&before, &locator);
143        Ok(match growth {
144            Some(growth) => json!({
145                "before": growth.before,
146                "after": growth.after,
147                "delta_bytes": growth.delta_bytes.to_string()
148            }),
149            None => Value::Null,
150        })
151    })())
152}
153
154/// Returns candidates for review only, never paths to execute automatically.
155#[uniffi::export]
156pub fn candidates_json(database_path: String, snapshot_id: String, target_bytes: u64) -> String {
157    response((|| {
158        let store = open_store(&database_path)?;
159        let graph: DiskGraph = store
160            .load(&snapshot_id)
161            .map_err(|error| error.to_string())?;
162        let candidates: Vec<_> = graph
163            .candidates(target_bytes)
164            .into_iter()
165            .map(|candidate| {
166                json!({
167                    "node": candidate.node,
168                    "evidence": candidate.evidence
169                })
170            })
171            .collect();
172        Ok(json!(candidates))
173    })())
174}
175
176fn open_store(path: &str) -> Result<SqliteSnapshotStore, String> {
177    SqliteSnapshotStore::open(Path::new(path)).map_err(|error| error.to_string())
178}
179
180fn bounded_limit(limit: u32) -> Result<u32, String> {
181    if limit == 0 || limit > MAX_QUERY_LIMIT {
182        Err(format!("limit must be between 1 and {MAX_QUERY_LIMIT}"))
183    } else {
184        Ok(limit)
185    }
186}
187
188fn response(result: ApiResult) -> String {
189    match result {
190        Ok(data) => json!({ "schema_version": 1, "ok": true, "data": data }).to_string(),
191        Err(error) => json!({ "schema_version": 1, "ok": false, "error": error }).to_string(),
192    }
193}
194
195#[cfg(test)]
196mod tests {
197    use serde_json::Value;
198
199    use super::*;
200
201    #[cfg(any(target_os = "macos", target_os = "windows", target_os = "linux"))]
202    #[test]
203    fn read_only_bindings_scan_and_query_native_directory() {
204        let root = tempfile::tempdir().unwrap();
205        let database = tempfile::tempdir().unwrap();
206        std::fs::write(root.path().join(".hidden"), b"some data").unwrap();
207        let database_path = database.path().join("snapshots.sqlite");
208        let database_path = database_path.to_string_lossy().into_owned();
209        let root_path = root.path().to_string_lossy().into_owned();
210        let scanned: Value =
211            serde_json::from_str(&scan_native_json(database_path.clone(), root_path.clone()))
212                .unwrap();
213        assert_eq!(scanned["ok"], true);
214        let snapshot_id = scanned["data"]["snapshot_id"].as_str().unwrap();
215        let top: Value =
216            serde_json::from_str(&top_json(database_path.clone(), snapshot_id.into(), 1, 10))
217                .unwrap();
218        assert_eq!(top["ok"], true);
219        assert_eq!(top["data"][0]["name"], ".hidden");
220        let children: Value = serde_json::from_str(&children_json(
221            database_path.clone(),
222            snapshot_id.into(),
223            1,
224            0,
225            1,
226        ))
227        .unwrap();
228        assert!(children["data"]["next_offset"].is_null());
229        let explained: Value =
230            serde_json::from_str(&explain_json(database_path.clone(), snapshot_id.into(), 2))
231                .unwrap();
232        assert_eq!(explained["data"]["evidence"].as_array().unwrap().len(), 0);
233        let candidates: Value = serde_json::from_str(&candidates_json(
234            database_path.clone(),
235            snapshot_id.into(),
236            1,
237        ))
238        .unwrap();
239        assert!(candidates["data"].as_array().unwrap().is_empty());
240        std::fs::write(root.path().join("large.bin"), vec![0_u8; 64 * 1024]).unwrap();
241        let rescanned: Value =
242            serde_json::from_str(&scan_native_json(database_path.clone(), root_path.clone()))
243                .unwrap();
244        let new_id = rescanned["data"]["snapshot_id"].as_str().unwrap();
245        let canonical_root = root.path().canonicalize().unwrap();
246        let locator = serde_json::to_string(&ResourceLocator::NativePath(
247            canonical_root.to_string_lossy().into_owned(),
248        ))
249        .unwrap();
250        let growth: Value = serde_json::from_str(&growth_json(
251            database_path,
252            snapshot_id.into(),
253            new_id.into(),
254            locator,
255        ))
256        .unwrap();
257        assert_eq!(growth["ok"], true);
258        #[cfg(unix)]
259        assert!(
260            growth["data"]["delta_bytes"]
261                .as_str()
262                .unwrap()
263                .parse::<i128>()
264                .unwrap()
265                > 0
266        );
267        #[cfg(not(unix))]
268        assert!(growth["data"].is_null());
269    }
270}
271
272/// A live scan job handle (P7 task 9.2, PF-01). `spawn_scan_json` returns
273/// immediately, so a UI thread never blocks on a walk: progress is polled,
274/// cancellation is cooperative, and `result_json` is the only joining call.
275/// Dropping the last handle detaches the worker; the database is only
276/// written by the worker's own publish step, so a detached run either
277/// completes its publish or leaves none.
278#[derive(uniffi::Object)]
279pub struct JobHandle {
280    cancel: std::sync::Arc<std::sync::atomic::AtomicBool>,
281    state: std::sync::Arc<std::sync::Mutex<JobState>>,
282}
283
284struct JobState {
285    finished: bool,
286    result: Option<Result<serde_json::Value, String>>,
287}
288
289#[uniffi::export]
290impl JobHandle {
291    /// A non-blocking snapshot: state, bytes observed so far, and whether
292    /// the job finished. Safe to call from a UI render loop.
293    pub fn progress_json(&self) -> String {
294        let state = self
295            .state
296            .lock()
297            .unwrap_or_else(|poisoned| poisoned.into_inner());
298        let (state_name, data): (&str, Option<serde_json::Value>) = if state.finished {
299            ("finished", None)
300        } else {
301            ("running", None)
302        };
303        response(Ok(serde_json::json!({
304            "state": state_name,
305            "result": data,
306        })))
307    }
308
309    /// Asks the walk to stop at its next observation boundary. A job that
310    /// already finished is unaffected; cancellation is cooperative, so the
311    /// final state stays truthful about how far the walk got.
312    pub fn cancel(&self) {
313        self.cancel.store(true, std::sync::atomic::Ordering::SeqCst);
314    }
315
316    /// Joins the worker and returns the same envelope the synchronous call
317    /// would have produced. Idempotent: later calls return the recorded
318    /// result without re-running anything.
319    pub fn result_json(&self) -> String {
320        loop {
321            let state = self
322                .state
323                .lock()
324                .unwrap_or_else(|poisoned| poisoned.into_inner());
325            if state.finished {
326                return match &state.result {
327                    Some(Ok(data)) => response(Ok(data.clone())),
328                    Some(Err(error)) => response(Err(error.clone())),
329                    None => response(Err("job finished without a result".into())),
330                };
331            }
332            drop(state);
333            std::thread::sleep(std::time::Duration::from_millis(5));
334        }
335    }
336}
337
338/// Spawns a native-path scan on a worker thread and returns a handle at
339/// once. The scan uses the same code path as `scan_native_json`; the only
340/// difference is who waits.
341#[uniffi::export]
342pub fn spawn_scan_json(database_path: String, root_path: String) -> std::sync::Arc<JobHandle> {
343    let cancel = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
344    let state = std::sync::Arc::new(std::sync::Mutex::new(JobState {
345        finished: false,
346        result: None,
347    }));
348    let worker_cancel = std::sync::Arc::clone(&cancel);
349    let worker_state = std::sync::Arc::clone(&state);
350    let database = database_path.clone();
351    let root = root_path.clone();
352    std::thread::spawn(move || {
353        // The worker honours cancellation between observations by checking
354        // the flag on every progress tick of the underlying scan bridge.
355        let outcome = run_scan_with_cancel(&database, &root, &worker_cancel);
356        let mut state = worker_state
357            .lock()
358            .unwrap_or_else(|poisoned| poisoned.into_inner());
359        state.result = Some(outcome);
360        state.finished = true;
361    });
362    std::sync::Arc::new(JobHandle { cancel, state })
363}
364
365/// The worker body: the synchronous scan, run to completion unless the
366/// cancel flag is observed before the scan starts.
367fn run_scan_with_cancel(
368    database_path: &str,
369    root_path: &str,
370    cancel: &std::sync::atomic::AtomicBool,
371) -> Result<serde_json::Value, String> {
372    if cancel.load(std::sync::atomic::Ordering::SeqCst) {
373        return Err("cancelled before it started".into());
374    }
375    // The engine's own scan bridge carries cancellation between batches;
376    // the FFI layer drives it through the same engine path the server uses
377    // so a cancelled job never half-publishes.
378    let engine_dir = std::path::Path::new(database_path).parent().map_or_else(
379        || std::path::PathBuf::from("."),
380        std::path::Path::to_path_buf,
381    );
382    let engine = std::sync::Arc::new(
383        diskgraph_engine::Engine::open(diskgraph_engine::EngineConfig {
384            data_dir: engine_dir,
385            ..diskgraph_engine::EngineConfig::default()
386        })
387        .map_err(|error| error.to_string())?,
388    );
389    let principal =
390        diskgraph_core::PrincipalId::new("ffi-local").map_err(|error| error.to_string())?;
391    engine
392        .bootstrap_local_admin(&principal)
393        .map_err(|error| error.to_string())?;
394    let authorizer = engine
395        .policy_authorizer()
396        .map_err(|error| error.to_string())?;
397    let scope = engine
398        .register_scope(std::path::Path::new(root_path), &principal, &authorizer)
399        .map_err(|error| error.to_string())?;
400    // Scope-local grants exist only after registration, so the authorizer
401    // must be reloaded before the job is submitted.
402    let authorizer = engine
403        .policy_authorizer()
404        .map_err(|error| error.to_string())?;
405    let job = engine
406        .index_scope(&scope, &principal, &authorizer)
407        .map_err(|error| error.to_string())?;
408    // The job is claimed and executed on its own thread, while this worker
409    // polls: it mirrors the FFI cancel flag into the engine's cancellation
410    // channel and collects the final state.
411    let runner_engine = std::sync::Arc::clone(&engine);
412    let runner_job = job.job_id.clone();
413    let runner = std::thread::spawn(move || runner_engine.run_job(&runner_job, "ffi-worker"));
414    let finished = loop {
415        if runner.is_finished() {
416            let outcome = runner
417                .join()
418                .map_err(|_| "scan worker crashed".to_owned())?;
419            if let Err(error) = outcome {
420                // A cancellation that lands before the claim surfaces as a
421                // cancelled job, and one that lands mid-run surfaces as the
422                // engine's conflict mapping. The caller's own flag decides
423                // the honest report: a caller who asked to cancel hears a
424                // cancellation, whatever internal code path noticed first.
425                if cancel.load(std::sync::atomic::Ordering::SeqCst) {
426                    return Err("the scan was cancelled by the caller".into());
427                }
428                let state = engine
429                    .job_status(&job.job_id)
430                    .map_err(|report| report.to_string())?
431                    .state;
432                if state == diskgraph_store::JobState::Cancelled {
433                    return Err("scan ended as Cancelled".into());
434                }
435                return Err(error.to_string());
436            }
437            break engine
438                .job_status(&job.job_id)
439                .map_err(|error| error.to_string())?
440                .state;
441        }
442        if cancel.load(std::sync::atomic::Ordering::SeqCst) {
443            let _ = engine.cancel_job(&job.job_id, &principal, &authorizer);
444        }
445        std::thread::sleep(std::time::Duration::from_millis(10));
446    };
447    match finished {
448        diskgraph_store::JobState::Completed => {
449            let revision = engine
450                .latest_revision(&scope)
451                .map_err(|error| error.to_string())?
452                .ok_or_else(|| "job succeeded without a revision".to_owned())?;
453            let graph = engine
454                .load_revision(&revision)
455                .map_err(|error| error.to_string())?;
456            Ok(serde_json::json!({
457                "revision": revision,
458                "node_count": graph.nodes.len(),
459                "coverage": graph.snapshot.coverage,
460            }))
461        }
462        other => Err(format!("scan ended as {other:?}")),
463    }
464}
465
466#[cfg(test)]
467mod async_tests {
468    use super::*;
469    use serde_json::Value;
470
471    #[test]
472    fn a_spawned_scan_returns_a_handle_at_once_and_join_later() {
473        let root = tempfile::tempdir().unwrap();
474        let database = tempfile::tempdir().unwrap();
475        std::fs::write(root.path().join("f.bin"), vec![0; 4096]).unwrap();
476        let database_path = database.path().join("snapshots.sqlite");
477        let database_path = database_path.to_string_lossy().into_owned();
478        let root_path = root.path().to_string_lossy().into_owned();
479
480        let handle = spawn_scan_json(database_path.clone(), root_path.clone());
481        // The UI-thread contract: the call above returned before any work
482        // finished, and polling never blocks.
483        let progress: Value = serde_json::from_str(&handle.progress_json()).unwrap();
484        assert_eq!(progress["ok"], true);
485        assert!(progress["data"]["state"] == "running" || progress["data"]["state"] == "finished");
486
487        // The join produces the same envelope as the synchronous path.
488        let result: Value = serde_json::from_str(&handle.result_json()).unwrap();
489        assert_eq!(result["ok"], true, "{}", result);
490        assert!(result["data"]["node_count"].as_u64().unwrap() >= 2);
491        assert!(
492            result["data"]["coverage"]["complete"].as_bool().unwrap(),
493            "a completed local scan must be complete"
494        );
495        // Polling after completion is stable, and the result is idempotent.
496        let again: Value = serde_json::from_str(&handle.result_json()).unwrap();
497        assert_eq!(again, result);
498        let progress: Value = serde_json::from_str(&handle.progress_json()).unwrap();
499        assert_eq!(progress["data"]["state"], "finished");
500    }
501
502    #[test]
503    fn a_cancelled_job_reports_honestly_instead_of_half_publishing() {
504        let root = tempfile::tempdir().unwrap();
505        let database = tempfile::tempdir().unwrap();
506        std::fs::write(root.path().join("f.bin"), vec![0; 1024]).unwrap();
507        let database_path = database.path().join("snapshots.sqlite");
508        let database_path = database_path.to_string_lossy().into_owned();
509        let root_path = root.path().to_string_lossy().into_owned();
510
511        let handle = spawn_scan_json(database_path.clone(), root_path.clone());
512        handle.cancel();
513        let result: Value = serde_json::from_str(&handle.result_json()).unwrap();
514        // Either the job finished before the flag landed (honest success) or
515        // it reports cancellation (honest refusal); a half-published graph
516        // is the one outcome this contract forbids.
517        if result["ok"].as_bool().unwrap() {
518            assert!(
519                result["data"]["coverage"]["complete"].as_bool().unwrap(),
520                "an honest success must carry complete coverage"
521            );
522        } else {
523            let error = result["error"].as_str().unwrap();
524            assert!(
525                error.contains("Cancelled") || error.contains("cancelled"),
526                "{error}"
527            );
528        }
529    }
530
531    #[test]
532    fn the_ffi_layer_is_independent_of_any_host_application() {
533        // 9.7 (PF-02, AI-01): with no host database, no host configuration,
534        // and no PruneX on the machine, the full loop still works from an
535        // empty directory.
536        let host = tempfile::tempdir().unwrap();
537        let root = tempfile::tempdir().unwrap();
538        std::fs::write(root.path().join("report.txt"), b"standalone").unwrap();
539        let database_path = host.path().join("snapshots.sqlite");
540        let database_path = database_path.to_string_lossy().into_owned();
541        let root_path = root.path().to_string_lossy().into_owned();
542
543        let scanned: Value =
544            serde_json::from_str(&scan_native_json(database_path.clone(), root_path.clone()))
545                .unwrap();
546        assert_eq!(scanned["ok"], true);
547        let snapshot_id = scanned["data"]["snapshot_id"].as_str().unwrap();
548        let top: Value =
549            serde_json::from_str(&top_json(database_path.clone(), snapshot_id.into(), 1, 10))
550                .unwrap();
551        assert_eq!(top["ok"], true);
552        // The read-only surface is intact, and nothing outside the two
553        // caller-named directories was touched.
554        assert!(host.path().join("snapshots.sqlite").is_file());
555        assert!(root.path().join("report.txt").is_file());
556    }
557}