1use 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#[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#[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#[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#[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#[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#[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 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 pub fn cancel(&self) {
313 self.cancel.store(true, std::sync::atomic::Ordering::SeqCst);
314 }
315
316 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#[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 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
365fn 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 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 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 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 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 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 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 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 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 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 assert!(host.path().join("snapshots.sqlite").is_file());
555 assert!(root.path().join("report.txt").is_file());
556 }
557}