1use crate::*;
2use rusqlite::{Connection, OptionalExtension, params};
3use sha1::{Digest, Sha1};
4use std::{
5 collections::{BTreeSet, HashSet},
6 os::unix::fs::MetadataExt,
7 sync::atomic::{AtomicBool, Ordering},
8};
9
10const TABLES: &str = "CREATE TABLE IF NOT EXISTS file (id INTEGER PRIMARY KEY,parent INTEGER REFERENCES file ON DELETE CASCADE,name TEXT NOT NULL,dir INTEGER NOT NULL,size INTEGER NOT NULL,alloc INTEGER NOT NULL,files INTEGER NOT NULL,mtime INTEGER NOT NULL); CREATE TABLE IF NOT EXISTS meta (snapshot TEXT,scanned INTEGER NOT NULL,updated INTEGER NOT NULL);";
11const INDEXES: &str = "CREATE UNIQUE INDEX IF NOT EXISTS file_child ON file (parent,name); CREATE INDEX IF NOT EXISTS file_parent_size ON file (parent,size); CREATE INDEX IF NOT EXISTS file_size ON file (size) WHERE NOT dir;";
12
13pub struct Index {
14 pool: String,
15 dir: PathBuf,
16 active: Mutex<BTreeSet<String>>,
17 scanning: AtomicBool,
18 wake: tokio::sync::Notify,
19}
20impl Index {
21 pub fn new(pool: String, dir: PathBuf) -> Self {
22 Self {
23 pool,
24 dir,
25 active: Mutex::new(BTreeSet::new()),
26 scanning: AtomicBool::new(false),
27 wake: tokio::sync::Notify::new(),
28 }
29 }
30 fn file(&self, dataset: &str) -> PathBuf {
31 self.dir.join(format!("{}.db", encoded(dataset)))
32 }
33 fn db(&self, dataset: &str) -> Result<Option<Connection>> {
34 if !self.active.lock().unwrap().contains(dataset) {
35 return Ok(None);
36 }
37 let file = self.file(dataset);
38 if !file.exists() {
39 return Ok(None);
40 }
41 let db = Connection::open_with_flags(file, rusqlite::OpenFlags::SQLITE_OPEN_READ_WRITE)?;
42 if meta(&db)?.is_none() {
43 return Ok(None);
44 }
45 Ok(Some(db))
46 }
47 fn names(&self) -> BTreeSet<String> {
48 let mut names = BTreeSet::new();
49 for name in self.active.lock().unwrap().iter() {
50 let parts: Vec<_> = name.split('/').collect();
51 for i in 1..=parts.len() {
52 names.insert(parts[..i].join("/"));
53 }
54 }
55 names
56 }
57 fn owner(&self, rel: &str) -> Option<(String, String)> {
58 let names = self.names();
59 let full = if rel.is_empty() {
60 self.pool.clone()
61 } else {
62 format!("{}/{rel}", self.pool)
63 };
64 let parts: Vec<_> = full.split('/').collect();
65 for i in (1..=parts.len()).rev() {
66 let dataset = parts[..i].join("/");
67 if names.contains(&dataset) {
68 return Some((dataset, parts[i..].join("/")));
69 }
70 }
71 None
72 }
73 fn subs(&self, dataset: &str) -> Vec<String> {
74 self.names()
75 .into_iter()
76 .filter_map(|name| {
77 let (parent, sub) = name.rsplit_once('/')?;
78 (parent == dataset).then(|| sub.to_owned())
79 })
80 .collect()
81 }
82 pub fn children(&self, rel: &str, limit: usize) -> Result<Option<Value>> {
83 let Some((dataset, inner)) = self.owner(rel) else {
84 return Ok(None);
85 };
86 let db = self.db(&dataset)?;
87 let subs = if inner.is_empty() {
88 self.subs(&dataset)
89 } else {
90 Vec::new()
91 };
92 let own = match db.as_ref() {
93 Some(db) => children(db, &inner, limit + subs.len())?,
94 None if inner.is_empty() => {
95 Some(json!({"size":0,"alloc":0,"files":0,"mtime":0,"count":0,"entries":[]}))
96 }
97 _ => None,
98 };
99 let Some(mut own) = own else {
100 return Ok(None);
101 };
102 let mut entries: Vec<_> = array(&own["entries"])
103 .iter()
104 .filter(|e| !subs.iter().any(|s| e["name"] == *s))
105 .cloned()
106 .collect();
107 for sub in subs {
108 let next = if rel.is_empty() {
109 sub.clone()
110 } else {
111 format!("{rel}/{sub}")
112 };
113 let Some(total) = self.children(&next, 0)? else {
114 continue;
115 };
116 let shadow = if let Some(db) = &db {
117 self::children(db, &sub, 0)?.unwrap_or(json!({}))
118 } else {
119 json!({})
120 };
121 for key in ["size", "alloc", "files"] {
122 own[key] = json!(number(&own[key]) - number(&shadow[key]) + number(&total[key]));
123 }
124 own["mtime"] = json!(number(&own["mtime"]).max(number(&total["mtime"])));
125 own["count"] = json!(
126 number(&own["count"]) + 1.0 - if shadow["size"].is_null() { 0.0 } else { 1.0 }
127 );
128 entries.push(json!({"name":sub,"dir":1,"size":total["size"],"alloc":total["alloc"],"files":total["files"],"mtime":total["mtime"]}));
129 }
130 entries.sort_by(|a, b| {
131 number(&b["size"])
132 .total_cmp(&number(&a["size"]))
133 .then_with(|| string(&a["name"]).cmp(string(&b["name"])))
134 });
135 entries.truncate(limit);
136 own["entries"] = json!(entries);
137 Ok(Some(own))
138 }
139 pub fn map(&self, rel: &str, min: f64) -> Result<Option<Value>> {
140 let Some((dataset, inner)) = self.owner(rel) else {
141 return Ok(None);
142 };
143 let db = self.db(&dataset)?;
144 let own = match db.as_ref() {
145 Some(db) => map(db, &inner, min)?,
146 None if inner.is_empty() => Some(json!(["", 0, 0, []])),
147 _ => None,
148 };
149 let Some(mut own) = own else {
150 return Ok(None);
151 };
152 if !inner.is_empty() {
153 return Ok(Some(own));
154 }
155 let subs = self.subs(&dataset);
156 let mut children: Vec<_> = array(&own[3])
157 .iter()
158 .filter(|n| !subs.iter().any(|s| n[0] == *s))
159 .cloned()
160 .collect();
161 for sub in subs {
162 let next = if rel.is_empty() {
163 sub.clone()
164 } else {
165 format!("{rel}/{sub}")
166 };
167 let Some(tree) = self.map(&next, min)? else {
168 continue;
169 };
170 let shadow = if let Some(db) = &db {
171 self::children(db, &sub, 0)?.unwrap_or(json!({}))
172 } else {
173 json!({})
174 };
175 own[1] = json!(number(&own[1]) + number(&tree[1]) - number(&shadow["size"]));
176 own[2] = json!(number(&own[2]) + number(&tree[2]) - number(&shadow["files"]));
177 if number(&tree[1]) >= min {
178 children.push(json!([sub, tree[1], tree[2], tree[3]]));
179 }
180 }
181 own[3] = json!(children);
182 Ok(Some(own))
183 }
184 pub fn largest(&self, limit: usize) -> Result<Value> {
185 let mut values = Vec::new();
186 let datasets = self.active.lock().unwrap().clone();
187 for dataset in datasets {
188 let Some(db) = self.db(&dataset)? else {
189 continue;
190 };
191 let mut query=db.prepare("SELECT (WITH RECURSIVE up(id,name,parent,depth) AS (SELECT id,name,parent,0 FROM file AS f WHERE f.id=top.id UNION ALL SELECT file.id,file.name,file.parent,depth+1 FROM file JOIN up ON file.id=up.parent WHERE file.parent IS NOT NULL) SELECT group_concat(name,'/' ORDER BY depth DESC) FROM up) AS path,size,alloc,mtime FROM file AS top WHERE NOT dir ORDER BY size DESC LIMIT ?")?;
192 let rows=query.query_map([limit],|r| Ok(json!({"path":r.get::<_,String>(0)?,"size":r.get::<_,i64>(1)?,"alloc":r.get::<_,i64>(2)?,"mtime":r.get::<_,i64>(3)?})))?;
193 let prefix = dataset
194 .strip_prefix(&format!("{}/", self.pool))
195 .map(|p| format!("{p}/"))
196 .unwrap_or_default();
197 for row in rows {
198 let mut row = row?;
199 row["path"] = json!(format!("{prefix}{}", string(&row["path"])));
200 values.push(row);
201 }
202 }
203 values.sort_by(|a, b| number(&b["size"]).total_cmp(&number(&a["size"])));
204 values.truncate(limit);
205 Ok(json!(values))
206 }
207 pub fn updated(&self) -> Result<Option<f64>> {
208 let mut values = Vec::new();
209 let datasets = self.active.lock().unwrap().clone();
210 for dataset in datasets {
211 let Some(db) = self.db(&dataset)? else {
212 return Ok(None);
213 };
214 let Some(meta) = meta(&db)? else {
215 return Ok(None);
216 };
217 values.push(number(&meta["updated"]));
218 }
219 Ok(values.into_iter().min_by(f64::total_cmp))
220 }
221 pub async fn start(self: Arc<Self>, app: Arc<App>) {
222 let mut failures = HashMap::<String, u32>::new();
223 loop {
224 let result = async {
225 tokio::fs::create_dir_all(&self.dir).await?;
226 let (listed, mounted) = tokio::try_join!(
227 crate::host::call(json!({"operation":"storage.datasets"})),
228 crate::host::call(json!({"operation":"storage.mounts"}))
229 )?;
230 let origins: HashMap<_, _> = array(&listed)
231 .iter()
232 .filter(|d| {
233 string(&d["name"]) == self.pool
234 || string(&d["name"]).starts_with(&format!("{}/", self.pool))
235 })
236 .map(|d| (string(&d["name"]), &d["origin"]))
237 .collect();
238 let visible =
239 mounted_datasets(&tokio::fs::read_to_string("/proc/self/mountinfo").await?);
240 let mounts: Vec<_> = array(&mounted["filesystems"])
241 .iter()
242 .filter(|m| {
243 origins
244 .get(string(&m["source"]))
245 .is_some_and(|v| v.is_null())
246 && !string(&m["source"]).starts_with(&format!("{}/staging/", self.pool))
247 && visible.contains(&(
248 string(&m["source"]).to_owned(),
249 string(&m["target"]).to_owned(),
250 ))
251 })
252 .map(|m| {
253 (
254 string(&m["source"]).to_owned(),
255 string(&m["target"]).to_owned(),
256 )
257 })
258 .collect();
259 *self.active.lock().unwrap() = mounts.iter().map(|(d, _)| d.clone()).collect();
260 self.scanning.store(true, Ordering::Relaxed);
261 for (dataset, mount) in mounts {
262 let result = self
263 .round(
264 app.clone(),
265 &dataset,
266 &mount,
267 *failures.get(&dataset).unwrap_or(&0),
268 )
269 .await;
270 if let Err(e) = result {
271 eprintln!("index {dataset}: {}", e.message);
272 *failures.entry(dataset).or_default() += 1;
273 } else {
274 failures.remove(&dataset);
275 }
276 }
277 let mut files = tokio::fs::read_dir(&self.dir).await?;
278 while let Some(file) = files.next_entry().await? {
279 if let Some(name) = file
280 .file_name()
281 .to_str()
282 .and_then(|s| s.strip_suffix(".db"))
283 {
284 let name = url::form_urlencoded::parse(format!("x={name}").as_bytes())
285 .next()
286 .map(|(_, v)| v.into_owned())
287 .unwrap_or_default();
288 if !origins.contains_key(name.as_str())
289 || name.starts_with(&format!("{}/staging/", self.pool))
290 {
291 tokio::fs::remove_file(file.path()).await?;
292 }
293 }
294 }
295 Ok::<_, Error>(())
296 }
297 .await;
298 self.scanning.store(false, Ordering::Relaxed);
299 if let Err(e) = result {
300 eprintln!("file index: {}", e.message);
301 }
302 tokio::select! {
303 _ = tokio::time::sleep(Duration::from_secs(600)) => (),
304 _ = self.wake.notified() => (),
305 }
306 }
307 }
308 async fn round(&self, app: Arc<App>, dataset: &str, mount: &str, failures: u32) -> Result<()> {
309 let snapshot: String = serde_json::from_value(
310 crate::host::call(json!({"operation":"index.snapshot","dataset":dataset})).await?,
311 )?;
312 let existing: Vec<String> = serde_json::from_value(
313 crate::host::call(json!({"operation":"index.snapshots","dataset":dataset})).await?,
314 )?;
315 let file = self.file(dataset);
316 let old = if let Some(db) = self.db(dataset)? {
317 meta(&db)?
318 } else {
319 None
320 };
321 let offset = u32::from_be_bytes(Sha1::digest(dataset.as_bytes())[..4].try_into().unwrap())
322 as f64
323 % (7.0 * 86400.0);
324 let incremental = old.as_ref().is_some_and(|m| {
325 m["snapshot"].is_string()
326 && existing.iter().any(|s| s == string(&m["snapshot"]))
327 && failures < 3
328 && ((now() - offset) / (7.0 * 86400.0)).floor()
329 == ((number(&m["scanned"]) - offset) / (7.0 * 86400.0)).floor()
330 });
331 let changes = if incremental {
332 let previous = string(&old.as_ref().unwrap()["snapshot"]);
333 let diff = crate::host::call(json!({"operation":"index.diff","dataset":dataset,
334 "from":previous.split_once('@').map(|(_, name)| name).unwrap_or_default(),
335 "to":snapshot.split_once('@').map(|(_, name)| name).unwrap_or_default()}))
336 .await?;
337 Some(parse_diff(string(&diff), mount)?)
338 } else {
339 None
340 };
341 let _slot = app.heavy.acquire().await?;
342 let frozen = PathBuf::from(mount)
343 .join(".zfs/snapshot")
344 .join(snapshot.split_once('@').unwrap().1);
345 let current = snapshot.clone();
346 tokio::task::spawn_blocking(move || {
347 if let Some(changes) = changes {
348 apply(&file, &frozen, &snapshot, &changes)
349 } else {
350 scan(&file, &frozen, Some(&snapshot))
351 }
352 })
353 .await??;
354 for previous in existing.iter().filter(|s| **s != current) {
355 let name = previous.split_once('@').unwrap().1;
356 if let Err(error) = crate::host::call(
357 json!({"operation":"index.discard","dataset":dataset,"snapshot":name}),
358 )
359 .await
360 {
361 eprintln!("index snapshot {previous}: {}", error.message);
362 }
363 }
364 Ok(())
365 }
366 pub fn changed(&self) {
367 self.wake.notify_one();
368 }
369 pub fn is_scanning(&self) -> bool {
370 self.scanning.load(Ordering::Relaxed)
371 }
372}
373fn mounted_datasets(text: &str) -> HashSet<(String, String)> {
374 fn unescape(value: &str) -> String {
375 value
376 .replace("\\040", " ")
377 .replace("\\011", "\t")
378 .replace("\\012", "\n")
379 .replace("\\134", "\\")
380 }
381 text.lines()
382 .filter_map(|line| {
383 let (mount, filesystem) = line.split_once(" - ")?;
384 let mut fields = mount.split_whitespace();
385 if fields.nth(3)? != "/" {
386 return None;
387 }
388 let target = fields.next()?;
389 let mut fields = filesystem.split_whitespace();
390 if fields.next()? != "zfs" {
391 return None;
392 }
393 let source = fields.next()?;
394 (!source.contains('@')).then(|| (unescape(source), unescape(target)))
395 })
396 .collect()
397}
398fn resolve(db: &Connection, rel: &str) -> Result<Option<Vec<i64>>> {
399 let mut chain = vec![1];
400 for name in rel.split('/').filter(|s| !s.is_empty()) {
401 let id = db
402 .query_row(
403 "SELECT id FROM file WHERE parent=? AND name=?",
404 params![chain.last().unwrap(), name],
405 |r| r.get::<_, i64>(0),
406 )
407 .optional()?;
408 let Some(id) = id else {
409 return Ok(None);
410 };
411 chain.push(id);
412 }
413 Ok(Some(chain))
414}
415fn meta(db: &Connection) -> Result<Option<Value>> {
416 Ok(db.query_row("SELECT snapshot,scanned,updated FROM meta",[],|r| Ok(json!({"snapshot":r.get::<_,Option<String>>(0)?,"scanned":r.get::<_,i64>(1)?,"updated":r.get::<_,i64>(2)?}))).optional()?)
417}
418fn children(db: &Connection, rel: &str, limit: usize) -> Result<Option<Value>> {
419 let Some(chain) = resolve(db, rel)? else {
420 return Ok(None);
421 };
422 let id = chain.last().unwrap();
423 let mut value=db.query_row("SELECT size,alloc,files,mtime,(SELECT count(*) FROM file WHERE parent=?) FROM file WHERE id=?",params![id,id],|r| Ok(json!({"size":r.get::<_,i64>(0)?,"alloc":r.get::<_,i64>(1)?,"files":r.get::<_,i64>(2)?,"mtime":r.get::<_,i64>(3)?,"count":r.get::<_,i64>(4)?})))?;
424 let mut query=db.prepare("SELECT name,dir,size,alloc,files,mtime FROM file WHERE parent=? ORDER BY size DESC,name LIMIT ?")?;
425 let rows=query.query_map(params![id,limit],|r| Ok(json!({"name":r.get::<_,String>(0)?,"dir":r.get::<_,i64>(1)?,"size":r.get::<_,i64>(2)?,"alloc":r.get::<_,i64>(3)?,"files":r.get::<_,i64>(4)?,"mtime":r.get::<_,i64>(5)?})))?.collect::<std::result::Result<Vec<_>,_>>()?;
426 value["entries"] = json!(rows);
427 Ok(Some(value))
428}
429fn map(db: &Connection, rel: &str, min: f64) -> Result<Option<Value>> {
430 let Some(chain) = resolve(db, rel)? else {
431 return Ok(None);
432 };
433 let id = *chain.last().unwrap();
434 let mut query=db.prepare("WITH RECURSIVE tree(id,parent,name,dir,size,files) AS (SELECT id,parent,name,dir,size,files FROM file WHERE id=? UNION ALL SELECT file.id,file.parent,file.name,file.dir,file.size,file.files FROM file JOIN tree ON file.parent=tree.id WHERE tree.dir AND file.size>=?) SELECT * FROM tree")?;
435 let rows = query
436 .query_map(params![id, min], |r| {
437 Ok((
438 r.get::<_, i64>(0)?,
439 r.get::<_, Option<i64>>(1)?,
440 r.get::<_, String>(2)?,
441 r.get::<_, bool>(3)?,
442 r.get::<_, i64>(4)?,
443 r.get::<_, i64>(5)?,
444 ))
445 })?
446 .collect::<std::result::Result<Vec<_>, _>>()?;
447 let nodes: HashMap<_, _> = rows
448 .iter()
449 .map(|(id, _, name, dir, size, files)| {
450 (
451 *id,
452 if *dir {
453 json!([name, size, files, []])
454 } else {
455 json!([name, size])
456 },
457 )
458 })
459 .collect();
460 let mut children = HashMap::<i64, Vec<i64>>::new();
461 for (id, parent, ..) in rows {
462 if let Some(parent) = parent {
463 children.entry(parent).or_default().push(id);
464 }
465 }
466 fn assemble(id: i64, nodes: &HashMap<i64, Value>, children: &HashMap<i64, Vec<i64>>) -> Value {
467 let mut value = nodes[&id].clone();
468 if value.as_array().unwrap().len() == 4 {
469 value[3] = json!(
470 children
471 .get(&id)
472 .into_iter()
473 .flatten()
474 .map(|id| assemble(*id, nodes, children))
475 .collect::<Vec<_>>()
476 );
477 }
478 value
479 }
480 Ok(Some(assemble(id, &nodes, &children)))
481}
482fn scan(file: &std::path::Path, root: &std::path::Path, snapshot: Option<&str>) -> Result<()> {
483 let next = file.with_extension("db.new");
484 let _ = std::fs::remove_file(&next);
485 let db = Connection::open(&next)?;
486 db.execute_batch(&format!(
487 "PRAGMA journal_mode=OFF; PRAGMA synchronous=OFF;{TABLES} BEGIN;"
488 ))?;
489 let handle = std::fs::File::open(root)?;
490 let top = handle.metadata()?;
491 db.execute(
492 "INSERT INTO file VALUES(1,NULL,'',1,0,0,0,?)",
493 [top.mtime()],
494 )?;
495 fn walk(
496 db: &Connection,
497 id: i64,
498 root: &std::path::Path,
499 device: u64,
500 ) -> Result<(u64, u64, u64)> {
501 let mut sum = (0, 0, 0);
502 for entry in std::fs::read_dir(root)? {
503 let entry = entry?;
504 let info = std::fs::symlink_metadata(entry.path())?;
505 if info.is_dir() && info.dev() != device {
506 continue;
507 }
508 db.execute(
509 "INSERT INTO file(parent,name,dir,size,alloc,files,mtime) VALUES(?,?,?,?,?,?,?)",
510 params![
511 id,
512 entry.file_name().to_string_lossy(),
513 info.is_dir(),
514 if info.is_dir() { 0 } else { info.len() },
515 if info.is_dir() {
516 0
517 } else {
518 info.blocks() * 512
519 },
520 if info.is_dir() { 0 } else { 1 },
521 info.mtime()
522 ],
523 )?;
524 let child = db.last_insert_rowid();
525 let (size, alloc, files) = if info.is_dir() {
526 walk(db, child, &entry.path(), device)?
527 } else {
528 (info.len(), info.blocks() * 512, 1)
529 };
530 if info.is_dir() {
531 db.execute(
532 "UPDATE file SET size=?,alloc=?,files=? WHERE id=?",
533 params![size, alloc, files, child],
534 )?;
535 }
536 sum.0 += size;
537 sum.1 += alloc;
538 sum.2 += files;
539 }
540 Ok(sum)
541 }
542 let (size, alloc, files) = walk(&db, 1, root, top.dev())?;
543 db.execute(
544 "UPDATE file SET size=?,alloc=?,files=? WHERE id=1",
545 params![size, alloc, files],
546 )?;
547 db.execute(
548 "INSERT INTO meta VALUES(?,?,?)",
549 params![snapshot, now() as u64, now() as u64],
550 )?;
551 db.execute_batch(&format!("{INDEXES} COMMIT;"))?;
552 db.close().map_err(|(_, e)| Error::from(e))?;
553 std::fs::rename(next, file)?;
554 Ok(())
555}
556#[derive(Debug)]
557struct Change {
558 kind: char,
559 path: String,
560 to: Option<String>,
561}
562fn parse_diff(text: &str, mount: &str) -> Result<Vec<Change>> {
563 let mut changes = Vec::new();
564 let mount = std::path::Path::new(mount);
565 for line in text.lines().filter(|s| !s.is_empty()) {
566 let fields: Vec<_> = line.split('\t').collect();
567 if fields.len() < 3 || !["+", "-", "M", "R"].contains(&fields[0]) {
568 return Err(Error::new(500, format!("unexpected zfs diff line: {line}")));
569 }
570 let relative = |s: &str| {
571 let path = crate::storage::unescape(s);
572 std::path::Path::new(&path)
573 .strip_prefix(mount)
574 .map(|p| p.to_string_lossy().into_owned())
575 .map_err(|_| {
576 Error::new(
577 500,
578 format!("zfs diff path outside {}: {path}", mount.display()),
579 )
580 })
581 };
582 let to = if fields[0] == "R" {
583 Some(relative(
584 fields
585 .get(3)
586 .ok_or_else(|| Error::new(500, "Invalid rename."))?,
587 )?)
588 } else {
589 None
590 };
591 changes.push(Change {
592 kind: fields[0].chars().next().unwrap(),
593 path: relative(fields[2])?,
594 to,
595 });
596 }
597 Ok(changes)
598}
599fn apply(
600 file: &std::path::Path,
601 root: &std::path::Path,
602 snapshot: &str,
603 changes: &[Change],
604) -> Result<()> {
605 let db = Connection::open(file)?;
606 db.execute_batch("PRAGMA foreign_keys=ON;")?;
607 let targets: HashSet<String> = changes
608 .iter()
609 .filter(|c| c.kind != '-')
610 .map(|c| c.to.as_ref().unwrap_or(&c.path).clone())
611 .collect();
612 let mut listings = HashMap::<String, Option<Vec<std::fs::DirEntry>>>::new();
613 let mut wanted = HashSet::new();
614 for target in &targets {
615 let entries = match std::fs::read_dir(root.join(target)) {
616 Ok(dir) => Some(dir.collect::<std::io::Result<Vec<_>>>()?),
617 Err(error)
618 if matches!(
619 error.kind(),
620 std::io::ErrorKind::NotFound | std::io::ErrorKind::NotADirectory
621 ) =>
622 {
623 None
624 }
625 Err(error) => return Err(error.into()),
626 };
627 if let Some(entries) = &entries {
628 for entry in entries {
629 if !entry.file_type()?.is_dir() {
630 wanted.insert(
631 format!("{}/{}", target, entry.file_name().to_string_lossy())
632 .trim_start_matches('/')
633 .to_owned(),
634 );
635 }
636 }
637 }
638 listings.insert(target.clone(), entries);
639 let mut path = std::path::Path::new(target);
640 while !path.as_os_str().is_empty() {
641 wanted.insert(path.to_string_lossy().into_owned());
642 path = path.parent().unwrap_or(std::path::Path::new(""));
643 }
644 }
645 let stats: HashMap<_, _> = wanted
646 .into_iter()
647 .map(|p| {
648 let info = match std::fs::symlink_metadata(root.join(&p)) {
649 Ok(info) => Some(info),
650 Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
651 Err(error) => return Err(Error::from(error)),
652 };
653 Ok((p, info))
654 })
655 .collect::<Result<_>>()?;
656 let mut dirty = HashSet::<i64>::new();
657 let moves: Vec<_> = changes
658 .iter()
659 .filter(|c| c.kind == 'R')
660 .map(|c| Ok((resolve(&db, &c.path)?, c.to.as_ref().unwrap().clone())))
661 .collect::<Result<_>>()?;
662 let removed: Vec<_> = changes
663 .iter()
664 .filter(|c| c.kind == '-')
665 .map(|c| resolve(&db, &c.path))
666 .collect::<Result<_>>()?;
667 let mut placed: Vec<_> = moves
668 .iter()
669 .map(|(chain, to)| (to.clone(), chain.as_ref().and_then(|c| c.last()).copied()))
670 .chain(
671 changes
672 .iter()
673 .filter(|c| c.kind == '+' || c.kind == 'M')
674 .map(|c| (c.path.clone(), None)),
675 )
676 .collect();
677 placed.sort_by_key(|(p, _)| p.split('/').count());
678 let upsert = "INSERT INTO file(parent,name,dir,size,alloc,files,mtime) VALUES(?,?,?,?,?,?,?) ON CONFLICT(parent,name) DO UPDATE SET dir=excluded.dir,size=excluded.size,alloc=excluded.alloc,files=excluded.files,mtime=excluded.mtime RETURNING id";
679 db.execute_batch("BEGIN;")?;
680 let result = || -> Result<()> {
681 for (chain, _) in &moves {
682 if let Some(chain) = chain {
683 dirty.extend(chain);
684 db.execute(
685 "UPDATE file SET parent=NULL WHERE id=?",
686 [chain.last().unwrap()],
687 )?;
688 }
689 }
690 for chain in removed.iter().flatten() {
691 dirty.extend(chain);
692 db.execute("DELETE FROM file WHERE id=?", [chain.last().unwrap()])?;
693 }
694 for (rel, id) in placed {
695 if rel.is_empty() {
696 continue;
697 }
698 let names: Vec<_> = rel.split('/').collect();
699 let mut parent = 1i64;
700 for (i, name) in names.iter().enumerate() {
701 let Some(Some(info)) = stats.get(&names[..=i].join("/")) else {
702 break;
703 };
704 if let Some(id) = id.filter(|_| i == names.len() - 1) {
705 db.execute(
706 "DELETE FROM file WHERE parent=? AND name=? AND id!=?",
707 params![parent, name, id],
708 )?;
709 db.execute(
710 "UPDATE file SET parent=?,name=? WHERE id=?",
711 params![parent, name, id],
712 )?;
713 }
714 dirty.insert(parent);
715 parent = db.query_row(
716 upsert,
717 params![
718 parent,
719 name,
720 info.is_dir(),
721 if info.is_dir() { 0 } else { info.len() },
722 if info.is_dir() {
723 0
724 } else {
725 info.blocks() * 512
726 },
727 if info.is_dir() { 0 } else { 1 },
728 info.mtime()
729 ],
730 |r| r.get(0),
731 )?;
732 dirty.insert(parent);
733 }
734 }
735 for (rel, entries) in listings {
736 let Some(entries) = entries else {
737 continue;
738 };
739 let Some(chain) = resolve(&db, &rel)? else {
740 continue;
741 };
742 dirty.extend(&chain);
743 let parent = *chain.last().unwrap();
744 let present: HashSet<_> = entries
745 .iter()
746 .map(|e| e.file_name().to_string_lossy().into_owned())
747 .collect();
748 let mut query = db.prepare("SELECT id,name FROM file WHERE parent=?")?;
749 let children = query
750 .query_map([parent], |r| {
751 Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?))
752 })?
753 .collect::<std::result::Result<Vec<_>, _>>()?;
754 for (id, name) in children {
755 if !present.contains(&name) {
756 db.execute("DELETE FROM file WHERE id=?", [id])?;
757 }
758 }
759 for entry in entries {
760 let name = entry.file_name().to_string_lossy().into_owned();
761 let path = if rel.is_empty() {
762 name.clone()
763 } else {
764 format!("{rel}/{name}")
765 };
766 let Some(Some(info)) = stats.get(&path) else {
767 continue;
768 };
769 if info.is_dir() {
770 continue;
771 }
772 let _: i64 = db.query_row(
773 upsert,
774 params![
775 parent,
776 name,
777 false,
778 info.len(),
779 info.blocks() * 512,
780 1,
781 info.mtime()
782 ],
783 |r| r.get(0),
784 )?;
785 }
786 }
787 db.execute("DELETE FROM file WHERE parent IS NULL AND id!=1", [])?;
788 let mut order = Vec::new();
789 for id in dirty {
790 let depth:i64=db.query_row("WITH RECURSIVE up(id) AS(SELECT ? UNION ALL SELECT file.parent FROM file JOIN up USING(id) WHERE file.parent IS NOT NULL) SELECT count(*) FROM up",[id],|r| r.get(0))?;
791 order.push((id, depth));
792 }
793 order.sort_by_key(|(_, depth)| std::cmp::Reverse(*depth));
794 for (id, _) in order {
795 db.execute("UPDATE file SET (size,alloc,files)=(SELECT coalesce(sum(size),0),coalesce(sum(alloc),0),coalesce(sum(files),0) FROM file AS child WHERE child.parent=file.id) WHERE id=? AND dir",[id])?;
796 }
797 db.execute(
798 "UPDATE meta SET snapshot=?,updated=?",
799 params![snapshot, now() as u64],
800 )?;
801 Ok(())
802 }();
803 if result.is_ok() {
804 db.execute_batch("COMMIT;")?;
805 } else {
806 db.execute_batch("ROLLBACK;")?;
807 }
808 result
809}
810
811pub async fn route(app: Arc<App>, kind: &str, query: &HashMap<String, String>) -> Result<Response> {
812 let index = app
813 .index
814 .as_ref()
815 .ok_or_else(|| {
816 Error::new(
817 501,
818 "The file index is off. Enable it on this home server, then retry.",
819 )
820 })?
821 .clone();
822 let _slot = app.heavy.acquire().await?;
823 let root = env("STUDIO_STORE_ROOT", "/srv");
824 let rel = crate::files::relative(
825 query.get("path").map(String::as_str).unwrap_or_default(),
826 false,
827 )?;
828 let largest = kind == "largest";
829 let (width, height) = if largest {
830 (1.0, 1.0)
831 } else {
832 (
833 telemetry::query_number(query, "width", 0.0, 1.0, 8192.0)?,
834 telemetry::query_number(query, "height", 0.0, 1.0, 8192.0)?,
835 )
836 };
837 let value=tokio::task::spawn_blocking(move || {
838 if largest {let mut files=index.largest(100)?;for file in files.as_array_mut().unwrap() {file["link"]=apps::file_link(&format!("{root}/{}",string(&file["path"])),false);}return Ok::<_,Error>(files);}
839 let total=index.children(&rel,0)?;let min=(total.as_ref().map(|v| number(&v["size"])).unwrap_or(0.0)*16.0/(width*height)).max(1.0).ceil();
840 Ok(json!({"root":root,"min":min,"scanning":index.is_scanning(),"tree":if total.is_some() {index.map(&rel,min)?} else {None}}))
841 }).await??;
842 Ok(Document::new(value).response())
843}
844
845#[cfg(test)]
846mod tests {
847 use super::*;
848 use std::fs;
849
850 #[test]
851 fn reads_recover_interrupted_updates() {
852 if let Ok(file) = std::env::var("STUDIO_INDEX_CRASH_TEST_DB") {
853 let db = Connection::open(file).unwrap();
854 db.execute_batch("PRAGMA cache_size=1; BEGIN IMMEDIATE; UPDATE file SET size=99;")
855 .unwrap();
856 std::process::exit(0);
857 }
858 let dir = std::env::temp_dir().join(format!("snowglobe-index-{}", uuid::Uuid::new_v4()));
859 let root = dir.join("tree");
860 fs::create_dir_all(&root).unwrap();
861 for n in 0..512 {
862 fs::write(root.join(format!("{n}-{}", "x".repeat(150))), [0]).unwrap();
863 }
864 let index = Index::new("tank".into(), dir.clone());
865 index.active.lock().unwrap().insert("tank".into());
866 let file = index.file("tank");
867 scan(&file, &root, Some("tank@1")).unwrap();
868 let thread = std::thread::current();
869 assert!(
870 std::process::Command::new(std::env::current_exe().unwrap())
871 .args(["--exact", thread.name().unwrap()])
872 .env("STUDIO_INDEX_CRASH_TEST_DB", &file)
873 .status()
874 .unwrap()
875 .success()
876 );
877 let journal = fs::read(file.with_extension("db-journal")).unwrap();
878 assert!(journal.len() > 512 && journal[..8].iter().any(|byte| *byte != 0));
879 let tree = index.children("", 0).unwrap().unwrap();
880 assert_eq!(tree["size"], 512);
881 assert_eq!(tree["files"], 512);
882 fs::remove_dir_all(dir).unwrap();
883 }
884
885 fn dump(file: &std::path::Path) -> Vec<Value> {
886 let db = Connection::open(file).unwrap();
887 let mut query = db.prepare("WITH RECURSIVE walk(id,path) AS (SELECT id,'' FROM file WHERE parent IS NULL UNION ALL SELECT file.id,walk.path||'/'||file.name FROM file JOIN walk ON file.parent=walk.id) SELECT path,dir,size,alloc,files,CASE WHEN path='' THEN 0 ELSE mtime END FROM walk JOIN file USING(id) ORDER BY path").unwrap();
888 query
889 .query_map([], |r| {
890 Ok(json!([
891 r.get::<_, String>(0)?,
892 r.get::<_, bool>(1)?,
893 r.get::<_, i64>(2)?,
894 r.get::<_, i64>(3)?,
895 r.get::<_, i64>(4)?,
896 r.get::<_, i64>(5)?
897 ]))
898 })
899 .unwrap()
900 .map(|r| r.unwrap())
901 .collect()
902 }
903 #[test]
904 fn incremental_changes_match_fresh_scans_in_any_order() {
905 let cases: Value = serde_json::from_str(include_str!("../tests/index-cases.json")).unwrap();
906 for case in array(&cases) {
907 for seed in 0..6u64 {
908 let dir =
909 std::env::temp_dir().join(format!("snowglobe-index-{}", uuid::Uuid::new_v4()));
910 let root = dir.join("tree");
911 fs::create_dir_all(&root).unwrap();
912 for (name, value) in case["before"].as_object().unwrap() {
913 let path = root.join(name);
914 fs::create_dir_all(if name.ends_with('/') {
915 path.as_path()
916 } else {
917 path.parent().unwrap()
918 })
919 .unwrap();
920 if let Some(link) = value.as_str() {
921 fs::hard_link(root.join(link), path).unwrap();
922 } else if !name.ends_with('/') {
923 fs::write(path, vec![0; number(value) as usize]).unwrap();
924 }
925 }
926 let old = dir.join("old.db");
927 let fresh = dir.join("fresh.db");
928 scan(&old, &root, Some("t@1")).unwrap();
929 for op in array(&case["ops"]) {
930 let from = root.join(string(&op[1]));
931 match string(&op[0]) {
932 "renameSync" => fs::rename(from, root.join(string(&op[2]))).unwrap(),
933 "rmSync" => {
934 if from.is_dir() {
935 fs::remove_dir_all(from).unwrap()
936 } else {
937 fs::remove_file(from).unwrap()
938 }
939 }
940 "mkdirSync" => fs::create_dir_all(from).unwrap(),
941 "writeFileSync" => {
942 fs::write(from, vec![0; number(&op[2]) as usize]).unwrap()
943 }
944 "linkSync" => fs::hard_link(from, root.join(string(&op[2]))).unwrap(),
945 other => panic!("unknown fixture operation {other}"),
946 }
947 }
948 let mut lines = array(&case["diff"]).to_vec();
949 let mut state = seed;
950 if seed > 0 {
951 for i in (1..lines.len()).rev() {
952 state = (state * 1103515245 + 12345) % 2147483648;
953 lines.swap(i, state as usize % (i + 1));
954 }
955 }
956 let diff = lines.iter().map(string).collect::<Vec<_>>().join("\n");
957 apply(&old, &root, "t@2", &parse_diff(&diff, "/t").unwrap()).unwrap();
958 scan(&fresh, &root, None).unwrap();
959 assert_eq!(
960 dump(&old),
961 dump(&fresh),
962 "{} seed {seed}: {diff}",
963 string(&case["name"])
964 );
965 assert_eq!(
966 meta(&Connection::open(&old).unwrap()).unwrap().unwrap()["snapshot"],
967 "t@2"
968 );
969 fs::remove_dir_all(dir).unwrap();
970 }
971 }
972 }
973 #[test]
974 fn diff_unescapes_names_and_rejects_paths_outside_mount() {
975 let changes=parse_diff("+\tF\t/tank/r/caf\\0303\\0251\\0040notes\nR\tF\t/tank/r/a\\0134b\t/tank/r/tab\\0011name", "/tank/r").unwrap();
976 assert_eq!(changes[0].path, "café notes");
977 assert_eq!(changes[1].path, "a\\b");
978 assert_eq!(changes[1].to.as_deref(), Some("tab\tname"));
979 assert!(parse_diff("+\tF\t/other/x", "/tank/r").is_err());
980 assert!(parse_diff("1789\t+\tF\t/tank/r/x", "/tank/r").is_err());
981 }
982 #[test]
983 fn namespace_mounts_exclude_subdirectory_binds_and_snapshots() {
984 let text = "1 0 0:1 / /srv/clover rw shared:1 - zfs tank/clover rw\n2 0 0:2 /data /srv/prod/yt-feed/data rw - zfs tank/prod/yt-feed rw\n3 0 0:3 / /srv/clover/Media\\040Files ro - zfs tank/media\\040files rw\n4 0 0:4 / /srv/clover/.zfs/snapshot/index-1 ro - zfs tank/clover@index-1 ro\n5 0 0:5 / /srv/prod rw - overlay overlay rw\nmalformed";
985 assert_eq!(
986 mounted_datasets(text),
987 HashSet::from([
988 ("tank/clover".into(), "/srv/clover".into()),
989 ("tank/media files".into(), "/srv/clover/Media Files".into()),
990 ])
991 );
992 }
993 #[test]
994 fn mounted_datasets_replace_shadowed_totals() {
995 let dir = std::env::temp_dir().join(format!("snowglobe-index-{}", uuid::Uuid::new_v4()));
996 let root = dir.join("tree");
997 let nested = dir.join("nested");
998 let dbs = dir.join("dbs");
999 fs::create_dir_all(root.join("media")).unwrap();
1000 fs::create_dir_all(&nested).unwrap();
1001 fs::create_dir_all(&dbs).unwrap();
1002 fs::write(root.join("media/shadow"), [0; 100]).unwrap();
1003 fs::write(root.join("other"), [0; 10]).unwrap();
1004 fs::write(nested.join("film"), [0; 500]).unwrap();
1005 let index = Index::new("tank".into(), dbs);
1006 index
1007 .active
1008 .lock()
1009 .unwrap()
1010 .extend(["tank".into(), "tank/media".into()]);
1011 scan(&index.file("tank"), &root, None).unwrap();
1012 scan(&index.file("tank/media"), &nested, None).unwrap();
1013 let top = index.children("", 10).unwrap().unwrap();
1014 assert_eq!(number(&top["size"]), 510.0);
1015 assert_eq!(number(&top["files"]), 2.0);
1016 assert_eq!(number(&index.map("", 50.0).unwrap().unwrap()[1]), 510.0);
1017 index.active.lock().unwrap().remove("tank");
1018 assert_eq!(
1019 number(&index.children("", 10).unwrap().unwrap()["size"]),
1020 500.0
1021 );
1022 assert_eq!(number(&index.map("", 50.0).unwrap().unwrap()[1]), 500.0);
1023 assert!(index.children("other", 10).unwrap().is_none());
1024 fs::remove_dir_all(dir).unwrap();
1025 }
1026}