1use std::collections::HashMap;
14use std::fs;
15use std::io::{Read, Write};
16use std::path::{Path, PathBuf};
17use std::sync::{Arc, Mutex as StdMutex, OnceLock};
18use std::time::{Duration, SystemTime};
19
20use async_trait::async_trait;
21
22use khive_storage::blob::{
23 BlobOrphanSweepConfig, BlobOrphanSweepResult, BlobStore, ContentRef, MAX_BLOB_WHOLE_BYTES,
24};
25use khive_storage::error::StorageError;
26use khive_storage::types::{SqlRow, SqlStatement, SqlValue, StorageResult};
27use khive_storage::{AtomicUnitOp, SqlAccess, StorageCapability};
28
29use crate::error::SqliteError;
30use uuid::Uuid;
31
32const ROOT_WRITE_LOCK_FILE: &str = ".khive-blob-write.lock";
33const DATABASE_GC_LOCK_SUFFIX: &str = ".khive-blob-gc.lock";
34const BLOB_GC_CLAIM_BATCH_SIZE: usize = 128;
38
39fn map_io_err(e: std::io::Error, op: &'static str) -> StorageError {
40 StorageError::driver(StorageCapability::Blob, op, e)
41}
42
43#[cfg(any(test, not(unix)))]
44fn shard_path(root: &Path, content_ref: &ContentRef) -> PathBuf {
45 let hex = content_ref.as_str();
46 root.join(&hex[0..2]).join(&hex[2..4]).join(hex)
47}
48
49#[cfg(unix)]
50fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
51 open_dir_no_follow(root)
52}
53
54#[cfg(windows)]
55fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
56 use std::fs::OpenOptions;
57 use std::os::windows::fs::OpenOptionsExt;
58
59 const FILE_SHARE_READ: u32 = 0x1;
60 const FILE_SHARE_WRITE: u32 = 0x2;
61 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
62 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
63
64 let handle = OpenOptions::new()
65 .read(true)
66 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
70 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
71 .open(root)?;
72 let file_type = handle.metadata()?.file_type();
73 if file_type.is_symlink() || !file_type.is_dir() {
74 return Err(std::io::Error::new(
75 std::io::ErrorKind::InvalidInput,
76 format!(
77 "blob store root is not a directory or is a reparse point: {}",
78 root.display()
79 ),
80 ));
81 }
82 Ok(handle)
83}
84
85#[cfg(not(any(unix, windows)))]
86fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
87 let handle = std::fs::File::open(root)?;
88 if !handle.metadata()?.is_dir() {
89 return Err(std::io::Error::new(
90 std::io::ErrorKind::InvalidInput,
91 format!("blob store root is not a directory: {}", root.display()),
92 ));
93 }
94 Ok(handle)
95}
96
97#[cfg(unix)]
103fn verify_blob_root_identity(root: &Path, root_handle: &std::fs::File) -> std::io::Result<()> {
104 use std::os::unix::fs::MetadataExt;
105
106 let current = open_dir_no_follow(root).map_err(|error| {
107 std::io::Error::new(
108 std::io::ErrorKind::InvalidInput,
109 format!(
110 "blob store root is no longer reachable as its initialization-time directory ({}): {error}",
111 root.display()
112 ),
113 )
114 })?;
115 let expected = root_handle.metadata()?;
116 let current = current.metadata()?;
117 if expected.dev() != current.dev() || expected.ino() != current.ino() {
118 return Err(std::io::Error::new(
119 std::io::ErrorKind::InvalidInput,
120 format!(
121 "blob store root no longer names its initialization-time directory: {}",
122 root.display()
123 ),
124 ));
125 }
126 Ok(())
127}
128
129#[cfg(not(unix))]
130fn verify_blob_root_identity(root: &Path, _root_handle: &std::fs::File) -> std::io::Result<()> {
131 let current = root.canonicalize().map_err(|error| {
137 std::io::Error::new(
138 std::io::ErrorKind::InvalidInput,
139 format!(
140 "blob store root is no longer reachable as its initialization-time directory ({}): {error}",
141 root.display()
142 ),
143 )
144 })?;
145 if current != root {
146 return Err(std::io::Error::new(
147 std::io::ErrorKind::InvalidInput,
148 format!(
149 "blob store root no longer names its initialization-time directory: {}",
150 root.display()
151 ),
152 ));
153 }
154 Ok(())
155}
156
157#[cfg(unix)]
175fn unlink_blob_shard_file_no_follow(
176 root: &Path,
177 root_handle: &std::fs::File,
178 content_ref: &ContentRef,
179) -> std::io::Result<()> {
180 use std::os::unix::io::AsRawFd;
181
182 verify_blob_root_identity(root, root_handle)?;
183 let hex = content_ref.as_str();
184 let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
185 let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
186 unlink_entry_at(shard2_dir.as_raw_fd(), hex)
189}
190
191#[cfg(windows)]
192fn unlink_blob_shard_file_no_follow(
193 root: &Path,
194 root_handle: &std::fs::File,
195 content_ref: &ContentRef,
196) -> std::io::Result<()> {
197 use std::fs::OpenOptions;
240 use std::os::windows::ffi::OsStringExt;
241 use std::os::windows::fs::OpenOptionsExt;
242 use std::os::windows::io::AsRawHandle;
243 use windows_sys::Win32::Storage::FileSystem::{
244 FileDispositionInfo, GetFinalPathNameByHandleW, SetFileInformationByHandle,
245 FILE_DISPOSITION_INFO,
246 };
247
248 const FILE_SHARE_READ: u32 = 0x1;
249 const FILE_SHARE_WRITE: u32 = 0x2;
250 const DELETE: u32 = 0x0001_0000;
251 const FILE_READ_ATTRIBUTES: u32 = 0x80;
252 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
253 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
254 const FINAL_PATH_FLAGS: u32 = 0x0;
258
259 fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
260 let dir = OpenOptions::new()
261 .read(true)
262 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
263 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
264 .open(path)?;
265 let file_type = dir.metadata()?.file_type();
266 if file_type.is_symlink() || !file_type.is_dir() {
267 return Err(std::io::Error::new(
268 std::io::ErrorKind::InvalidInput,
269 format!(
270 "refusing to unlink blob shard file through non-directory or \
271 reparse-point path component: {}",
272 path.display()
273 ),
274 ));
275 }
276 Ok(dir)
277 }
278
279 fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<std::path::PathBuf> {
283 let handle = file.as_raw_handle();
284 let mut buf: Vec<u16> = vec![0; 512];
285 loop {
286 let len = unsafe {
287 GetFinalPathNameByHandleW(
288 handle as _,
289 buf.as_mut_ptr(),
290 buf.len() as u32,
291 FINAL_PATH_FLAGS,
292 )
293 };
294 if len == 0 {
295 return Err(std::io::Error::last_os_error());
296 }
297 let len = len as usize;
298 if len <= buf.len() {
299 buf.truncate(len);
300 return Ok(std::path::PathBuf::from(std::ffi::OsString::from_wide(
301 &buf,
302 )));
303 }
304 buf.resize(len, 0);
307 }
308 }
309
310 verify_blob_root_identity(root, root_handle)?;
311 let hex = content_ref.as_str();
312 let shard1 = root.join(&hex[0..2]);
313 let shard2 = shard1.join(&hex[2..4]);
314 let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
315 let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
316
317 let expected = final_path_by_handle(root_handle)?
318 .join(&hex[0..2])
319 .join(&hex[2..4])
320 .join(hex);
321
322 let target = OpenOptions::new()
328 .access_mode(DELETE | FILE_READ_ATTRIBUTES)
329 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
330 .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
331 .open(shard2.join(hex))?;
332
333 let resolved = final_path_by_handle(&target)?;
334 if resolved != expected {
335 return Err(std::io::Error::new(
336 std::io::ErrorKind::InvalidInput,
337 format!(
338 "refusing blob delete: handle resolved outside the verified blob root \
339 (expected {}, resolved {})",
340 expected.display(),
341 resolved.display()
342 ),
343 ));
344 }
345
346 let disposition = FILE_DISPOSITION_INFO { DeleteFile: true };
347 let ok = unsafe {
348 SetFileInformationByHandle(
349 target.as_raw_handle() as _,
350 FileDispositionInfo,
351 std::ptr::from_ref(&disposition).cast(),
352 std::mem::size_of::<FILE_DISPOSITION_INFO>() as u32,
353 )
354 };
355 if ok == 0 {
356 return Err(std::io::Error::last_os_error());
357 }
358 Ok(())
359}
360
361#[cfg(not(any(unix, windows)))]
362fn unlink_blob_shard_file_no_follow(
363 root: &Path,
364 root_handle: &std::fs::File,
365 content_ref: &ContentRef,
366) -> std::io::Result<()> {
367 verify_blob_root_identity(root, root_handle)?;
381 let hex = content_ref.as_str();
382 let shard1 = root.join(&hex[0..2]);
383 let shard2 = shard1.join(&hex[2..4]);
384 for component in [root, shard1.as_path(), shard2.as_path()] {
385 let metadata = fs::symlink_metadata(component)?;
386 if metadata.file_type().is_symlink() {
387 return Err(std::io::Error::new(
388 std::io::ErrorKind::InvalidInput,
389 format!(
390 "refusing to unlink blob shard file through symlinked path component: {}",
391 component.display()
392 ),
393 ));
394 }
395 }
396 fs::remove_file(shard2.join(hex))
397}
398
399#[cfg(unix)]
403fn open_blob_shard_file_no_follow(
404 root: &Path,
405 root_handle: &std::fs::File,
406 content_ref: &ContentRef,
407) -> std::io::Result<std::fs::File> {
408 verify_blob_root_identity(root, root_handle)?;
409 open_blob_shard_file_at_no_follow(root_handle, content_ref, libc::O_RDONLY)
410}
411
412#[cfg(windows)]
413fn open_blob_shard_file_no_follow(
414 root: &Path,
415 root_handle: &std::fs::File,
416 content_ref: &ContentRef,
417) -> std::io::Result<std::fs::File> {
418 use std::fs::OpenOptions;
419 use std::os::windows::ffi::OsStringExt;
420 use std::os::windows::fs::OpenOptionsExt;
421 use std::os::windows::io::AsRawHandle;
422 use windows_sys::Win32::Storage::FileSystem::GetFinalPathNameByHandleW;
423
424 const FILE_SHARE_READ: u32 = 0x1;
425 const FILE_SHARE_WRITE: u32 = 0x2;
426 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
427 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
428 const FINAL_PATH_FLAGS: u32 = 0x0;
429
430 fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
431 let dir = OpenOptions::new()
432 .read(true)
433 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
436 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
437 .open(path)?;
438 let file_type = dir.metadata()?.file_type();
439 if file_type.is_symlink() || !file_type.is_dir() {
440 return Err(std::io::Error::new(
441 std::io::ErrorKind::InvalidInput,
442 format!(
443 "refusing to read a blob through non-directory or reparse-point component: {}",
444 path.display()
445 ),
446 ));
447 }
448 Ok(dir)
449 }
450
451 fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<PathBuf> {
452 let mut buf: Vec<u16> = vec![0; 512];
453 loop {
454 let len = unsafe {
455 GetFinalPathNameByHandleW(
456 file.as_raw_handle() as _,
457 buf.as_mut_ptr(),
458 buf.len() as u32,
459 FINAL_PATH_FLAGS,
460 )
461 };
462 if len == 0 {
463 return Err(std::io::Error::last_os_error());
464 }
465 let len = len as usize;
466 if len <= buf.len() {
467 buf.truncate(len);
468 return Ok(PathBuf::from(std::ffi::OsString::from_wide(&buf)));
469 }
470 buf.resize(len, 0);
471 }
472 }
473
474 verify_blob_root_identity(root, root_handle)?;
475 let hex = content_ref.as_str();
476 let shard1 = root.join(&hex[0..2]);
477 let shard2 = shard1.join(&hex[2..4]);
478 let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
479 let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
480 let expected = final_path_by_handle(root_handle)?
481 .join(&hex[0..2])
482 .join(&hex[2..4])
483 .join(hex);
484
485 let target = OpenOptions::new()
486 .read(true)
487 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
488 .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
489 .open(shard2.join(hex))?;
490 let file_type = target.metadata()?.file_type();
491 if file_type.is_symlink() || !file_type.is_file() {
492 return Err(std::io::Error::new(
493 std::io::ErrorKind::InvalidInput,
494 "refusing to read a blob leaf that is not a regular file",
495 ));
496 }
497 let resolved = final_path_by_handle(&target)?;
498 if resolved != expected {
499 return Err(std::io::Error::new(
500 std::io::ErrorKind::InvalidInput,
501 format!(
502 "refusing blob read: handle resolved outside the verified blob root (expected {}, resolved {})",
503 expected.display(),
504 resolved.display()
505 ),
506 ));
507 }
508 Ok(target)
509}
510
511#[cfg(not(any(unix, windows)))]
512fn open_blob_shard_file_no_follow(
513 root: &Path,
514 root_handle: &std::fs::File,
515 _content_ref: &ContentRef,
516) -> std::io::Result<std::fs::File> {
517 verify_blob_root_identity(root, root_handle)?;
518 Err(std::io::Error::new(
519 std::io::ErrorKind::Unsupported,
520 "bounded verified blob reads require handle-relative no-follow file APIs",
521 ))
522}
523
524#[cfg(unix)]
525fn open_dir_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
526 use std::os::unix::ffi::OsStrExt;
527 use std::os::unix::io::FromRawFd;
528
529 let c_path = std::ffi::CString::new(path.as_os_str().as_bytes())
530 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
531 let fd = unsafe {
534 libc::open(
535 c_path.as_ptr(),
536 libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
537 )
538 };
539 if fd < 0 {
540 return Err(std::io::Error::last_os_error());
541 }
542 Ok(unsafe { std::fs::File::from_raw_fd(fd) })
545}
546
547#[cfg(unix)]
548fn openat_dir_no_follow(
549 parent_fd: std::os::unix::io::RawFd,
550 name: &str,
551) -> std::io::Result<std::fs::File> {
552 use std::os::unix::io::FromRawFd;
553
554 let c_name = std::ffi::CString::new(name)
555 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
556 let fd = unsafe {
560 libc::openat(
561 parent_fd,
562 c_name.as_ptr(),
563 libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
564 )
565 };
566 if fd < 0 {
567 return Err(std::io::Error::last_os_error());
568 }
569 Ok(unsafe { std::fs::File::from_raw_fd(fd) })
572}
573
574#[cfg(unix)]
575fn openat_regular_file_no_follow(
576 parent_fd: std::os::unix::io::RawFd,
577 name: &str,
578 access_flags: libc::c_int,
579) -> std::io::Result<std::fs::File> {
580 use std::os::unix::io::FromRawFd;
581
582 let c_name = std::ffi::CString::new(name)
583 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
584 let fd = unsafe {
588 libc::openat(
589 parent_fd,
590 c_name.as_ptr(),
591 access_flags | libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK,
592 )
593 };
594 if fd < 0 {
595 return Err(std::io::Error::last_os_error());
596 }
597 let file = unsafe { std::fs::File::from_raw_fd(fd) };
599 if !file.metadata()?.file_type().is_file() {
600 return Err(std::io::Error::new(
601 std::io::ErrorKind::InvalidInput,
602 format!("refusing non-regular blob store entry: {name}"),
603 ));
604 }
605 Ok(file)
606}
607
608#[cfg(unix)]
609fn open_blob_shard_file_at_no_follow(
610 root_handle: &std::fs::File,
611 content_ref: &ContentRef,
612 access_flags: libc::c_int,
613) -> std::io::Result<std::fs::File> {
614 use std::os::unix::io::AsRawFd;
615
616 let hex = content_ref.as_str();
617 let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
618 let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
619 openat_regular_file_no_follow(shard2_dir.as_raw_fd(), hex, access_flags)
620}
621
622#[cfg(unix)]
623fn open_or_create_dir_at_no_follow(
624 parent_fd: std::os::unix::io::RawFd,
625 name: &str,
626) -> std::io::Result<std::fs::File> {
627 let c_name = std::ffi::CString::new(name)
628 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
629 match openat_dir_no_follow(parent_fd, name) {
630 Ok(dir) => Ok(dir),
631 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
632 let rc = unsafe { libc::mkdirat(parent_fd, c_name.as_ptr(), 0o777) };
634 if rc != 0 {
635 let error = std::io::Error::last_os_error();
636 if error.kind() != std::io::ErrorKind::AlreadyExists {
637 return Err(error);
638 }
639 }
640 openat_dir_no_follow(parent_fd, name)
644 }
645 Err(error) => Err(error),
646 }
647}
648
649#[cfg(unix)]
650fn create_regular_file_at_no_follow(
651 parent_fd: std::os::unix::io::RawFd,
652 name: &str,
653 mode: libc::mode_t,
654) -> std::io::Result<std::fs::File> {
655 use std::os::unix::io::FromRawFd;
656
657 let c_name = std::ffi::CString::new(name)
658 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
659 let fd = unsafe {
662 libc::openat(
663 parent_fd,
664 c_name.as_ptr(),
665 libc::O_RDWR | libc::O_CREAT | libc::O_EXCL | libc::O_NOFOLLOW | libc::O_CLOEXEC,
666 mode as libc::c_uint,
667 )
668 };
669 if fd < 0 {
670 return Err(std::io::Error::last_os_error());
671 }
672 Ok(unsafe { std::fs::File::from_raw_fd(fd) })
674}
675
676#[cfg(unix)]
677fn unlink_entry_at(parent_fd: std::os::unix::io::RawFd, name: &str) -> std::io::Result<()> {
678 let c_name = std::ffi::CString::new(name)
679 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
680 let rc = unsafe { libc::unlinkat(parent_fd, c_name.as_ptr(), 0) };
682 if rc != 0 {
683 return Err(std::io::Error::last_os_error());
684 }
685 Ok(())
686}
687
688#[cfg(unix)]
689fn rename_entry_at(
690 parent_fd: std::os::unix::io::RawFd,
691 from: &str,
692 to: &str,
693) -> std::io::Result<()> {
694 let c_from = std::ffi::CString::new(from)
695 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
696 let c_to = std::ffi::CString::new(to)
697 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
698 let rc = unsafe { libc::renameat(parent_fd, c_from.as_ptr(), parent_fd, c_to.as_ptr()) };
701 if rc != 0 {
702 return Err(std::io::Error::last_os_error());
703 }
704 Ok(())
705}
706
707#[cfg(unix)]
708fn available_space_at(root_handle: &std::fs::File) -> std::io::Result<u64> {
709 use std::os::unix::io::AsRawFd;
710
711 let mut stat: libc::statvfs = unsafe { std::mem::zeroed() };
712 let rc = unsafe { libc::fstatvfs(root_handle.as_raw_fd(), &mut stat) };
714 if rc != 0 {
715 return Err(std::io::Error::last_os_error());
716 }
717 #[allow(clippy::useless_conversion)]
721 Ok(stat.f_frsize.saturating_mul(u64::from(stat.f_bavail)))
722}
723
724#[cfg(unix)]
725fn acquire_root_write_lock_at(root_handle: &std::fs::File) -> StorageResult<std::fs::File> {
726 use std::os::unix::io::{AsRawFd, FromRawFd};
727
728 let c_name = std::ffi::CString::new(ROOT_WRITE_LOCK_FILE)
729 .expect("the static blob root lock name contains no NUL");
730 let fd = unsafe {
733 libc::openat(
734 root_handle.as_raw_fd(),
735 c_name.as_ptr(),
736 libc::O_RDWR | libc::O_CREAT | libc::O_NOFOLLOW | libc::O_CLOEXEC,
737 0o666,
738 )
739 };
740 if fd < 0 {
741 return Err(map_io_err(
742 std::io::Error::last_os_error(),
743 "root_write_lock_open",
744 ));
745 }
746 let lock_file = unsafe { std::fs::File::from_raw_fd(fd) };
748 if !lock_file
749 .metadata()
750 .map_err(|e| map_io_err(e, "root_write_lock_metadata"))?
751 .file_type()
752 .is_file()
753 {
754 return Err(map_io_err(
755 std::io::Error::new(
756 std::io::ErrorKind::InvalidInput,
757 "blob root write lock is not a regular file",
758 ),
759 "root_write_lock_open",
760 ));
761 }
762 fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
763 Ok(lock_file)
764}
765
766fn blob_root_key(root: &Path) -> String {
774 #[cfg(unix)]
775 let bytes = {
776 use std::os::unix::ffi::OsStrExt;
777 root.as_os_str().as_bytes().to_vec()
778 };
779 #[cfg(windows)]
780 let bytes = {
781 use std::os::windows::ffi::OsStrExt;
782 root.as_os_str()
783 .encode_wide()
784 .flat_map(u16::to_le_bytes)
785 .collect::<Vec<_>>()
786 };
787 #[cfg(not(any(unix, windows)))]
788 let bytes = root.to_string_lossy().as_bytes().to_vec();
789 blake3::hash(&bytes).to_hex().to_string()
790}
791
792pub fn resolve_blob_root(
801 db_dir: Option<&Path>,
802 config_root: Option<&Path>,
803) -> Result<PathBuf, SqliteError> {
804 if let Ok(env_root) = std::env::var("KHIVE_BLOB_ROOT") {
805 if !env_root.trim().is_empty() {
806 return Ok(PathBuf::from(env_root));
807 }
808 }
809 if let Some(root) = config_root {
810 return Ok(root.to_path_buf());
811 }
812 if let Some(dir) = db_dir {
813 return Ok(dir.join("blobs"));
814 }
815 Err(SqliteError::InvalidData(
816 "cannot resolve a blob store root: no KHIVE_BLOB_ROOT env var, no configured \
817 root, and the database has no on-disk directory to default beside (in-memory \
818 backend)"
819 .to_string(),
820 ))
821}
822
823fn crosses_floor(available: u64, required_write_bytes: u64, floor_bytes: u64) -> bool {
837 available.saturating_sub(required_write_bytes) < floor_bytes
838}
839
840#[cfg(any(test, not(unix)))]
841fn put_blocking_with_space_probe<F>(
842 root: &Path,
843 floor_bytes: u64,
844 bytes: Vec<u8>,
845 available_space: F,
846) -> StorageResult<ContentRef>
847where
848 F: FnOnce(&Path) -> std::io::Result<u64>,
849{
850 let digest = blake3::hash(&bytes);
851 let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
852 let target = shard_path(root, &content_ref);
853
854 if target.exists() {
865 let file = fs::OpenOptions::new()
866 .write(true)
867 .open(&target)
868 .map_err(|e| map_io_err(e, "put_touch_open"))?;
869 file.set_modified(SystemTime::now())
870 .map_err(|e| map_io_err(e, "put_touch_mtime"))?;
871 return Ok(content_ref);
872 }
873
874 let required_write_bytes = bytes.len() as u64;
875 let available = available_space(root).map_err(|e| map_io_err(e, "put_check_space"))?;
876 if crosses_floor(available, required_write_bytes, floor_bytes) {
877 return Err(StorageError::CapacityFloor {
878 capability: StorageCapability::Blob,
879 volume: root.display().to_string(),
880 available_bytes: available,
881 floor_bytes,
882 });
883 }
884
885 let shard_dir = target
886 .parent()
887 .expect("shard_path always nests under two directory levels");
888 fs::create_dir_all(shard_dir).map_err(|e| map_io_err(e, "put_mkdir"))?;
889
890 let mut tmp = tempfile::Builder::new()
891 .prefix(".tmp-")
892 .tempfile_in(shard_dir)
893 .map_err(|e| map_io_err(e, "put_tempfile"))?;
894 tmp.write_all(&bytes)
895 .map_err(|e| map_io_err(e, "put_write"))?;
896 tmp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
897 tmp.as_file()
898 .sync_all()
899 .map_err(|e| map_io_err(e, "put_fsync"))?;
900
901 let written_len = tmp
902 .as_file()
903 .metadata()
904 .map_err(|e| map_io_err(e, "put_verify"))?
905 .len();
906 if written_len != bytes.len() as u64 {
907 return Err(map_io_err(
908 std::io::Error::other(format!(
909 "temp file length {written_len} does not match {} written bytes",
910 bytes.len()
911 )),
912 "put_verify",
913 ));
914 }
915
916 tmp.persist(&target)
917 .map_err(|e| map_io_err(e.error, "put_persist"))?;
918
919 Ok(content_ref)
920}
921
922#[cfg(any(test, not(unix)))]
923fn put_blocking(root: &Path, floor_bytes: u64, bytes: Vec<u8>) -> StorageResult<ContentRef> {
924 let _root_write_guard = acquire_root_write_lock(root)?;
925 put_blocking_with_space_probe(root, floor_bytes, bytes, |path| fs4::available_space(path))
926}
927
928fn acquire_root_write_lock_anchored(
929 root: &Path,
930 root_handle: &std::fs::File,
931) -> StorageResult<std::fs::File> {
932 verify_blob_root_identity(root, root_handle)
933 .map_err(|e| map_io_err(e, "root_write_lock_identity"))?;
934 #[cfg(unix)]
935 {
936 acquire_root_write_lock_at(root_handle)
937 }
938 #[cfg(not(unix))]
939 {
940 acquire_root_write_lock(root)
941 }
942}
943
944#[cfg(unix)]
945fn put_blocking_from_root_handle(
946 root: &Path,
947 root_handle: &std::fs::File,
948 floor_bytes: u64,
949 bytes: Vec<u8>,
950) -> StorageResult<ContentRef> {
951 use std::os::unix::io::AsRawFd;
952
953 let _root_write_guard = acquire_root_write_lock_anchored(root, root_handle)?;
954 let digest = blake3::hash(&bytes);
955 let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
956
957 match open_blob_shard_file_at_no_follow(root_handle, &content_ref, libc::O_WRONLY) {
962 Ok(file) => {
963 file.set_modified(SystemTime::now())
964 .map_err(|e| map_io_err(e, "put_touch_mtime"))?;
965 return Ok(content_ref);
966 }
967 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
968 Err(error) => return Err(map_io_err(error, "put_touch_open")),
969 }
970
971 let required_write_bytes = bytes.len() as u64;
972 let available =
973 available_space_at(root_handle).map_err(|e| map_io_err(e, "put_check_space"))?;
974 if crosses_floor(available, required_write_bytes, floor_bytes) {
975 return Err(StorageError::CapacityFloor {
976 capability: StorageCapability::Blob,
977 volume: root.display().to_string(),
978 available_bytes: available,
979 floor_bytes,
980 });
981 }
982
983 let hex = content_ref.as_str();
984 let shard1_dir = open_or_create_dir_at_no_follow(root_handle.as_raw_fd(), &hex[0..2])
985 .map_err(|e| map_io_err(e, "put_mkdir"))?;
986 let shard2_dir = open_or_create_dir_at_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])
987 .map_err(|e| map_io_err(e, "put_mkdir"))?;
988
989 let temp_name = format!(".tmp-{}", Uuid::new_v4());
990 let mut temp = create_regular_file_at_no_follow(shard2_dir.as_raw_fd(), &temp_name, 0o600)
991 .map_err(|e| map_io_err(e, "put_tempfile"))?;
992 let write_result = (|| -> StorageResult<()> {
993 temp.write_all(&bytes)
994 .map_err(|e| map_io_err(e, "put_write"))?;
995 temp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
996 temp.sync_all().map_err(|e| map_io_err(e, "put_fsync"))?;
997
998 let written_len = temp
999 .metadata()
1000 .map_err(|e| map_io_err(e, "put_verify"))?
1001 .len();
1002 if written_len != bytes.len() as u64 {
1003 return Err(map_io_err(
1004 std::io::Error::other(format!(
1005 "temp file length {written_len} does not match {} written bytes",
1006 bytes.len()
1007 )),
1008 "put_verify",
1009 ));
1010 }
1011 Ok(())
1012 })();
1013 drop(temp);
1014 if let Err(error) = write_result {
1015 let _ = unlink_entry_at(shard2_dir.as_raw_fd(), &temp_name);
1016 return Err(error);
1017 }
1018
1019 if let Err(error) = rename_entry_at(shard2_dir.as_raw_fd(), &temp_name, hex) {
1020 let _ = unlink_entry_at(shard2_dir.as_raw_fd(), &temp_name);
1021 return Err(map_io_err(error, "put_persist"));
1022 }
1023 Ok(content_ref)
1024}
1025
1026#[cfg(not(unix))]
1027fn put_blocking_from_root_handle(
1028 root: &Path,
1029 root_handle: &std::fs::File,
1030 floor_bytes: u64,
1031 bytes: Vec<u8>,
1032) -> StorageResult<ContentRef> {
1033 verify_blob_root_identity(root, root_handle).map_err(|e| map_io_err(e, "put_root_identity"))?;
1034 put_blocking(root, floor_bytes, bytes)
1035}
1036
1037#[cfg(any(test, not(unix)))]
1038fn acquire_root_write_lock(root: &Path) -> StorageResult<fs::File> {
1039 let lock_file = fs::OpenOptions::new()
1040 .read(true)
1041 .write(true)
1042 .create(true)
1043 .truncate(false)
1044 .open(root.join(ROOT_WRITE_LOCK_FILE))
1045 .map_err(|e| map_io_err(e, "root_write_lock_open"))?;
1046 fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
1047 Ok(lock_file)
1048}
1049
1050fn database_gc_lock_path(database_path: &Path) -> PathBuf {
1051 let mut lock_path = database_path.as_os_str().to_os_string();
1052 lock_path.push(DATABASE_GC_LOCK_SUFFIX);
1053 PathBuf::from(lock_path)
1054}
1055
1056fn acquire_database_gc_lock(database_path: Option<&Path>) -> StorageResult<Option<fs::File>> {
1057 let Some(database_path) = database_path else {
1058 return Ok(None);
1061 };
1062 let lock_path = database_gc_lock_path(database_path);
1063 let lock_file = fs::OpenOptions::new()
1064 .read(true)
1065 .write(true)
1066 .create(true)
1067 .truncate(false)
1068 .open(&lock_path)
1069 .map_err(|e| map_io_err(e, "database_gc_lock_open"))?;
1070 fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "database_gc_lock_acquire"))?;
1071 Ok(Some(lock_file))
1072}
1073
1074#[cfg(all(unix, target_os = "macos"))]
1075fn errno_location() -> *mut libc::c_int {
1076 unsafe { libc::__error() }
1080}
1081
1082#[cfg(all(unix, not(target_os = "macos")))]
1083fn errno_location() -> *mut libc::c_int {
1084 unsafe { libc::__errno_location() }
1087}
1088
1089#[cfg(unix)]
1093fn clear_errno() {
1094 unsafe { *errno_location() = 0 };
1097}
1098
1099#[cfg(unix)]
1100fn current_errno() -> libc::c_int {
1101 unsafe { *errno_location() }
1103}
1104
1105#[cfg(unix)]
1130fn read_dir_names_no_follow(dir_fd: std::os::unix::io::RawFd) -> std::io::Result<Vec<String>> {
1131 use std::os::unix::io::IntoRawFd;
1132
1133 let reopened = openat_dir_no_follow(dir_fd, ".")?;
1134 let owned_fd = reopened.into_raw_fd();
1138 let dirp = unsafe { libc::fdopendir(owned_fd) };
1141 if dirp.is_null() {
1142 let err = std::io::Error::last_os_error();
1143 unsafe { libc::close(owned_fd) };
1145 return Err(err);
1146 }
1147 let mut names = Vec::new();
1148 loop {
1149 clear_errno();
1152 let entry = unsafe { libc::readdir(dirp) };
1154 if entry.is_null() {
1155 if current_errno() != 0 {
1156 let err = std::io::Error::last_os_error();
1157 unsafe { libc::closedir(dirp) };
1161 return Err(err);
1162 }
1163 break;
1164 }
1165 let first = unsafe { *(*entry).d_name.as_ptr() };
1170 if first == b'.' as libc::c_char {
1171 continue;
1172 }
1173 let name = unsafe { std::ffi::CStr::from_ptr((*entry).d_name.as_ptr()) }
1176 .to_string_lossy()
1177 .into_owned();
1178 names.push(name);
1179 }
1180 unsafe { libc::closedir(dirp) };
1182 Ok(names)
1183}
1184
1185#[cfg(unix)]
1203fn walk_blob_files_from_root_handle(
1204 root_handle: &std::fs::File,
1205 _root: &Path,
1209) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
1210 use std::os::unix::io::AsRawFd;
1211
1212 let mut out = Vec::new();
1213 for l1_name in read_dir_names_no_follow(root_handle.as_raw_fd())? {
1214 let l1_dir = match openat_dir_no_follow(root_handle.as_raw_fd(), &l1_name) {
1215 Ok(dir) => dir,
1216 Err(_) => continue,
1217 };
1218 for l2_name in read_dir_names_no_follow(l1_dir.as_raw_fd())? {
1219 let l2_dir = match openat_dir_no_follow(l1_dir.as_raw_fd(), &l2_name) {
1220 Ok(dir) => dir,
1221 Err(_) => continue,
1222 };
1223 for leaf_name in read_dir_names_no_follow(l2_dir.as_raw_fd())? {
1224 let Ok(content_ref) = ContentRef::from_hex(leaf_name.clone()) else {
1227 continue;
1228 };
1229 #[cfg(all(test, unix))]
1230 if let Some(hook) = walk_leaf_sync_hook::take(_root) {
1231 let _ = hook.reached.send(());
1232 let _ = hook.release.recv();
1233 }
1234 let file = match openat_regular_file_no_follow(
1235 l2_dir.as_raw_fd(),
1236 &leaf_name,
1237 libc::O_RDONLY,
1238 ) {
1239 Ok(file) => file,
1240 Err(_) => continue,
1241 };
1242 let mtime = file.metadata().ok().and_then(|meta| meta.modified().ok());
1243 out.push((content_ref, mtime));
1244 }
1245 }
1246 }
1247 Ok(out)
1248}
1249
1250#[cfg(not(unix))]
1258fn walk_blob_files_from_root_handle(
1259 _root_handle: &std::fs::File,
1260 _root: &Path,
1261) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
1262 Err(std::io::Error::new(
1263 std::io::ErrorKind::Unsupported,
1264 "orphan sweep candidate enumeration requires descriptor-relative directory \
1265 reads, available only on unix in this release; refusing to classify via \
1266 path-based reads",
1267 ))
1268}
1269
1270fn within_publish_grace(
1282 mtime: Option<SystemTime>,
1283 now: SystemTime,
1284 grace_period: Duration,
1285) -> bool {
1286 let age = mtime.and_then(|mtime| now.duration_since(mtime).ok());
1287 match age {
1288 Some(age) => age < grace_period,
1289 None => true,
1290 }
1291}
1292
1293#[derive(Debug)]
1294struct PreparedTransactionalSweep {
1295 result: BlobOrphanSweepResult,
1296 candidates: Vec<(ContentRef, bool)>,
1297}
1298
1299fn prepare_transactional_sweep(
1302 files: Vec<(ContentRef, Option<SystemTime>)>,
1303 grace_period: Duration,
1304) -> PreparedTransactionalSweep {
1305 let now = SystemTime::now();
1306 let mut result = BlobOrphanSweepResult::default();
1307 let mut candidates = Vec::with_capacity(files.len());
1308 for (content_ref, mtime) in files {
1309 result.scanned += 1;
1310 let within_grace = within_publish_grace(mtime, now, grace_period);
1311 candidates.push((content_ref, within_grace));
1312 }
1313 PreparedTransactionalSweep { result, candidates }
1314}
1315
1316#[derive(Debug)]
1317struct BlobGcBatchRows {
1318 grace_period_skipped: u64,
1319 would_delete: u64,
1320 claimed_rows: Vec<SqlRow>,
1321}
1322
1323fn required_nonnegative_count(
1324 value: Option<SqlValue>,
1325 operation: &'static str,
1326) -> StorageResult<u64> {
1327 match value {
1328 Some(SqlValue::Integer(value)) if value >= 0 => Ok(value as u64),
1329 other => Err(StorageError::Internal(format!(
1330 "{operation} returned an invalid count: {other:?}"
1331 ))),
1332 }
1333}
1334
1335fn invalid_content_ref(message: String) -> StorageError {
1336 StorageError::InvalidInput {
1337 capability: StorageCapability::Blob,
1338 operation: "transactional_orphan_sweep".into(),
1339 message,
1340 }
1341}
1342
1343async fn blob_gc_fencing_complete(sql: &dyn SqlAccess) -> StorageResult<bool> {
1357 let mut reader = sql.reader().await?;
1358 let present = required_nonnegative_count(
1359 reader
1360 .query_scalar(SqlStatement {
1361 sql: "SELECT COUNT(*) FROM sqlite_master \
1362 WHERE (type = 'table' AND name IN ( \
1363 'blob_gc_claims', 'attachments', \
1364 'attachment_cutover_state')) \
1365 OR (type = 'index' AND name IN ( \
1366 'idx_blob_gc_claims_content_ref', \
1367 'idx_attachments_content_ref')) \
1368 OR (type = 'trigger' AND name IN ( \
1369 'attachments_reject_claimed_blob_insert', \
1370 'attachments_reject_claimed_blob_update'))"
1371 .to_string(),
1372 params: vec![],
1373 label: Some("blob_gc_fencing_complete".to_string()),
1374 })
1375 .await?,
1376 "blob_gc_fencing_complete",
1377 )?;
1378 if present != 7 {
1379 return Ok(false);
1380 }
1381
1382 let legacy_objects = required_nonnegative_count(
1383 reader
1384 .query_scalar(SqlStatement {
1385 sql: "SELECT \
1386 (SELECT COUNT(*) FROM pragma_table_info('entities') \
1387 WHERE name = 'content_ref') \
1388 + (SELECT COUNT(*) FROM sqlite_master \
1389 WHERE (type = 'index' AND name = 'idx_entities_content_ref') \
1390 OR (type = 'trigger' AND name IN ( \
1391 'entities_reject_claimed_blob_insert', \
1392 'entities_reject_claimed_blob_update')))"
1393 .to_string(),
1394 params: vec![],
1395 label: Some("blob_gc_legacy_fencing_absent".to_string()),
1396 })
1397 .await?,
1398 "blob_gc_legacy_fencing_absent",
1399 )?;
1400 if legacy_objects != 0 {
1401 return Ok(false);
1402 }
1403
1404 let complete = required_nonnegative_count(
1429 reader
1430 .query_scalar(SqlStatement {
1431 sql: "SELECT COUNT(*) FROM attachment_cutover_state AS cutover \
1432 WHERE cutover.singleton = 1 \
1433 AND cutover.state = 'complete' \
1434 AND cutover.completed_at IS NOT NULL \
1435 AND (SELECT COUNT(*) FROM _schema_migrations \
1436 WHERE version = ?1 \
1437 AND name = 'attachments_first_class') = 1 \
1438 AND (SELECT COUNT(*) FROM _schema_migrations) = ?1 \
1439 AND (SELECT MIN(version) FROM _schema_migrations) = 1 \
1440 AND (SELECT MAX(version) FROM _schema_migrations) = ?1"
1441 .to_string(),
1442 params: vec![SqlValue::Integer(i64::from(
1443 crate::migrations::ATTACHMENT_CUTOVER_VERSION,
1444 ))],
1445 label: Some("blob_gc_cutover_complete".to_string()),
1446 })
1447 .await?,
1448 "blob_gc_cutover_complete",
1449 )?;
1450 Ok(complete == 1)
1451}
1452
1453fn unsupported_blob_gc_epoch() -> StorageError {
1454 StorageError::Unsupported {
1455 capability: StorageCapability::Blob,
1456 operation: "transactional_orphan_sweep".into(),
1457 message: "transactional blob GC requires a complete V21 attachment cutover with \
1458 the attachment claim-fencing set; refusing both report-only and \
1459 destructive sweep in this database epoch"
1460 .into(),
1461 }
1462}
1463
1464const BLOB_GC_FENCE_PROBE_REF: &str =
1468 "0000000000000000000000000000000000000000000000000000000000000000";
1469
1470const BLOB_GC_FENCE_TRIGGER_MESSAGE: &str = "content_ref is reserved by an active blob sweep";
1473
1474async fn blob_gc_fence_probe(sql: &dyn SqlAccess) -> StorageResult<()> {
1489 let run = Uuid::new_v4().simple().to_string();
1490 blob_gc_fence_probe_with_ids(
1491 sql,
1492 format!("__blob-gc-fence-probe-insert-{run}__"),
1493 format!("__blob-gc-fence-probe-update-{run}__"),
1494 format!("__blob-gc-fence-probe-insert2-{run}__"),
1495 format!("__blob-gc-fence-probe-update2-{run}__"),
1496 format!("__fence_probe-{run}__"),
1497 )
1498 .await
1499}
1500
1501async fn blob_gc_fence_probe_with_ids(
1506 sql: &dyn SqlAccess,
1507 insert_id: String,
1508 update_id: String,
1509 insert2_id: String,
1510 update2_id: String,
1511 claim_key: String,
1512) -> StorageResult<()> {
1513 fn fence_rejection(result: Result<u64, StorageError>) -> Result<bool, String> {
1514 match result {
1515 Ok(_) => Ok(false),
1516 Err(error) => {
1517 let text = error.to_string();
1518 if text.contains(BLOB_GC_FENCE_TRIGGER_MESSAGE) {
1519 Ok(true)
1520 } else {
1521 Err(text)
1522 }
1523 }
1524 }
1525 }
1526
1527 fn required_seed(value: Option<SqlValue>) -> StorageResult<String> {
1528 match value {
1529 Some(SqlValue::Text(seed)) => Ok(seed),
1530 _ => Err(StorageError::Unsupported {
1531 capability: StorageCapability::Blob,
1532 operation: "transactional_orphan_sweep".into(),
1533 message: "the blob GC fence probe could not select an unclaimed \
1534 canonical seed; refusing deletion so a later sweep can retry"
1535 .into(),
1536 }),
1537 }
1538 }
1539
1540 let op: AtomicUnitOp = Box::new(move |writer| {
1541 Box::pin(async move {
1542 let preexisting = writer
1546 .query_row(SqlStatement {
1547 sql: "SELECT (SELECT COUNT(*) FROM attachments \
1548 WHERE record_uuid IN (?1, ?2, ?3, ?4)) \
1549 + (SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?5)"
1550 .to_string(),
1551 params: vec![
1552 SqlValue::Text(insert_id.clone()),
1553 SqlValue::Text(update_id.clone()),
1554 SqlValue::Text(insert2_id.clone()),
1555 SqlValue::Text(update2_id.clone()),
1556 SqlValue::Text(claim_key.clone()),
1557 ],
1558 label: Some("blob_gc_fence_probe_ownership_guard".to_string()),
1559 })
1560 .await?
1561 .and_then(|row| row.columns.first().map(|c| c.value.clone()));
1562 match preexisting {
1563 Some(SqlValue::Integer(0)) => {}
1564 Some(SqlValue::Integer(_)) => {
1565 return Err(StorageError::Unsupported {
1566 capability: StorageCapability::Blob,
1567 operation: "transactional_orphan_sweep".into(),
1568 message: "the blob GC fence probe's row ids collide with existing \
1569 rows; refusing to probe rather than delete data the \
1570 probe does not own"
1571 .into(),
1572 });
1573 }
1574 _ => {
1575 return Err(StorageError::Internal(
1576 "blob GC fence probe ownership guard returned no count".into(),
1577 ));
1578 }
1579 }
1580
1581 let seed_ref = writer
1587 .query_row(SqlStatement {
1588 sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1589 SELECT 1, lower(hex(randomblob(32))) \
1590 UNION ALL \
1591 SELECT attempt + 1, lower(hex(randomblob(32))) \
1592 FROM candidates WHERE attempt < 8 \
1593 ) \
1594 SELECT candidate.content_ref FROM candidates AS candidate \
1595 WHERE candidate.content_ref <> ?1 \
1596 AND NOT EXISTS ( \
1597 SELECT 1 FROM blob_gc_claims \
1598 WHERE content_ref = candidate.content_ref \
1599 ) \
1600 LIMIT 1"
1601 .to_string(),
1602 params: vec![SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string())],
1603 label: Some("blob_gc_fence_probe_select_seed".to_string()),
1604 })
1605 .await?
1606 .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1607 let seed_ref = required_seed(seed_ref)?;
1608
1609 let seed2_ref = writer
1610 .query_row(SqlStatement {
1611 sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1612 SELECT 1, lower(hex(randomblob(32))) \
1613 UNION ALL \
1614 SELECT attempt + 1, lower(hex(randomblob(32))) \
1615 FROM candidates WHERE attempt < 8 \
1616 ) \
1617 SELECT candidate.content_ref FROM candidates AS candidate \
1618 WHERE candidate.content_ref NOT IN (?1, ?2) \
1619 AND NOT EXISTS ( \
1620 SELECT 1 FROM blob_gc_claims \
1621 WHERE content_ref = candidate.content_ref \
1622 ) \
1623 LIMIT 1"
1624 .to_string(),
1625 params: vec![
1626 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1627 SqlValue::Text(seed_ref.clone()),
1628 ],
1629 label: Some("blob_gc_fence_probe_select_seed2".to_string()),
1630 })
1631 .await?
1632 .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1633 let seed2_ref = required_seed(seed2_ref)?;
1634
1635 let probe2_ref = writer
1639 .query_row(SqlStatement {
1640 sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1641 SELECT 1, lower(hex(randomblob(32))) \
1642 UNION ALL \
1643 SELECT attempt + 1, lower(hex(randomblob(32))) \
1644 FROM candidates WHERE attempt < 8 \
1645 ) \
1646 SELECT candidate.content_ref FROM candidates AS candidate \
1647 WHERE candidate.content_ref NOT IN (?1, ?2, ?3) \
1648 AND NOT EXISTS ( \
1649 SELECT 1 FROM blob_gc_claims \
1650 WHERE content_ref = candidate.content_ref \
1651 ) \
1652 LIMIT 1"
1653 .to_string(),
1654 params: vec![
1655 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1656 SqlValue::Text(seed_ref.clone()),
1657 SqlValue::Text(seed2_ref.clone()),
1658 ],
1659 label: Some("blob_gc_fence_probe_select_probe2".to_string()),
1660 })
1661 .await?
1662 .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1663 let probe2_ref = required_seed(probe2_ref)?;
1664
1665 writer
1666 .execute(SqlStatement {
1667 sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
1668 VALUES (?1, ?2, 0), (?1, ?3, 0)"
1669 .to_string(),
1670 params: vec![
1671 SqlValue::Text(claim_key.clone()),
1672 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1673 SqlValue::Text(probe2_ref.clone()),
1674 ],
1675 label: Some("blob_gc_fence_probe_claim".to_string()),
1676 })
1677 .await?;
1678
1679 let insert_attempt = writer
1680 .execute(SqlStatement {
1681 sql: "INSERT INTO attachments \
1682 (record_uuid, substrate, role, content_ref, created_at) \
1683 VALUES (?1, 'entity', 'content', ?2, 0)"
1684 .to_string(),
1685 params: vec![
1686 SqlValue::Text(insert_id.clone()),
1687 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1688 ],
1689 label: Some("blob_gc_fence_probe_insert_arm".to_string()),
1690 })
1691 .await;
1692 let insert_fenced = fence_rejection(insert_attempt);
1693
1694 writer
1695 .execute(SqlStatement {
1696 sql: "INSERT INTO attachments \
1697 (record_uuid, substrate, role, content_ref, created_at) \
1698 VALUES (?1, 'entity', 'content', ?2, 0)"
1699 .to_string(),
1700 params: vec![SqlValue::Text(update_id.clone()), SqlValue::Text(seed_ref)],
1701 label: Some("blob_gc_fence_probe_update_arm_seed".to_string()),
1702 })
1703 .await?;
1704 let update_attempt = writer
1705 .execute(SqlStatement {
1706 sql: "UPDATE attachments SET content_ref = ?1 \
1707 WHERE record_uuid = ?2 AND role = 'content'"
1708 .to_string(),
1709 params: vec![
1710 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1711 SqlValue::Text(update_id.clone()),
1712 ],
1713 label: Some("blob_gc_fence_probe_update_arm".to_string()),
1714 })
1715 .await;
1716 let update_fenced = fence_rejection(update_attempt);
1717
1718 let insert2_attempt = writer
1719 .execute(SqlStatement {
1720 sql: "INSERT INTO attachments \
1721 (record_uuid, substrate, role, content_ref, created_at) \
1722 VALUES (?1, 'note', 'evidence', ?2, 0)"
1723 .to_string(),
1724 params: vec![
1725 SqlValue::Text(insert2_id.clone()),
1726 SqlValue::Text(probe2_ref.clone()),
1727 ],
1728 label: Some("blob_gc_fence_probe_insert2_arm".to_string()),
1729 })
1730 .await;
1731 let insert2_fenced = fence_rejection(insert2_attempt);
1732
1733 writer
1734 .execute(SqlStatement {
1735 sql: "INSERT INTO attachments \
1736 (record_uuid, substrate, role, content_ref, created_at) \
1737 VALUES (?1, 'note', 'evidence', ?2, 0)"
1738 .to_string(),
1739 params: vec![
1740 SqlValue::Text(update2_id.clone()),
1741 SqlValue::Text(seed2_ref),
1742 ],
1743 label: Some("blob_gc_fence_probe_update2_arm_seed".to_string()),
1744 })
1745 .await?;
1746 let update2_attempt = writer
1747 .execute(SqlStatement {
1748 sql: "UPDATE attachments SET content_ref = ?1 \
1749 WHERE record_uuid = ?2 AND role = 'evidence'"
1750 .to_string(),
1751 params: vec![
1752 SqlValue::Text(probe2_ref),
1753 SqlValue::Text(update2_id.clone()),
1754 ],
1755 label: Some("blob_gc_fence_probe_update2_arm".to_string()),
1756 })
1757 .await;
1758 let update2_fenced = fence_rejection(update2_attempt);
1759
1760 writer
1763 .execute(SqlStatement {
1764 sql: "DELETE FROM attachments WHERE record_uuid IN (?1, ?2, ?3, ?4)"
1765 .to_string(),
1766 params: vec![
1767 SqlValue::Text(insert_id.clone()),
1768 SqlValue::Text(update_id.clone()),
1769 SqlValue::Text(insert2_id.clone()),
1770 SqlValue::Text(update2_id.clone()),
1771 ],
1772 label: Some("blob_gc_fence_probe_cleanup_attachments".to_string()),
1773 })
1774 .await?;
1775 writer
1776 .execute(SqlStatement {
1777 sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
1778 params: vec![SqlValue::Text(claim_key)],
1779 label: Some("blob_gc_fence_probe_cleanup_claim".to_string()),
1780 })
1781 .await?;
1782
1783 Ok(
1784 Box::new((insert_fenced, update_fenced, insert2_fenced, update2_fenced))
1785 as Box<dyn std::any::Any + Send>,
1786 )
1787 })
1788 });
1789 let outcome = sql.atomic_unit(op).await?;
1790 let (insert_fenced, update_fenced, insert2_fenced, update2_fenced) = *outcome
1791 .downcast::<(
1792 Result<bool, String>,
1793 Result<bool, String>,
1794 Result<bool, String>,
1795 Result<bool, String>,
1796 )>()
1797 .map_err(|_| {
1798 StorageError::Internal("blob GC fence probe returned an unexpected outcome type".into())
1799 })?;
1800 let arm_verdict = |arm: &str, fenced: Result<bool, String>| -> StorageResult<()> {
1801 match fenced {
1802 Ok(true) => Ok(()),
1803 Ok(false) => Err(StorageError::Unsupported {
1804 capability: StorageCapability::Blob,
1805 operation: "transactional_orphan_sweep".into(),
1806 message: format!(
1807 "the V21 fencing triggers exist by name but did not reject a claimed \
1808 content_ref on the attachment {arm} path; refusing unfenced deletion"
1809 ),
1810 }),
1811 Err(other) => Err(StorageError::Unsupported {
1812 capability: StorageCapability::Blob,
1813 operation: "transactional_orphan_sweep".into(),
1814 message: format!(
1815 "the blob GC fence probe could not verify the attachment {arm} fence \
1816 (unexpected rejection: {other}); refusing unfenced deletion"
1817 ),
1818 }),
1819 }
1820 };
1821 arm_verdict("INSERT", insert_fenced)?;
1822 arm_verdict("UPDATE", update_fenced)?;
1823 arm_verdict("second-digest INSERT", insert2_fenced)?;
1824 arm_verdict("second-digest UPDATE", update2_fenced)
1825}
1826
1827async fn validate_blob_gc_evidence(sql: &dyn SqlAccess) -> StorageResult<()> {
1828 let mut reader = sql.reader().await?;
1833 let canonical_bytes = match reader
1842 .query_row(SqlStatement {
1843 sql: "SELECT length(CAST('x' AS BLOB))".to_string(),
1844 params: vec![],
1845 label: Some("blob_gc_validate_encoding_width".to_string()),
1846 })
1847 .await?
1848 .and_then(|row| row.columns.first().map(|column| column.value.clone()))
1849 {
1850 Some(SqlValue::Integer(width)) if (1..=4).contains(&width) => width * 64,
1851 other => {
1852 return Err(invalid_content_ref(format!(
1853 "the text-encoding width probe returned {other:?}; refusing GC validation"
1854 )));
1855 }
1856 };
1857 let invalid_claim = reader
1858 .query_row(SqlStatement {
1859 sql: "SELECT content_ref FROM blob_gc_claims \
1860 WHERE typeof(content_ref) <> 'text' \
1861 OR length(content_ref) <> 64 \
1862 OR length(CAST(content_ref AS BLOB)) <> ?1 \
1863 OR content_ref GLOB '*[^0-9a-f]*' \
1864 LIMIT 1"
1865 .to_string(),
1866 params: vec![SqlValue::Integer(canonical_bytes)],
1867 label: Some("blob_gc_validate_existing_claims".to_string()),
1868 })
1869 .await?;
1870 if invalid_claim.is_some() {
1871 return Err(invalid_content_ref(
1872 "blob_gc_claims.content_ref contained a non-canonical value".into(),
1873 ));
1874 }
1875
1876 let invalid_live = reader
1877 .query_row(SqlStatement {
1878 sql: "SELECT content_ref FROM attachments \
1879 WHERE typeof(content_ref) <> 'text' \
1880 OR length(content_ref) <> 64 \
1881 OR length(CAST(content_ref AS BLOB)) <> ?1 \
1882 OR content_ref GLOB '*[^0-9a-f]*' \
1883 LIMIT 1"
1884 .to_string(),
1885 params: vec![SqlValue::Integer(canonical_bytes)],
1886 label: Some("blob_gc_validate_live_refs".to_string()),
1887 })
1888 .await?;
1889 if invalid_live.is_some() {
1890 return Err(invalid_content_ref(
1891 "attachments.content_ref contained a non-canonical value".into(),
1892 ));
1893 }
1894 Ok(())
1895}
1896
1897async fn release_abandoned_blob_gc_claim_batch(sql: &dyn SqlAccess) -> StorageResult<u64> {
1898 let op: AtomicUnitOp = Box::new(move |writer| {
1899 Box::pin(async move {
1900 let released = writer
1901 .execute(SqlStatement {
1902 sql: "DELETE FROM blob_gc_claims \
1903 WHERE rowid IN ( \
1904 SELECT rowid FROM blob_gc_claims \
1905 ORDER BY rowid LIMIT ?1 \
1906 )"
1907 .to_string(),
1908 params: vec![SqlValue::Integer(BLOB_GC_CLAIM_BATCH_SIZE as i64)],
1909 label: Some("blob_gc_release_abandoned_claim_batch".to_string()),
1910 })
1911 .await?;
1912 Ok(Box::new(released) as Box<dyn std::any::Any + Send>)
1913 })
1914 });
1915 let released = sql.atomic_unit(op).await?;
1916 released.downcast::<u64>().map(|count| *count).map_err(|_| {
1917 StorageError::Internal(
1918 "transactional orphan sweep returned an unexpected recovery count type".into(),
1919 )
1920 })
1921}
1922
1923async fn claim_blob_gc_batch(
1924 sql: &dyn SqlAccess,
1925 root_key: String,
1926 candidates: &[(ContentRef, bool)],
1927 dry_run: bool,
1928) -> StorageResult<BlobGcBatchRows> {
1929 debug_assert!(candidates.len() <= BLOB_GC_CLAIM_BATCH_SIZE);
1930 let eligible_refs = candidates
1931 .iter()
1932 .filter(|(_, within_grace)| !within_grace)
1933 .map(|(content_ref, _)| content_ref.to_string())
1934 .collect::<Vec<_>>();
1935 let grace_refs = candidates
1936 .iter()
1937 .filter(|(_, within_grace)| *within_grace)
1938 .map(|(content_ref, _)| content_ref.to_string())
1939 .collect::<Vec<_>>();
1940 let eligible_json = serde_json::to_string(&eligible_refs).map_err(|error| {
1941 StorageError::Internal(format!(
1942 "failed to prepare blob GC eligible candidate batch: {error}"
1943 ))
1944 })?;
1945 let grace_json = serde_json::to_string(&grace_refs).map_err(|error| {
1946 StorageError::Internal(format!(
1947 "failed to prepare blob GC grace candidate batch: {error}"
1948 ))
1949 })?;
1950 let claimed_at = chrono::Utc::now().timestamp_micros();
1951 let op: AtomicUnitOp = Box::new(move |writer| {
1952 Box::pin(async move {
1953 let grace_period_skipped = required_nonnegative_count(
1954 writer
1955 .query_scalar(SqlStatement {
1956 sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
1957 WHERE NOT EXISTS ( \
1958 SELECT 1 FROM attachments \
1959 WHERE content_ref = candidate.value \
1960 )"
1961 .to_string(),
1962 params: vec![SqlValue::Text(grace_json)],
1963 label: Some("blob_gc_count_grace_candidates_batch".to_string()),
1964 })
1965 .await?,
1966 "blob_gc_count_grace_candidates_batch",
1967 )?;
1968
1969 if dry_run {
1970 let would_delete = required_nonnegative_count(
1971 writer
1972 .query_scalar(SqlStatement {
1973 sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
1974 WHERE NOT EXISTS ( \
1975 SELECT 1 FROM attachments \
1976 WHERE content_ref = candidate.value \
1977 )"
1978 .to_string(),
1979 params: vec![SqlValue::Text(eligible_json)],
1980 label: Some("blob_gc_count_dry_run_candidates_batch".to_string()),
1981 })
1982 .await?,
1983 "blob_gc_count_dry_run_candidates_batch",
1984 )?;
1985 return Ok(Box::new(BlobGcBatchRows {
1986 grace_period_skipped,
1987 would_delete,
1988 claimed_rows: Vec::new(),
1989 }) as Box<dyn std::any::Any + Send>);
1990 }
1991
1992 writer
1993 .execute(SqlStatement {
1994 sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
1995 SELECT ?1, candidate.value, ?3 \
1996 FROM json_each(?2) AS candidate \
1997 WHERE NOT EXISTS ( \
1998 SELECT 1 FROM attachments \
1999 WHERE content_ref = candidate.value \
2000 )"
2001 .to_string(),
2002 params: vec![
2003 SqlValue::Text(root_key.clone()),
2004 SqlValue::Text(eligible_json),
2005 SqlValue::Integer(claimed_at),
2006 ],
2007 label: Some("blob_gc_claim_candidate_batch".to_string()),
2008 })
2009 .await?;
2010
2011 let claimed_rows = writer
2012 .query_all(SqlStatement {
2013 sql: "SELECT content_ref FROM blob_gc_claims \
2014 WHERE root_key = ?1 ORDER BY content_ref"
2015 .to_string(),
2016 params: vec![SqlValue::Text(root_key)],
2017 label: Some("blob_gc_claimed_candidate_batch".to_string()),
2018 })
2019 .await?;
2020 Ok(Box::new(BlobGcBatchRows {
2021 grace_period_skipped,
2022 would_delete: claimed_rows.len() as u64,
2023 claimed_rows,
2024 }) as Box<dyn std::any::Any + Send>)
2025 })
2026 });
2027 let rows = sql.atomic_unit(op).await?;
2028 rows.downcast::<BlobGcBatchRows>()
2029 .map(|rows| *rows)
2030 .map_err(|_| {
2031 StorageError::Internal(
2032 "transactional orphan sweep returned an unexpected batch-row type".into(),
2033 )
2034 })
2035}
2036
2037fn parse_blob_gc_claim_rows(rows: Vec<SqlRow>) -> StorageResult<Vec<ContentRef>> {
2038 let mut claimed = Vec::with_capacity(rows.len());
2039 for row in rows {
2040 let raw = match row.get("content_ref") {
2041 Some(SqlValue::Text(raw)) => raw.clone(),
2042 _ => {
2043 return Err(invalid_content_ref(
2044 "blob_gc_claims.content_ref contained a non-text value".into(),
2045 ));
2046 }
2047 };
2048 claimed.push(ContentRef::from_hex(raw).map_err(invalid_content_ref)?);
2049 }
2050 Ok(claimed)
2051}
2052
2053async fn release_blob_gc_batch(sql: &dyn SqlAccess, root_key: String) -> StorageResult<()> {
2054 let cleanup: AtomicUnitOp = Box::new(move |writer| {
2055 Box::pin(async move {
2056 writer
2057 .execute(SqlStatement {
2058 sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
2059 params: vec![SqlValue::Text(root_key)],
2060 label: Some("blob_gc_release_claim_batch".to_string()),
2061 })
2062 .await?;
2063 Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
2064 })
2065 });
2066 sql.atomic_unit(cleanup).await?;
2067 Ok(())
2068}
2069
2070type SweepLockMap = HashMap<Option<PathBuf>, Arc<DatabaseGcProcessLock>>;
2077
2078#[derive(Debug, Default)]
2079struct DatabaseGcProcessLock {
2080 held: StdMutex<bool>,
2081 released: std::sync::Condvar,
2082 #[cfg(test)]
2083 waiters: std::sync::atomic::AtomicUsize,
2084}
2085
2086impl DatabaseGcProcessLock {
2087 fn acquire(self: &Arc<Self>) -> DatabaseGcProcessGuard {
2088 let mut held = self
2089 .held
2090 .lock()
2091 .unwrap_or_else(std::sync::PoisonError::into_inner);
2092 while *held {
2093 #[cfg(test)]
2094 self.waiters
2095 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2096 held = self
2097 .released
2098 .wait(held)
2099 .unwrap_or_else(std::sync::PoisonError::into_inner);
2100 #[cfg(test)]
2101 self.waiters
2102 .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
2103 }
2104 *held = true;
2105 DatabaseGcProcessGuard {
2106 lock: Arc::clone(self),
2107 }
2108 }
2109
2110 fn try_acquire(self: &Arc<Self>) -> Option<DatabaseGcProcessGuard> {
2111 let mut held = self
2112 .held
2113 .lock()
2114 .unwrap_or_else(std::sync::PoisonError::into_inner);
2115 if *held {
2116 return None;
2117 }
2118 *held = true;
2119 Some(DatabaseGcProcessGuard {
2120 lock: Arc::clone(self),
2121 })
2122 }
2123}
2124
2125#[derive(Debug)]
2126struct DatabaseGcProcessGuard {
2127 lock: Arc<DatabaseGcProcessLock>,
2128}
2129
2130impl Drop for DatabaseGcProcessGuard {
2131 fn drop(&mut self) {
2132 let mut held = self
2133 .lock
2134 .held
2135 .lock()
2136 .unwrap_or_else(std::sync::PoisonError::into_inner);
2137 debug_assert!(*held, "database GC process owner released twice");
2138 *held = false;
2139 self.lock.released.notify_one();
2140 }
2141}
2142
2143fn database_sweep_locks() -> &'static StdMutex<SweepLockMap> {
2144 static REGISTRY: OnceLock<StdMutex<SweepLockMap>> = OnceLock::new();
2145 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2146}
2147
2148fn sweep_lock_for_database(database_path: Option<&Path>) -> Arc<DatabaseGcProcessLock> {
2149 let key = database_path.map(Path::to_path_buf);
2150 let mut locks = database_sweep_locks()
2151 .lock()
2152 .unwrap_or_else(std::sync::PoisonError::into_inner);
2153 locks
2154 .entry(key)
2155 .or_insert_with(|| Arc::new(DatabaseGcProcessLock::default()))
2156 .clone()
2157}
2158
2159#[cfg(test)]
2160pub(crate) fn database_gc_waiter_count(database_path: Option<&Path>) -> usize {
2161 sweep_lock_for_database(database_path)
2162 .waiters
2163 .load(std::sync::atomic::Ordering::SeqCst)
2164}
2165
2166pub struct DatabaseGcOwnerGuard {
2174 _process_guard: DatabaseGcProcessGuard,
2175 _advisory_guard: Option<fs::File>,
2176 database_path: Option<PathBuf>,
2177}
2178
2179pub(crate) fn acquire_database_gc_owner_for_path_blocking(
2180 database_path: Option<PathBuf>,
2181) -> StorageResult<DatabaseGcOwnerGuard> {
2182 let process_guard = sweep_lock_for_database(database_path.as_deref()).acquire();
2183 let advisory_guard = acquire_database_gc_lock(database_path.as_deref())?;
2184 Ok(DatabaseGcOwnerGuard {
2185 _process_guard: process_guard,
2186 _advisory_guard: advisory_guard,
2187 database_path,
2188 })
2189}
2190
2191pub(crate) fn try_acquire_database_gc_owner_for_path(
2199 database_path: PathBuf,
2200) -> StorageResult<DatabaseGcOwnerGuard> {
2201 let process_guard = sweep_lock_for_database(Some(&database_path))
2202 .try_acquire()
2203 .ok_or_else(|| {
2204 StorageError::Internal(format!(
2205 "database GC owner for {} is already held; retry schema migration through the \
2206 coordinated backend boot path",
2207 database_path.display()
2208 ))
2209 })?;
2210 let lock_path = database_gc_lock_path(&database_path);
2211 let advisory_guard = fs::OpenOptions::new()
2212 .read(true)
2213 .write(true)
2214 .create(true)
2215 .truncate(false)
2216 .open(&lock_path)
2217 .map_err(|error| map_io_err(error, "database_gc_lock_open"))?;
2218 fs4::FileExt::try_lock(&advisory_guard)
2219 .map_err(|error| map_io_err(error.into(), "database_gc_lock_try_acquire"))?;
2220 Ok(DatabaseGcOwnerGuard {
2221 _process_guard: process_guard,
2222 _advisory_guard: Some(advisory_guard),
2223 database_path: Some(database_path),
2224 })
2225}
2226
2227impl std::fmt::Debug for DatabaseGcOwnerGuard {
2228 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2229 formatter
2230 .debug_struct("DatabaseGcOwnerGuard")
2231 .field("database_path", &self.database_path)
2232 .finish_non_exhaustive()
2233 }
2234}
2235
2236impl DatabaseGcOwnerGuard {
2237 pub fn database_path(&self) -> Option<&Path> {
2240 self.database_path.as_deref()
2241 }
2242}
2243
2244pub async fn acquire_database_gc_owner(sql: &dyn SqlAccess) -> StorageResult<DatabaseGcOwnerGuard> {
2248 let database_path = sql.database_path();
2249 tokio::task::spawn_blocking(move || acquire_database_gc_owner_for_path_blocking(database_path))
2250 .await
2251 .map_err(|error| {
2252 StorageError::driver(StorageCapability::Blob, "acquire_database_gc_owner", error)
2253 })?
2254}
2255
2256fn root_write_locks() -> &'static StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>> {
2267 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>>> =
2268 OnceLock::new();
2269 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2270}
2271
2272fn write_lock_for_root(root: &Path) -> std::io::Result<Arc<tokio::sync::Mutex<()>>> {
2281 let canonical = root.canonicalize()?;
2282 let mut locks = root_write_locks()
2283 .lock()
2284 .unwrap_or_else(std::sync::PoisonError::into_inner);
2285 Ok(locks
2286 .entry(canonical)
2287 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
2288 .clone())
2289}
2290
2291#[derive(Debug)]
2293pub struct FsBlobStore {
2294 root: PathBuf,
2295 root_handle: Arc<fs::File>,
2299 floor_bytes: u64,
2300 write_lock: Arc<tokio::sync::Mutex<()>>,
2317 orphan_sweep_grace: Duration,
2324}
2325
2326impl FsBlobStore {
2327 pub const DEFAULT_FLOOR_BYTES: u64 = 100_000_000_000;
2330
2331 pub const DEFAULT_ORPHAN_SWEEP_GRACE: Duration = Duration::from_secs(3600);
2336
2337 pub fn new(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
2339 fs::create_dir_all(&root)?;
2340 Self::open_existing(root, floor_bytes)
2341 }
2342
2343 pub fn open_existing(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
2348 let root = root.canonicalize()?;
2352 let metadata = fs::metadata(&root)?;
2353 if !metadata.is_dir() {
2354 return Err(SqliteError::InvalidData(format!(
2355 "blob store root is not a directory: {}",
2356 root.display()
2357 )));
2358 }
2359 let root_handle = Arc::new(open_blob_root_handle(&root)?);
2360 let write_lock = write_lock_for_root(&root)?;
2361 verify_blob_root_identity(&root, &root_handle)?;
2362 Ok(Self {
2363 root,
2364 root_handle,
2365 floor_bytes,
2366 write_lock,
2367 orphan_sweep_grace: Self::DEFAULT_ORPHAN_SWEEP_GRACE,
2368 })
2369 }
2370
2371 pub fn with_orphan_sweep_grace(mut self, grace_period: Duration) -> Self {
2374 self.orphan_sweep_grace = grace_period;
2375 self
2376 }
2377
2378 pub fn root(&self) -> &Path {
2380 &self.root
2381 }
2382}
2383
2384#[async_trait]
2385impl BlobStore for FsBlobStore {
2386 async fn put(&self, bytes: Vec<u8>) -> StorageResult<ContentRef> {
2387 let owned_guard = self.write_lock.clone().lock_owned().await;
2397 let root = self.root.clone();
2398 let root_handle = Arc::clone(&self.root_handle);
2399 let floor_bytes = self.floor_bytes;
2400 #[cfg(test)]
2407 let hook = sync_hook::take(&root);
2408 tokio::task::spawn_blocking(move || {
2409 #[cfg_attr(not(test), allow(clippy::let_and_return))]
2414 let result = {
2415 let _owned_guard = owned_guard;
2416 #[cfg(test)]
2417 if let Some(h) = &hook {
2418 let _ = h.reached.send(());
2419 let _ = h.release.recv();
2420 }
2421 put_blocking_from_root_handle(&root, &root_handle, floor_bytes, bytes)
2422 };
2423 #[cfg(test)]
2424 if let Some(h) = &hook {
2425 let _ = h.done.send(());
2426 }
2427 result
2428 })
2429 .await
2430 .map_err(|e| StorageError::driver(StorageCapability::Blob, "put", e))?
2431 }
2432
2433 async fn get_bounded_verified(
2434 &self,
2435 content_ref: &ContentRef,
2436 max_bytes: u64,
2437 ) -> StorageResult<Vec<u8>> {
2438 if max_bytes > MAX_BLOB_WHOLE_BYTES {
2442 return Err(StorageError::InvalidInput {
2443 capability: StorageCapability::Blob,
2444 operation: "get_bounded_verified".into(),
2445 message: format!(
2446 "max_bytes {max_bytes} exceeds the {MAX_BLOB_WHOLE_BYTES}-byte portable whole-buffer envelope"
2447 ),
2448 });
2449 }
2450
2451 let root = self.root.clone();
2452 let root_handle = Arc::clone(&self.root_handle);
2453 let content_ref = content_ref.clone();
2454 #[cfg(test)]
2455 let read_hook = bounded_read_sync_hook::take(&root);
2456 tokio::task::spawn_blocking(move || {
2457 let mut file = open_blob_shard_file_no_follow(&root, &root_handle, &content_ref)
2461 .map_err(|e| {
2462 if e.kind() == std::io::ErrorKind::NotFound {
2463 StorageError::NotFound {
2464 capability: StorageCapability::Blob,
2465 resource: "blob",
2466 key: content_ref.to_string(),
2467 }
2468 } else if e.kind() == std::io::ErrorKind::Unsupported {
2469 StorageError::Unsupported {
2470 capability: StorageCapability::Blob,
2471 operation: "get_bounded_verified".into(),
2472 message: e.to_string(),
2473 }
2474 } else {
2475 map_io_err(e, "get_bounded_verified.open")
2476 }
2477 })?;
2478
2479 let metadata_bytes = file
2480 .metadata()
2481 .map_err(|e| map_io_err(e, "get_bounded_verified.metadata"))?
2482 .len();
2483 #[cfg(test)]
2484 if let Some(hook) = &read_hook {
2485 let _ = hook.reached.send(());
2486 let _ = hook.release.recv();
2487 }
2488 if metadata_bytes > max_bytes {
2489 return Err(StorageError::BlobTooLarge {
2490 content_ref,
2491 max_bytes,
2492 observed_at_least: metadata_bytes,
2493 });
2494 }
2495
2496 let mut bytes = Vec::with_capacity(metadata_bytes as usize);
2499 (&mut file)
2500 .take(max_bytes + 1)
2501 .read_to_end(&mut bytes)
2502 .map_err(|e| map_io_err(e, "get_bounded_verified.read"))?;
2503 let actual_bytes = bytes.len() as u64;
2504 if actual_bytes > max_bytes {
2505 return Err(StorageError::BlobTooLarge {
2506 content_ref,
2507 max_bytes,
2508 observed_at_least: actual_bytes,
2509 });
2510 }
2511 if metadata_bytes != actual_bytes {
2512 return Err(StorageError::BlobSizeMismatch {
2513 content_ref,
2514 metadata_bytes,
2515 actual_bytes,
2516 });
2517 }
2518
2519 let actual = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
2520 if actual != content_ref {
2521 return Err(StorageError::BlobDigestMismatch {
2522 expected: content_ref,
2523 actual,
2524 });
2525 }
2526 Ok(bytes)
2527 })
2528 .await
2529 .map_err(|e| StorageError::driver(StorageCapability::Blob, "get_bounded_verified", e))?
2530 }
2531
2532 async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool> {
2533 let root = self.root.clone();
2534 let root_handle = Arc::clone(&self.root_handle);
2535 let content_ref = content_ref.clone();
2536 tokio::task::spawn_blocking(move || {
2537 match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2538 Ok(_) => Ok(true),
2539 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
2540 Err(error) => Err(map_io_err(error, "exists")),
2541 }
2542 })
2543 .await
2544 .map_err(|e| StorageError::driver(StorageCapability::Blob, "exists", e))?
2545 }
2546
2547 async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>> {
2548 let root = self.root.clone();
2549 let root_handle = Arc::clone(&self.root_handle);
2550 let content_ref = content_ref.clone();
2551 tokio::task::spawn_blocking(move || {
2552 match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2553 Ok(file) => file
2554 .metadata()
2555 .map(|metadata| Some(metadata.len()))
2556 .map_err(|error| map_io_err(error, "size")),
2557 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
2558 Err(error) => Err(map_io_err(error, "size")),
2559 }
2560 })
2561 .await
2562 .map_err(|e| StorageError::driver(StorageCapability::Blob, "size", e))?
2563 }
2564
2565 async fn delete(&self, content_ref: &ContentRef) -> StorageResult<bool> {
2566 let root = self.root.clone();
2567 let root_handle = Arc::clone(&self.root_handle);
2568 let content_ref = content_ref.clone();
2569 tokio::task::spawn_blocking(move || {
2570 match unlink_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2571 Ok(()) => Ok(true),
2572 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
2573 Err(e) => Err(map_io_err(e, "delete")),
2574 }
2575 })
2576 .await
2577 .map_err(|e| StorageError::driver(StorageCapability::Blob, "delete", e))?
2578 }
2579
2580 async fn orphan_sweep(
2590 &self,
2591 config: &BlobOrphanSweepConfig,
2592 ) -> StorageResult<BlobOrphanSweepResult> {
2593 let _ = config;
2594 Err(StorageError::Unsupported {
2595 capability: StorageCapability::Blob,
2596 operation: "orphan_sweep".into(),
2597 message: "caller-snapshot orphan_sweep is disabled in this compatibility release; \
2598 it cannot prove a completed V21 attachment epoch, use \
2599 transactional_orphan_sweep instead"
2600 .into(),
2601 })
2602 }
2603
2604 async fn transactional_orphan_sweep(
2626 &self,
2627 sql: &dyn SqlAccess,
2628 dry_run: bool,
2629 ) -> StorageResult<BlobOrphanSweepResult> {
2630 if !blob_gc_fencing_complete(sql).await? {
2635 return Err(unsupported_blob_gc_epoch());
2636 }
2637
2638 let database_path = sql.database_path();
2646 let lock_database_path = database_path.clone();
2647 #[cfg(test)]
2648 let hook_database_path = database_path.clone();
2649 let (database_guard, database_file_guard) = tokio::task::spawn_blocking(move || {
2650 let process_guard = sweep_lock_for_database(lock_database_path.as_deref()).acquire();
2651 let file_guard = acquire_database_gc_lock(lock_database_path.as_deref())?;
2652 #[cfg(test)]
2653 if let Some(hook) = db_ownership_sync_hook::take(hook_database_path.as_deref()) {
2654 let _ = hook.reached.send(());
2655 let _ = hook.release.recv();
2656 }
2657 Ok::<_, StorageError>((process_guard, file_guard))
2658 })
2659 .await
2660 .map_err(|e| {
2661 StorageError::driver(
2662 StorageCapability::Blob,
2663 "transactional_orphan_sweep_lock",
2664 e,
2665 )
2666 })??;
2667
2668 if !blob_gc_fencing_complete(sql).await? {
2674 return Err(unsupported_blob_gc_epoch());
2675 }
2676
2677 let root_guard = self.write_lock.clone().lock_owned().await;
2678 let root = self.root.clone();
2679 let root_handle = Arc::clone(&self.root_handle);
2680 let scan_root = root.clone();
2681 let scan_root_handle = Arc::clone(&root_handle);
2682 let grace_period = self.orphan_sweep_grace;
2683 let (write_guards, canonical_root, prepared) = tokio::task::spawn_blocking(move || {
2684 verify_blob_root_identity(&scan_root, &scan_root_handle)
2685 .map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
2686 let canonical_root = scan_root;
2689 let root_write_guard =
2690 acquire_root_write_lock_anchored(&canonical_root, &scan_root_handle)?;
2691 let candidates = walk_blob_files_from_root_handle(&scan_root_handle, &canonical_root)
2692 .map_err(|e| map_io_err(e, "transactional_orphan_sweep_walk"))?;
2693 verify_blob_root_identity(&canonical_root, &scan_root_handle)
2694 .map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
2695 let prepared = prepare_transactional_sweep(candidates, grace_period);
2696 Ok::<_, StorageError>((
2697 (
2698 database_guard,
2699 database_file_guard,
2700 root_guard,
2701 root_write_guard,
2702 ),
2703 canonical_root,
2704 prepared,
2705 ))
2706 })
2707 .await
2708 .map_err(|e| {
2709 StorageError::driver(
2710 StorageCapability::Blob,
2711 "transactional_orphan_sweep_walk",
2712 e,
2713 )
2714 })??;
2715 let root_key = blob_root_key(&canonical_root);
2716 validate_blob_gc_evidence(sql).await?;
2717 blob_gc_fence_probe(sql).await?;
2718 if !dry_run {
2719 loop {
2720 let released = release_abandoned_blob_gc_claim_batch(sql).await?;
2721 if released < BLOB_GC_CLAIM_BATCH_SIZE as u64 {
2722 break;
2723 }
2724 }
2725 }
2726
2727 let mut write_guards = write_guards;
2728 let mut result = prepared.result;
2729 let mut delete_error = None;
2730 #[cfg(test)]
2731 let mut hook: Option<sync_hook::Hook> = None;
2732 #[cfg(not(test))]
2733 let mut hook: Option<()> = None;
2734 #[cfg(test)]
2735 let mut hook_paused = false;
2736
2737 for candidates in prepared.candidates.chunks(BLOB_GC_CLAIM_BATCH_SIZE) {
2744 let batch = claim_blob_gc_batch(sql, root_key.clone(), candidates, dry_run).await?;
2745 result.grace_period_skipped += batch.grace_period_skipped;
2746 result.would_delete += batch.would_delete;
2747 if dry_run {
2748 continue;
2749 }
2750
2751 let claimed_refs = parse_blob_gc_claim_rows(batch.claimed_rows)?;
2752 if claimed_refs.is_empty() {
2753 continue;
2754 }
2755
2756 #[cfg(test)]
2757 if hook.is_none() {
2758 hook = sync_hook::take(&root);
2759 }
2760 #[cfg(test)]
2761 let pause_hook = hook.is_some() && !hook_paused;
2762 #[cfg(test)]
2763 if pause_hook {
2764 hook_paused = true;
2765 }
2766
2767 let delete_root = canonical_root.clone();
2768 let delete_root_handle = Arc::clone(&root_handle);
2769 let (returned_guards, deleted, batch_delete_error, returned_hook) =
2770 tokio::task::spawn_blocking(move || {
2771 #[cfg(test)]
2772 if pause_hook {
2773 if let Some(hook) = &hook {
2774 let _ = hook.reached.send(());
2775 let _ = hook.release.recv();
2776 }
2777 }
2778
2779 let mut deleted = 0_u64;
2780 let mut first_error = None;
2781 for content_ref in claimed_refs {
2782 match unlink_blob_shard_file_no_follow(
2783 &delete_root,
2784 &delete_root_handle,
2785 &content_ref,
2786 ) {
2787 Ok(()) => deleted += 1,
2788 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
2789 Err(error) => {
2790 first_error =
2791 Some(map_io_err(error, "transactional_orphan_sweep_delete"));
2792 break;
2793 }
2794 }
2795 }
2796 (write_guards, deleted, first_error, hook)
2801 })
2802 .await
2803 .map_err(|error| {
2804 StorageError::driver(
2805 StorageCapability::Blob,
2806 "transactional_orphan_sweep_delete",
2807 error,
2808 )
2809 })?;
2810 write_guards = returned_guards;
2811 hook = returned_hook;
2812 result.deleted += deleted;
2813
2814 release_blob_gc_batch(sql, root_key.clone()).await?;
2819 if batch_delete_error.is_some() {
2820 delete_error = batch_delete_error;
2821 break;
2822 }
2823 }
2824
2825 drop(write_guards);
2826 #[cfg(test)]
2827 if let Some(hook) = hook {
2828 let _ = hook.done.send(());
2829 }
2830 #[cfg(not(test))]
2831 let _ = hook;
2832 if let Some(error) = delete_error {
2833 return Err(error);
2834 }
2835 Ok(result)
2836 }
2837}
2838
2839#[cfg(test)]
2854mod sync_hook {
2855 use std::collections::{HashMap, VecDeque};
2856 use std::path::{Path, PathBuf};
2857 use std::sync::mpsc::{Receiver, Sender};
2858 use std::sync::{Mutex as StdMutex, OnceLock};
2859
2860 pub(super) struct Hook {
2861 pub(super) reached: Sender<()>,
2862 pub(super) release: Receiver<()>,
2863 pub(super) done: Sender<()>,
2864 }
2865
2866 fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
2867 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
2868 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2869 }
2870
2871 pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>, Receiver<()>) {
2874 let canonical = root
2875 .canonicalize()
2876 .expect("root must exist before installing a sync_hook");
2877 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
2878 let (release_tx, release_rx) = std::sync::mpsc::channel();
2879 let (done_tx, done_rx) = std::sync::mpsc::channel();
2880 registry()
2881 .lock()
2882 .unwrap_or_else(std::sync::PoisonError::into_inner)
2883 .entry(canonical)
2884 .or_default()
2885 .push_back(Hook {
2886 reached: reached_tx,
2887 release: release_rx,
2888 done: done_tx,
2889 });
2890 (reached_rx, release_tx, done_rx)
2891 }
2892
2893 pub(super) fn take(root: &Path) -> Option<Hook> {
2899 let canonical = root.canonicalize().ok()?;
2900 registry()
2901 .lock()
2902 .unwrap_or_else(std::sync::PoisonError::into_inner)
2903 .get_mut(&canonical)
2904 .and_then(VecDeque::pop_front)
2905 }
2906}
2907
2908#[cfg(all(test, unix))]
2924mod walk_leaf_sync_hook {
2925 use std::collections::HashMap;
2926 use std::path::{Path, PathBuf};
2927 use std::sync::mpsc::{Receiver, Sender};
2928 use std::sync::{Mutex as StdMutex, OnceLock};
2929
2930 pub(super) struct Hook {
2931 pub(super) reached: Sender<()>,
2932 pub(super) release: Receiver<()>,
2933 }
2934
2935 fn registry() -> &'static StdMutex<HashMap<PathBuf, Hook>> {
2936 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Hook>>> = OnceLock::new();
2937 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2938 }
2939
2940 pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
2941 let canonical = root
2942 .canonicalize()
2943 .expect("root must exist before installing a walk_leaf_sync_hook");
2944 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
2945 let (release_tx, release_rx) = std::sync::mpsc::channel();
2946 registry()
2947 .lock()
2948 .unwrap_or_else(std::sync::PoisonError::into_inner)
2949 .insert(
2950 canonical,
2951 Hook {
2952 reached: reached_tx,
2953 release: release_rx,
2954 },
2955 );
2956 (reached_rx, release_tx)
2957 }
2958
2959 pub(super) fn take(root: &Path) -> Option<Hook> {
2962 let canonical = root.canonicalize().ok()?;
2963 registry()
2964 .lock()
2965 .unwrap_or_else(std::sync::PoisonError::into_inner)
2966 .remove(&canonical)
2967 }
2968}
2969
2970#[cfg(test)]
2975mod bounded_read_sync_hook {
2976 use std::collections::{HashMap, VecDeque};
2977 use std::path::{Path, PathBuf};
2978 use std::sync::mpsc::{Receiver, Sender};
2979 use std::sync::{Mutex as StdMutex, OnceLock};
2980
2981 pub(super) struct Hook {
2982 pub(super) reached: Sender<()>,
2983 pub(super) release: Receiver<()>,
2984 }
2985
2986 fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
2987 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
2988 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2989 }
2990
2991 pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
2992 let canonical = root
2993 .canonicalize()
2994 .expect("root must exist before installing a bounded-read hook");
2995 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
2996 let (release_tx, release_rx) = std::sync::mpsc::channel();
2997 registry()
2998 .lock()
2999 .unwrap_or_else(std::sync::PoisonError::into_inner)
3000 .entry(canonical)
3001 .or_default()
3002 .push_back(Hook {
3003 reached: reached_tx,
3004 release: release_rx,
3005 });
3006 (reached_rx, release_tx)
3007 }
3008
3009 pub(super) fn take(root: &Path) -> Option<Hook> {
3010 let canonical = root.canonicalize().ok()?;
3011 registry()
3012 .lock()
3013 .unwrap_or_else(std::sync::PoisonError::into_inner)
3014 .get_mut(&canonical)
3015 .and_then(VecDeque::pop_front)
3016 }
3017}
3018
3019#[cfg(test)]
3026mod db_ownership_sync_hook {
3027 use std::collections::{HashMap, VecDeque};
3028 use std::path::{Path, PathBuf};
3029 use std::sync::mpsc::{Receiver, Sender};
3030 use std::sync::{Mutex as StdMutex, OnceLock};
3031
3032 pub(super) struct Hook {
3033 pub(super) reached: Sender<()>,
3034 pub(super) release: Receiver<()>,
3035 }
3036
3037 fn registry() -> &'static StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>> {
3038 static REGISTRY: OnceLock<StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>>> =
3039 OnceLock::new();
3040 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
3041 }
3042
3043 pub(super) fn install(database_path: Option<&Path>) -> (Receiver<()>, Sender<()>) {
3044 let key = database_path.map(Path::to_path_buf);
3045 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
3046 let (release_tx, release_rx) = std::sync::mpsc::channel();
3047 registry()
3048 .lock()
3049 .unwrap_or_else(std::sync::PoisonError::into_inner)
3050 .entry(key)
3051 .or_default()
3052 .push_back(Hook {
3053 reached: reached_tx,
3054 release: release_rx,
3055 });
3056 (reached_rx, release_tx)
3057 }
3058
3059 pub(super) fn take(database_path: Option<&Path>) -> Option<Hook> {
3060 let key = database_path.map(Path::to_path_buf);
3061 registry()
3062 .lock()
3063 .unwrap_or_else(std::sync::PoisonError::into_inner)
3064 .get_mut(&key)
3065 .and_then(VecDeque::pop_front)
3066 }
3067}
3068
3069#[cfg(test)]
3070mod tests {
3071 use super::*;
3072
3073 fn store(floor_bytes: u64) -> (tempfile::TempDir, FsBlobStore) {
3074 let dir = tempfile::tempdir().unwrap();
3075 let root = dir.path().join("blobs");
3076 let store = FsBlobStore::new(root, floor_bytes)
3080 .unwrap()
3081 .with_orphan_sweep_grace(Duration::ZERO);
3082 (dir, store)
3083 }
3084
3085 fn prepare_v20_gc_fixture(conn: &mut rusqlite::Connection) {
3088 conn.execute_batch(include_str!("../../sql/schema-migrations-table.sql"))
3089 .expect("create migration ledger");
3090 for migration in crate::MIGRATIONS
3091 .iter()
3092 .filter(|migration| migration.version <= 20)
3093 {
3094 let tx = conn.transaction().expect("begin historical migration");
3095 tx.execute_batch(migration.up)
3096 .expect("apply historical migration body");
3097 tx.execute(
3098 "INSERT INTO _schema_migrations (version, name, applied_at) \
3099 VALUES (?1, ?2, 0)",
3100 rusqlite::params![migration.version, migration.name],
3101 )
3102 .expect("record historical migration");
3103 tx.commit().expect("commit historical migration");
3104 }
3105 }
3106
3107 fn prepare_completed_v21_gc_fixture(conn: &mut rusqlite::Connection) {
3111 let version = crate::run_migrations(conn).expect("prepare canonical completed V21");
3112 assert_eq!(
3118 version,
3119 crate::migrations::ATTACHMENT_CUTOVER_VERSION,
3120 "GC-gate fixtures need a completed-V21 ledger; a later migration chain \
3121 must provide a pinned through-V21 fixture builder for these tests"
3122 );
3123 }
3124
3125 #[tokio::test]
3126 async fn completed_v21_gc_gate_requires_new_indexes_and_absent_legacy_column() {
3127 let dir = tempfile::tempdir().unwrap();
3128 let db_path = dir.path().join("khive.db");
3129 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3130 {
3131 let mut writer = backend.pool().writer().unwrap();
3132 prepare_completed_v21_gc_fixture(writer.conn_mut());
3133 }
3134 assert!(blob_gc_fencing_complete(backend.sql().as_ref())
3135 .await
3136 .unwrap());
3137
3138 {
3139 let writer = backend.pool().writer().unwrap();
3140 writer
3141 .conn()
3142 .execute_batch("DROP INDEX idx_attachments_content_ref")
3143 .unwrap();
3144 }
3145 assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
3146 .await
3147 .unwrap());
3148
3149 {
3150 let writer = backend.pool().writer().unwrap();
3151 writer
3152 .conn()
3153 .execute_batch(
3154 "CREATE INDEX idx_attachments_content_ref \
3155 ON attachments(content_ref); \
3156 ALTER TABLE entities ADD COLUMN content_ref TEXT",
3157 )
3158 .unwrap();
3159 }
3160 assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
3161 .await
3162 .unwrap());
3163 }
3164
3165 #[tokio::test]
3175 async fn completed_v21_gate_acceptance_matrix_rejects_each_removed_fact_independently() {
3176 let cases: &[(&str, &str)] = &[
3177 (
3178 "blob_gc_claims_content_ref_index_dropped",
3179 "DROP INDEX idx_blob_gc_claims_content_ref",
3180 ),
3181 (
3182 "v21_ledger_row_deleted",
3183 "DELETE FROM _schema_migrations WHERE version = 21",
3184 ),
3185 (
3186 "v21_ledger_row_renamed",
3187 "UPDATE _schema_migrations SET name = 'not_attachments_first_class' \
3188 WHERE version = 21",
3189 ),
3190 (
3200 "below_v21_ledger_rows_deleted",
3201 "DELETE FROM _schema_migrations WHERE version < 21",
3202 ),
3203 (
3204 "marker_row_deleted",
3205 "DELETE FROM attachment_cutover_state WHERE singleton = 1",
3206 ),
3207 (
3208 "marker_completed_at_null_while_state_complete",
3215 "DROP TABLE attachment_cutover_state; \
3216 CREATE TABLE attachment_cutover_state ( \
3217 singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
3218 state TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
3219 started_at INTEGER NOT NULL, \
3220 completed_at INTEGER \
3221 ) STRICT; \
3222 INSERT INTO attachment_cutover_state \
3223 (singleton, state, started_at, completed_at) \
3224 VALUES (1, 'complete', 21, NULL)",
3225 ),
3226 (
3227 "insert_fence_dropped_alone",
3228 "DROP TRIGGER attachments_reject_claimed_blob_insert",
3229 ),
3230 ];
3231
3232 for (case, mutation_sql) in cases {
3233 let dir = tempfile::tempdir().unwrap();
3234 let db_path = dir.path().join("khive.db");
3235 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3236 {
3237 let mut writer = backend.pool().writer().unwrap();
3238 prepare_completed_v21_gc_fixture(writer.conn_mut());
3239 writer
3240 .conn_mut()
3241 .execute_batch(mutation_sql)
3242 .unwrap_or_else(|e| panic!("case {case}: failed to apply mutation: {e}"));
3243 }
3244
3245 assert!(
3246 !blob_gc_fencing_complete(backend.sql().as_ref())
3247 .await
3248 .unwrap(),
3249 "case {case}: gate must reject with this fact removed"
3250 );
3251
3252 let store = Arc::new(
3253 FsBlobStore::new(dir.path().join("blobs"), 0)
3254 .unwrap()
3255 .with_orphan_sweep_grace(Duration::ZERO),
3256 );
3257 let orphan = store
3258 .put(format!("gate matrix orphan for {case}").into_bytes())
3259 .await
3260 .unwrap();
3261 let abandoned_ref = "c".repeat(64);
3262 {
3263 let writer = backend.pool().writer().unwrap();
3264 writer
3265 .conn()
3266 .execute(
3267 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
3268 VALUES ('gate-matrix-abandoned', ?1, 1)",
3269 [abandoned_ref.as_str()],
3270 )
3271 .unwrap();
3272 }
3273
3274 let _root_guard = store.write_lock.clone().lock_owned().await;
3275 for dry_run in [true, false] {
3276 let outcome = tokio::time::timeout(
3277 Duration::from_secs(1),
3278 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
3279 )
3280 .await
3281 .unwrap_or_else(|_| panic!("case {case}: refusal must precede the root wait"));
3282 assert!(
3283 matches!(outcome, Err(StorageError::Unsupported { .. })),
3284 "case {case} dry_run={dry_run}: expected Unsupported, got {outcome:?}"
3285 );
3286 }
3287
3288 assert!(
3289 store.exists(&orphan).await.unwrap(),
3290 "case {case}: a refused sweep must not delete anything"
3291 );
3292 let reader = backend.pool().reader().unwrap();
3293 let remaining: i64 = reader
3294 .conn()
3295 .query_row(
3296 "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = 'gate-matrix-abandoned'",
3297 [],
3298 |row| row.get(0),
3299 )
3300 .unwrap();
3301 assert_eq!(remaining, 1, "case {case}: refusal must not recover claims");
3302
3303 let probe_claims: i64 = reader
3308 .conn()
3309 .query_row(
3310 "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key GLOB '__fence_probe-*'",
3311 [],
3312 |row| row.get(0),
3313 )
3314 .unwrap();
3315 assert_eq!(
3316 probe_claims, 0,
3317 "case {case}: refusal must leave no fence-probe claim residue"
3318 );
3319 let probe_attachments: i64 = reader
3320 .conn()
3321 .query_row(
3322 "SELECT COUNT(*) FROM attachments \
3323 WHERE record_uuid GLOB '__blob-gc-fence-probe-*'",
3324 [],
3325 |row| row.get(0),
3326 )
3327 .unwrap();
3328 assert_eq!(
3329 probe_attachments, 0,
3330 "case {case}: refusal must leave no fence-probe attachment residue"
3331 );
3332 }
3333 }
3334
3335 #[tokio::test]
3351 async fn gate_rejects_ledger_ahead_of_binary_latest() {
3352 let dir = tempfile::tempdir().unwrap();
3353 let db_path = dir.path().join("khive.db");
3354 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3355 {
3356 let mut writer = backend.pool().writer().unwrap();
3357 prepare_completed_v21_gc_fixture(writer.conn_mut());
3358 }
3360 assert!(
3361 blob_gc_fencing_complete(backend.sql().as_ref())
3362 .await
3363 .unwrap(),
3364 "control: completed fixture at the binary's latest version must pass"
3365 );
3366
3367 {
3368 let writer = backend.pool().writer().unwrap();
3369 writer
3370 .conn()
3371 .execute(
3372 "INSERT INTO _schema_migrations (version, name, applied_at) \
3373 VALUES (?1, 'post_cutover_feature', unixepoch())",
3374 [i64::from(crate::migrations::latest_schema_version()) + 1],
3375 )
3376 .unwrap();
3377 }
3378 assert!(
3379 !blob_gc_fencing_complete(backend.sql().as_ref())
3380 .await
3381 .unwrap(),
3382 "a ledger ahead of the binary's latest schema version must fail closed"
3383 );
3384 }
3385
3386 #[tokio::test]
3394 async fn gate_rejects_incomplete_ledger_behind_v21() {
3395 let dir = tempfile::tempdir().unwrap();
3396 let db_path = dir.path().join("khive.db");
3397 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3398 {
3399 let mut writer = backend.pool().writer().unwrap();
3400 prepare_completed_v21_gc_fixture(writer.conn_mut());
3401 }
3402 assert!(
3403 blob_gc_fencing_complete(backend.sql().as_ref())
3404 .await
3405 .unwrap(),
3406 "control: the untouched completed fixture must pass"
3407 );
3408
3409 {
3410 let writer = backend.pool().writer().unwrap();
3411 writer
3412 .conn()
3413 .execute_batch("DELETE FROM _schema_migrations WHERE version < 21")
3414 .unwrap();
3415 let (v21_named, max_version): (i64, i64) = writer
3417 .conn()
3418 .query_row(
3419 "SELECT (SELECT COUNT(*) FROM _schema_migrations \
3420 WHERE version = 21 AND name = 'attachments_first_class'), \
3421 (SELECT MAX(version) FROM _schema_migrations)",
3422 [],
3423 |row| Ok((row.get(0)?, row.get(1)?)),
3424 )
3425 .unwrap();
3426 assert_eq!(
3427 (v21_named, max_version),
3428 (1, 21),
3429 "fixture must keep the named V21 row and MAX(version) = 21 so the \
3430 rejection can only come from the contiguity clause"
3431 );
3432 }
3433 assert!(
3434 !blob_gc_fencing_complete(backend.sql().as_ref())
3435 .await
3436 .unwrap(),
3437 "an incomplete ledger behind V21 must fail closed despite a valid terminal row"
3438 );
3439 }
3440
3441 #[tokio::test]
3448 async fn completed_v21_gate_fails_closed_when_marker_read_errors() {
3449 let dir = tempfile::tempdir().unwrap();
3450 let db_path = dir.path().join("khive.db");
3451 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3452 {
3453 let mut writer = backend.pool().writer().unwrap();
3454 prepare_completed_v21_gc_fixture(writer.conn_mut());
3455 writer
3460 .conn_mut()
3461 .execute_batch(
3462 "DROP TABLE attachment_cutover_state; \
3463 CREATE TABLE attachment_cutover_state ( \
3464 singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
3465 state TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
3466 started_at INTEGER NOT NULL \
3467 ) STRICT; \
3468 INSERT INTO attachment_cutover_state (singleton, state, started_at) \
3469 VALUES (1, 'complete', 21)",
3470 )
3471 .unwrap();
3472 }
3473
3474 let gate_error = blob_gc_fencing_complete(backend.sql().as_ref())
3475 .await
3476 .expect_err("a marker read error must propagate, not silently resolve to false");
3477 assert!(
3478 !matches!(gate_error, StorageError::Unsupported { .. }),
3479 "a read error is a distinct failure from the typed epoch refusal: {gate_error:?}"
3480 );
3481
3482 let store = Arc::new(
3483 FsBlobStore::new(dir.path().join("blobs"), 0)
3484 .unwrap()
3485 .with_orphan_sweep_grace(Duration::ZERO),
3486 );
3487 let orphan = store
3488 .put(b"marker read error orphan".to_vec())
3489 .await
3490 .unwrap();
3491 let _root_guard = store.write_lock.clone().lock_owned().await;
3492 for dry_run in [true, false] {
3493 let outcome = tokio::time::timeout(
3494 Duration::from_secs(1),
3495 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
3496 )
3497 .await
3498 .expect("a marker read error must fail before waiting on the root lock");
3499 assert!(
3500 outcome.is_err(),
3501 "dry_run={dry_run}: expected the sweep to fail closed, got {outcome:?}"
3502 );
3503 }
3504 assert!(store.exists(&orphan).await.unwrap());
3505 }
3506
3507 #[test]
3508 fn database_sweep_owner_is_keyed_by_database_not_blob_root() {
3509 let dir = tempfile::tempdir().unwrap();
3510 let database = dir.path().join("khive.db");
3511 let same_database_a = sweep_lock_for_database(Some(&database));
3512 let same_database_b = sweep_lock_for_database(Some(&database));
3513 let other_database = sweep_lock_for_database(Some(&dir.path().join("other.db")));
3514 let mut expected_lock_path = database.as_os_str().to_os_string();
3515 expected_lock_path.push(DATABASE_GC_LOCK_SUFFIX);
3516
3517 assert!(Arc::ptr_eq(&same_database_a, &same_database_b));
3518 assert!(!Arc::ptr_eq(&same_database_a, &other_database));
3519 assert_eq!(
3520 database_gc_lock_path(&database),
3521 PathBuf::from(expected_lock_path)
3522 );
3523 }
3524
3525 #[tokio::test]
3526 async fn database_gc_owner_holds_process_and_advisory_fences_until_drop() {
3527 let dir = tempfile::tempdir().unwrap();
3528 let database = dir.path().join("owner.db");
3529 let backend = crate::StorageBackend::sqlite(&database).unwrap();
3530 let owner = acquire_database_gc_owner(backend.sql().as_ref())
3531 .await
3532 .unwrap();
3533 let canonical_database = owner
3534 .database_path()
3535 .expect("file-backed owner path")
3536 .to_path_buf();
3537
3538 assert!(
3539 sweep_lock_for_database(Some(&canonical_database))
3540 .try_acquire()
3541 .is_none(),
3542 "boot and sweep must share one process-local database owner"
3543 );
3544 let external = fs::OpenOptions::new()
3545 .read(true)
3546 .write(true)
3547 .open(database_gc_lock_path(&canonical_database))
3548 .unwrap();
3549 assert!(
3550 matches!(
3551 fs4::FileExt::try_lock(&external),
3552 Err(fs4::TryLockError::WouldBlock)
3553 ),
3554 "the reusable owner must also retain the cross-process advisory fence"
3555 );
3556
3557 drop(owner);
3558 fs4::FileExt::try_lock(&external).expect("owner drop releases advisory fence");
3559 }
3560
3561 #[cfg(unix)]
3562 #[test]
3563 fn database_gc_lock_path_preserves_non_utf8_identity() {
3564 use std::os::unix::ffi::{OsStrExt, OsStringExt};
3565
3566 let database = PathBuf::from(std::ffi::OsString::from_vec(
3567 b"khive-non-utf8-\xff.db".to_vec(),
3568 ));
3569 let lock_path = database_gc_lock_path(&database);
3570 let mut expected = database.as_os_str().as_bytes().to_vec();
3571 expected.extend_from_slice(DATABASE_GC_LOCK_SUFFIX.as_bytes());
3572 assert_eq!(lock_path.as_os_str().as_bytes(), expected);
3573 }
3574
3575 #[cfg(unix)]
3584 #[tokio::test]
3585 async fn unlink_blob_shard_refuses_symlinked_shard_dir_and_still_sweeps_real_shard() {
3586 let (dir, store) = store(0);
3587 let root = dir.path().join("blobs");
3588
3589 let real = store.put(b"real blob content".to_vec()).await.unwrap();
3592
3593 let outside = tempfile::tempdir().unwrap();
3598 let victim = outside.path().join("victim.txt");
3599 fs::write(&victim, b"do not delete me").unwrap();
3600
3601 let real_prefix = &real.as_str()[0..2];
3602 let attack_prefix = if real_prefix == "aa" { "bb" } else { "aa" };
3603 let fake_ref = ContentRef::from_hex(format!("{attack_prefix}{}", "0".repeat(62))).unwrap();
3604 let fake_hex = fake_ref.as_str().to_string();
3605
3606 let shard1 = root.join(attack_prefix);
3607 std::os::unix::fs::symlink(outside.path(), &shard1).unwrap();
3608 fs::create_dir_all(outside.path().join(&fake_hex[2..4])).unwrap();
3614 fs::write(
3615 outside.path().join(&fake_hex[2..4]).join(&fake_hex),
3616 b"decoy",
3617 )
3618 .unwrap();
3619
3620 let root_handle = open_blob_root_handle(&root).unwrap();
3621 let error = unlink_blob_shard_file_no_follow(&root, &root_handle, &fake_ref).unwrap_err();
3622 assert!(
3628 matches!(
3629 error.raw_os_error(),
3630 Some(libc::ELOOP) | Some(libc::ENOTDIR)
3631 ),
3632 "opening a symlinked shard directory must be refused, not followed; got: {error}"
3633 );
3634 assert!(
3635 victim.exists(),
3636 "the file outside the blob root must never be touched by a refused shard-dir open"
3637 );
3638
3639 unlink_blob_shard_file_no_follow(&root, &root_handle, &real).unwrap();
3645 assert!(!store.exists(&real).await.unwrap());
3646 }
3647
3648 async fn recv_blocking(rx: std::sync::mpsc::Receiver<()>) -> bool {
3654 tokio::task::spawn_blocking(move || rx.recv().is_ok())
3655 .await
3656 .expect("recv_blocking thread panicked")
3657 }
3658
3659 #[tokio::test]
3660 async fn put_bounded_get_roundtrip() {
3661 let (_dir, store) = store(0);
3662 let bytes = b"hello blob store".to_vec();
3663 let content_ref = store.put(bytes.clone()).await.unwrap();
3664 let fetched = store
3665 .get_bounded_verified(&content_ref, bytes.len() as u64)
3666 .await
3667 .unwrap();
3668 assert_eq!(fetched, bytes);
3669 }
3670
3671 #[tokio::test]
3672 async fn bounded_verified_get_accepts_exact_and_portable_maximum_limits() {
3673 let (_dir, store) = store(0);
3674 let bytes = b"bounded fs blob".to_vec();
3675 let content_ref = store.put(bytes.clone()).await.unwrap();
3676
3677 assert_eq!(
3678 store
3679 .get_bounded_verified(&content_ref, bytes.len() as u64)
3680 .await
3681 .unwrap(),
3682 bytes
3683 );
3684 assert_eq!(
3685 store
3686 .get_bounded_verified(&content_ref, MAX_BLOB_WHOLE_BYTES)
3687 .await
3688 .unwrap(),
3689 bytes
3690 );
3691 }
3692
3693 #[cfg(unix)]
3694 #[tokio::test]
3695 async fn bounded_verified_get_resolves_a_configured_symlink_root_once() {
3696 use std::os::unix::fs::symlink;
3697
3698 let dir = tempfile::tempdir().unwrap();
3699 let target = dir.path().join("blob-target");
3700 fs::create_dir(&target).unwrap();
3701 let configured = dir.path().join("blob-configured");
3702 symlink(&target, &configured).unwrap();
3703
3704 let store = FsBlobStore::new(configured, 0).unwrap();
3705 assert_eq!(store.root(), target.canonicalize().unwrap());
3706 let bytes = b"symlink-configured root".to_vec();
3707 let content_ref = store.put(bytes.clone()).await.unwrap();
3708 assert_eq!(
3709 store
3710 .get_bounded_verified(&content_ref, bytes.len() as u64)
3711 .await
3712 .unwrap(),
3713 bytes
3714 );
3715 }
3716
3717 #[cfg(unix)]
3718 #[tokio::test]
3719 async fn fs_blob_store_refuses_root_replacement_before_put() {
3720 use std::os::unix::fs::symlink;
3721
3722 let dir = tempfile::tempdir().unwrap();
3723 let root = dir.path().join("blobs");
3724 let store = FsBlobStore::new(root.clone(), 0).unwrap();
3725 let original_root = dir.path().join("blobs-original");
3726 fs::rename(&root, &original_root).unwrap();
3727
3728 let redirected_root = tempfile::tempdir().unwrap();
3729 symlink(redirected_root.path(), &root).unwrap();
3730
3731 let bytes = b"must not land in a replaced root".to_vec();
3732 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
3733 store
3734 .put(bytes)
3735 .await
3736 .expect_err("a root replaced after construction must be refused");
3737
3738 assert!(
3739 !shard_path(redirected_root.path(), &content_ref).exists(),
3740 "put must not publish into the tree selected by the replacement symlink"
3741 );
3742 assert!(
3743 !shard_path(&original_root, &content_ref).exists(),
3744 "a refused put must not mutate the initialization-time root either"
3745 );
3746 }
3747
3748 #[cfg(unix)]
3749 #[tokio::test]
3750 async fn fs_blob_store_refuses_ancestor_replacement_for_existing_blob_operations() {
3751 use std::os::unix::fs::symlink;
3752
3753 let dir = tempfile::tempdir().unwrap();
3754 let ancestor = dir.path().join("store-parent");
3755 let root = ancestor.join("blobs");
3756 let store = FsBlobStore::new(root.clone(), 0).unwrap();
3757 let bytes = b"same bytes in both trees".to_vec();
3758 let content_ref = store.put(bytes.clone()).await.unwrap();
3759
3760 let original_ancestor = dir.path().join("store-parent-original");
3761 fs::rename(&ancestor, &original_ancestor).unwrap();
3762 let original_blob = shard_path(&original_ancestor.join("blobs"), &content_ref);
3763
3764 let redirected_ancestor = tempfile::tempdir().unwrap();
3765 let redirected_blob = shard_path(&redirected_ancestor.path().join("blobs"), &content_ref);
3766 fs::create_dir_all(redirected_blob.parent().unwrap()).unwrap();
3767 fs::write(&redirected_blob, &bytes).unwrap();
3768 symlink(redirected_ancestor.path(), &ancestor).unwrap();
3769
3770 store
3771 .get_bounded_verified(&content_ref, bytes.len() as u64)
3772 .await
3773 .expect_err("a read through a replaced root ancestor must be refused");
3774 store
3775 .exists(&content_ref)
3776 .await
3777 .expect_err("exists through a replaced root ancestor must be refused");
3778 store
3779 .size(&content_ref)
3780 .await
3781 .expect_err("size through a replaced root ancestor must be refused");
3782 store
3783 .delete(&content_ref)
3784 .await
3785 .expect_err("delete through a replaced root ancestor must be refused");
3786
3787 assert!(
3788 original_blob.exists(),
3789 "refusal must preserve the initialization-time blob"
3790 );
3791 assert!(
3792 redirected_blob.exists(),
3793 "refusal must not read as authority or delete the redirected blob"
3794 );
3795 }
3796
3797 #[tokio::test]
3798 async fn bounded_verified_get_rejects_same_size_digest_corruption() {
3799 let (_dir, store) = store(0);
3800 let expected_bytes = b"expected".to_vec();
3801 let actual_bytes = b"mutated!".to_vec();
3802 let expected = store.put(expected_bytes).await.unwrap();
3803 let actual = ContentRef::from_digest_bytes(blake3::hash(&actual_bytes).as_bytes());
3804 fs::write(shard_path(store.root(), &expected), actual_bytes).unwrap();
3805
3806 let err = store.get_bounded_verified(&expected, 8).await.unwrap_err();
3807 assert!(matches!(
3808 err,
3809 StorageError::BlobDigestMismatch {
3810 expected: ref got_expected,
3811 actual: ref got_actual,
3812 } if got_expected == &expected && got_actual == &actual
3813 ));
3814 }
3815
3816 #[tokio::test]
3817 async fn bounded_verified_get_stops_at_max_plus_one_after_file_growth() {
3818 let (_dir, store) = store(0);
3819 let store = Arc::new(store);
3820 let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
3821 let path = shard_path(store.root(), &content_ref);
3822 let (reached, release) = bounded_read_sync_hook::install(store.root());
3823
3824 let read_store = Arc::clone(&store);
3825 let read_ref = content_ref.clone();
3826 let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
3827 assert!(
3828 recv_blocking(reached).await,
3829 "read must reach the metadata seam"
3830 );
3831 let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
3832 writer.write_all(b"efgh-poison-tail").unwrap();
3833 writer.flush().unwrap();
3834 release.send(()).unwrap();
3835
3836 let err = read.await.unwrap().unwrap_err();
3837 assert!(matches!(
3838 err,
3839 StorageError::BlobTooLarge {
3840 content_ref: ref got,
3841 max_bytes: 4,
3842 observed_at_least: 5,
3843 } if got == &content_ref
3844 ));
3845 }
3846
3847 #[tokio::test]
3848 async fn bounded_verified_get_reports_growth_within_limit_as_size_mismatch() {
3849 let (_dir, store) = store(0);
3850 let store = Arc::new(store);
3851 let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
3852 let path = shard_path(store.root(), &content_ref);
3853 let (reached, release) = bounded_read_sync_hook::install(store.root());
3854
3855 let read_store = Arc::clone(&store);
3856 let read_ref = content_ref.clone();
3857 let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 8).await });
3858 assert!(
3859 recv_blocking(reached).await,
3860 "read must reach the metadata seam"
3861 );
3862 let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
3863 writer.write_all(b"ef").unwrap();
3864 writer.flush().unwrap();
3865 release.send(()).unwrap();
3866
3867 let err = read.await.unwrap().unwrap_err();
3868 assert!(matches!(
3869 err,
3870 StorageError::BlobSizeMismatch {
3871 content_ref: ref got,
3872 metadata_bytes: 4,
3873 actual_bytes: 6,
3874 } if got == &content_ref
3875 ));
3876 }
3877
3878 #[tokio::test]
3879 async fn bounded_verified_get_reports_truncation_as_size_mismatch() {
3880 let (_dir, store) = store(0);
3881 let store = Arc::new(store);
3882 let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
3883 let path = shard_path(store.root(), &content_ref);
3884 let (reached, release) = bounded_read_sync_hook::install(store.root());
3885
3886 let read_store = Arc::clone(&store);
3887 let read_ref = content_ref.clone();
3888 let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
3889 assert!(
3890 recv_blocking(reached).await,
3891 "read must reach the metadata seam"
3892 );
3893 let mut writer = fs::OpenOptions::new()
3894 .write(true)
3895 .truncate(true)
3896 .open(&path)
3897 .unwrap();
3898 writer.write_all(b"abc").unwrap();
3899 writer.flush().unwrap();
3900 release.send(()).unwrap();
3901
3902 let err = read.await.unwrap().unwrap_err();
3903 assert!(matches!(
3904 err,
3905 StorageError::BlobSizeMismatch {
3906 content_ref: ref got,
3907 metadata_bytes: 4,
3908 actual_bytes: 3,
3909 } if got == &content_ref
3910 ));
3911 }
3912
3913 #[cfg(unix)]
3914 #[tokio::test]
3915 async fn bounded_verified_get_keeps_the_opened_inode_when_the_path_is_replaced() {
3916 let (_dir, store) = store(0);
3917 let store = Arc::new(store);
3918 let original = b"original".to_vec();
3919 let replacement = b"replaced".to_vec();
3920 let content_ref = store.put(original.clone()).await.unwrap();
3921 let path = shard_path(store.root(), &content_ref);
3922 let moved_path = path.with_extension("opened-inode");
3923 let max_bytes = original.len() as u64;
3924 let (reached, release) = bounded_read_sync_hook::install(store.root());
3925
3926 let read_store = Arc::clone(&store);
3927 let read_ref = content_ref.clone();
3928 let read =
3929 tokio::spawn(
3930 async move { read_store.get_bounded_verified(&read_ref, max_bytes).await },
3931 );
3932 assert!(
3933 recv_blocking(reached).await,
3934 "read must reach the metadata seam"
3935 );
3936 fs::rename(&path, &moved_path).unwrap();
3937 fs::write(&path, replacement).unwrap();
3938 release.send(()).unwrap();
3939
3940 assert_eq!(read.await.unwrap().unwrap(), original);
3941 }
3942
3943 #[cfg(unix)]
3944 #[tokio::test]
3945 async fn bounded_verified_get_refuses_a_symlink_leaf() {
3946 use std::os::unix::fs::symlink;
3947
3948 let dir = tempfile::tempdir().unwrap();
3949 let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
3950 let outside_bytes = b"outside but digest matching".to_vec();
3951 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
3952 let outside = dir.path().join("outside");
3953 fs::write(&outside, &outside_bytes).unwrap();
3954 let leaf = shard_path(store.root(), &content_ref);
3955 fs::create_dir_all(leaf.parent().unwrap()).unwrap();
3956 symlink(&outside, &leaf).unwrap();
3957
3958 let err = store
3959 .get_bounded_verified(&content_ref, outside_bytes.len() as u64)
3960 .await
3961 .unwrap_err();
3962 assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
3963 }
3964
3965 #[cfg(unix)]
3966 #[tokio::test]
3967 async fn bounded_verified_get_refuses_a_symlinked_shard_component() {
3968 use std::os::unix::fs::symlink;
3969
3970 let dir = tempfile::tempdir().unwrap();
3971 let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
3972 let outside_bytes = b"outside through shard link".to_vec();
3973 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
3974 let hex = content_ref.as_str();
3975 let outside_shard1 = dir.path().join("outside-shard1");
3976 let outside_shard2 = outside_shard1.join(&hex[2..4]);
3977 fs::create_dir_all(&outside_shard2).unwrap();
3978 fs::write(outside_shard2.join(hex), &outside_bytes).unwrap();
3979 symlink(&outside_shard1, store.root().join(&hex[0..2])).unwrap();
3980
3981 let err = store
3982 .get_bounded_verified(&content_ref, outside_bytes.len() as u64)
3983 .await
3984 .unwrap_err();
3985 assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
3986 }
3987
3988 #[tokio::test]
3989 async fn put_content_ref_matches_blake3_digest() {
3990 let (_dir, store) = store(0);
3991 let bytes = b"digest check".to_vec();
3992 let content_ref = store.put(bytes.clone()).await.unwrap();
3993 let expected = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
3994 assert_eq!(content_ref, expected);
3995 }
3996
3997 #[tokio::test]
3998 async fn put_dedups_identical_content() {
3999 let (_dir, store) = store(0);
4000 let bytes = b"same bytes twice".to_vec();
4001 let first = store.put(bytes.clone()).await.unwrap();
4002 let second = store.put(bytes.clone()).await.unwrap();
4003 assert_eq!(first, second);
4004 assert_eq!(
4005 store
4006 .get_bounded_verified(&first, bytes.len() as u64)
4007 .await
4008 .unwrap(),
4009 bytes
4010 );
4011 }
4012
4013 #[tokio::test]
4014 async fn exists_reflects_put_and_delete() {
4015 let (_dir, store) = store(0);
4016 let bytes = b"exists check".to_vec();
4017 let content_ref = store.put(bytes).await.unwrap();
4018 assert!(store.exists(&content_ref).await.unwrap());
4019
4020 assert!(store.delete(&content_ref).await.unwrap());
4021 assert!(!store.exists(&content_ref).await.unwrap());
4022 }
4023
4024 #[tokio::test]
4025 async fn delete_missing_content_ref_returns_false() {
4026 let (_dir, store) = store(0);
4027 let missing = ContentRef::from_hex("f".repeat(64)).unwrap();
4028 assert!(!store.delete(&missing).await.unwrap());
4029 }
4030
4031 #[tokio::test]
4032 async fn size_reports_byte_length_for_a_present_object() {
4033 let (_dir, store) = store(0);
4034 let bytes = b"size check".to_vec();
4035 let content_ref = store.put(bytes.clone()).await.unwrap();
4036 assert_eq!(
4037 store.size(&content_ref).await.unwrap(),
4038 Some(bytes.len() as u64)
4039 );
4040 }
4041
4042 #[tokio::test]
4043 async fn size_returns_none_for_an_absent_object() {
4044 let (_dir, store) = store(0);
4045 let missing = ContentRef::from_hex("9".repeat(64)).unwrap();
4046 assert_eq!(store.size(&missing).await.unwrap(), None);
4047 }
4048
4049 #[tokio::test]
4050 async fn bounded_get_missing_content_ref_returns_not_found() {
4051 let (_dir, store) = store(0);
4052 let missing = ContentRef::from_hex("e".repeat(64)).unwrap();
4053 let err = store
4054 .get_bounded_verified(&missing, MAX_BLOB_WHOLE_BYTES)
4055 .await
4056 .unwrap_err();
4057 assert!(matches!(err, StorageError::NotFound { .. }));
4058 }
4059
4060 #[tokio::test]
4061 async fn put_refuses_below_free_space_floor() {
4062 let (_dir, store) = store(u64::MAX);
4065 let err = store.put(b"too big a floor".to_vec()).await.unwrap_err();
4066 match err {
4067 StorageError::CapacityFloor {
4068 floor_bytes,
4069 available_bytes,
4070 ..
4071 } => {
4072 assert_eq!(floor_bytes, u64::MAX);
4073 assert!(available_bytes < u64::MAX);
4074 }
4075 other => panic!("expected CapacityFloor, got {other:?}"),
4076 }
4077 }
4078
4079 #[tokio::test]
4080 async fn capacity_floor_error_names_the_floor_and_volume() {
4081 let (_dir, store) = store(u64::MAX);
4082 let err = store.put(b"x".to_vec()).await.unwrap_err();
4083 let msg = err.to_string();
4084 assert!(
4085 msg.contains(&u64::MAX.to_string()),
4086 "must name the floor: {msg}"
4087 );
4088 assert!(msg.contains("Blob"), "must name the capability: {msg}");
4089 }
4090
4091 #[test]
4092 fn crosses_floor_is_write_size_aware_at_the_exact_boundary() {
4093 assert!(crosses_floor(101, 2, 100));
4098 assert!(!crosses_floor(101, 1, 100));
4099 }
4100
4101 #[test]
4102 fn crosses_floor_accepts_a_write_that_lands_exactly_on_the_floor() {
4103 assert!(!crosses_floor(100, 0, 100));
4104 }
4105
4106 #[test]
4107 fn crosses_floor_rejects_a_write_that_lands_one_byte_under_the_floor() {
4108 assert!(crosses_floor(100, 1, 100));
4109 }
4110
4111 #[test]
4112 fn crosses_floor_saturates_instead_of_underflowing_when_write_exceeds_available() {
4113 assert!(crosses_floor(10, 100, 50));
4114 assert!(!crosses_floor(10, 100, 0));
4120 }
4121
4122 #[test]
4123 fn put_refuses_a_write_that_would_cross_the_floor_even_though_available_alone_clears_it() {
4124 let dir = tempfile::tempdir().unwrap();
4125 let root = dir.path().join("blobs");
4126 fs::create_dir_all(&root).unwrap();
4127
4128 let err = put_blocking_with_space_probe(&root, 100, vec![7u8; 2], |_| Ok(101)).unwrap_err();
4135 assert!(
4136 matches!(err, StorageError::CapacityFloor { .. }),
4137 "a write-size-aware floor check must reject a write that pushes the volume \
4138 below the floor even though available space alone still clears it: {err:?}"
4139 );
4140 }
4141
4142 #[test]
4143 fn a_later_put_checks_a_fresh_capacity_snapshot() {
4144 let dir = tempfile::tempdir().unwrap();
4145 let root = dir.path().join("blobs");
4146 fs::create_dir_all(&root).unwrap();
4147
4148 let first = put_blocking_with_space_probe(&root, 100, vec![1u8; 2], |_| Ok(102));
4156 let second = put_blocking_with_space_probe(&root, 100, vec![2u8; 2], |_| Ok(101));
4157
4158 assert!(
4159 first.is_ok(),
4160 "the first put may land on the floor: {first:?}"
4161 );
4162 assert!(
4163 matches!(second, Err(StorageError::CapacityFloor { .. })),
4164 "the later put must use its lower capacity snapshot: {second:?}"
4165 );
4166 }
4167
4168 #[tokio::test]
4169 async fn concurrent_puts_from_two_independently_constructed_stores_share_the_root_lock() {
4170 let dir = tempfile::tempdir().unwrap();
4223 let root = dir.path().join("blobs");
4224 fs::create_dir_all(&root).unwrap();
4225 let canonical_root = root.canonicalize().unwrap();
4226
4227 let store_a = std::sync::Arc::new(FsBlobStore::new(root.clone(), 0).unwrap());
4230 let store_b = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
4231
4232 let (a_reached, a_release, _a_done) = sync_hook::install(&canonical_root);
4233 let a = {
4234 let store_a = store_a.clone();
4235 tokio::spawn(async move { store_a.put(b"store_a payload".to_vec()).await })
4236 };
4237 assert!(
4238 recv_blocking(a_reached).await,
4239 "store_a's put must reach the sync_hook checkpoint"
4240 );
4241
4242 assert!(
4247 store_b.write_lock.try_lock().is_err(),
4248 "store_b's write_lock was NOT held while store_a's put held its guard -- the two \
4249 independently constructed stores do NOT share one lock"
4250 );
4251
4252 a_release.send(()).unwrap();
4256 let result_a = a.await.unwrap();
4257 assert!(result_a.is_ok(), "store_a's put must succeed: {result_a:?}");
4258
4259 let result_b = store_b.put(b"store_b payload".to_vec()).await;
4262 assert!(result_b.is_ok(), "store_b's put must succeed: {result_b:?}");
4263 }
4264
4265 #[tokio::test]
4266 async fn aborting_the_outer_put_future_does_not_release_the_guard_before_persist_completes() {
4267 let dir = tempfile::tempdir().unwrap();
4287 let root = dir.path().join("blobs");
4288 fs::create_dir_all(&root).unwrap();
4289 let canonical_root = root.canonicalize().unwrap();
4290
4291 let store = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
4292 let (reached, release, done) = sync_hook::install(&canonical_root);
4293 let handle = {
4294 let store = store.clone();
4295 tokio::spawn(async move { store.put(b"cancellation race payload".to_vec()).await })
4296 };
4297
4298 assert!(
4299 recv_blocking(reached).await,
4300 "put must reach the sync_hook checkpoint -- owned guard already moved into the \
4301 closure -- before this test can mean anything"
4302 );
4303
4304 handle.abort();
4305 let abort_result = handle.await;
4306 match &abort_result {
4307 Err(e) if e.is_cancelled() => {}
4308 other => panic!(
4309 "the outer task must actually have been cancelled for this test to be \
4310 meaningful: {other:?}"
4311 ),
4312 }
4313
4314 let shared_lock = write_lock_for_root(&canonical_root).unwrap();
4315 assert!(
4316 shared_lock.try_lock().is_err(),
4317 "the guard must still be held by the detached blocking write immediately after \
4318 the outer future was cancelled -- if this is free, the guard was released with \
4319 the aborted frame instead of moving into the spawn_blocking closure"
4320 );
4321
4322 release.send(()).unwrap();
4327 assert!(
4328 recv_blocking(done).await,
4329 "the detached write must signal completion once it actually persists"
4330 );
4331 assert!(
4332 shared_lock.try_lock().is_ok(),
4333 "the guard must be free once the detached write's completion was observed"
4334 );
4335 }
4336
4337 #[tokio::test]
4345 async fn orphan_sweep_is_disabled_in_both_modes_regardless_of_live_refs() {
4346 let (_dir, store) = store(0);
4347 let blob = store
4348 .put(b"never swept by this API".to_vec())
4349 .await
4350 .unwrap();
4351 let mut live_refs = std::collections::HashSet::new();
4352 live_refs.insert(blob.clone());
4353
4354 for dry_run in [true, false] {
4355 let error = store
4356 .orphan_sweep(&BlobOrphanSweepConfig {
4357 live_refs: live_refs.clone(),
4358 dry_run,
4359 })
4360 .await
4361 .expect_err("caller-snapshot orphan_sweep must be disabled");
4362 assert!(
4363 matches!(error, StorageError::Unsupported { .. }),
4364 "expected typed Unsupported, got {error:?}"
4365 );
4366 }
4367 assert!(store.exists(&blob).await.unwrap());
4368 }
4369
4370 #[tokio::test]
4375 async fn transactional_orphan_sweep_refuses_v20_before_root_or_claim_mutation() {
4376 let dir = tempfile::tempdir().unwrap();
4377 let db_path = dir.path().join("khive.db");
4378 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4379 {
4380 let mut writer = backend.pool().writer().unwrap();
4381 prepare_v20_gc_fixture(writer.conn_mut());
4382 }
4383
4384 let root = dir.path().join("blobs");
4385 let store = Arc::new(
4386 FsBlobStore::new(root, 0)
4387 .unwrap()
4388 .with_orphan_sweep_grace(Duration::ZERO),
4389 );
4390 let bundle = store.put(b"legacy model bundle".to_vec()).await.unwrap();
4391 let network = store.put(b"legacy FANN network".to_vec()).await.unwrap();
4392 let orphan = store.put(b"ordinary old orphan".to_vec()).await.unwrap();
4393 let abandoned_ref = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
4394 {
4395 let writer = backend.pool().writer().unwrap();
4396 writer
4397 .conn()
4398 .execute(
4399 "INSERT INTO entities \
4400 (id, namespace, kind, entity_type, name, tags, created_at, updated_at, \
4401 content_ref) \
4402 VALUES ('legacy-model', 'local', 'artifact', 'moodboard_model', \
4403 'legacy model', '[]', 1, 1, ?1)",
4404 [bundle.as_str()],
4405 )
4406 .unwrap();
4407 writer
4408 .conn()
4409 .execute(
4410 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4411 VALUES ('abandoned-before-compat', ?1, 1)",
4412 [abandoned_ref],
4413 )
4414 .unwrap();
4415 }
4416
4417 let _root_guard = store.write_lock.clone().lock_owned().await;
4420 for dry_run in [true, false] {
4421 let outcome = tokio::time::timeout(
4422 Duration::from_secs(1),
4423 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
4424 )
4425 .await
4426 .expect("V20 refusal must happen before waiting for the held root lock");
4427 let error = outcome.expect_err("V20 transactional sweep must be disabled");
4428 match error {
4429 StorageError::Unsupported {
4430 capability: StorageCapability::Blob,
4431 operation,
4432 message,
4433 } => {
4434 assert_eq!(operation, "transactional_orphan_sweep");
4435 assert!(
4436 message.contains("complete V21 attachment cutover"),
4437 "unexpected compatibility diagnostic: {message}"
4438 );
4439 }
4440 other => panic!("expected typed Unsupported refusal, got {other:?}"),
4441 }
4442 }
4443
4444 let reader = backend.pool().reader().unwrap();
4445 let abandoned: i64 = reader
4446 .conn()
4447 .query_row(
4448 "SELECT COUNT(*) FROM blob_gc_claims \
4449 WHERE root_key = 'abandoned-before-compat' AND content_ref = ?1",
4450 [abandoned_ref],
4451 |row| row.get(0),
4452 )
4453 .unwrap();
4454 assert_eq!(abandoned, 1, "V20 refusal must not clean abandoned claims");
4455 drop(reader);
4456 assert!(store.exists(&bundle).await.unwrap());
4457 assert!(store.exists(&network).await.unwrap());
4458 assert!(store.exists(&orphan).await.unwrap());
4459 }
4460
4461 #[tokio::test]
4466 async fn transactional_orphan_sweep_refuses_incomplete_v21_marker_without_mutation() {
4467 let dir = tempfile::tempdir().unwrap();
4468 let db_path = dir.path().join("khive.db");
4469 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4470 let abandoned_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
4471 {
4472 let mut writer = backend.pool().writer().unwrap();
4473 prepare_completed_v21_gc_fixture(writer.conn_mut());
4474 writer
4475 .conn_mut()
4476 .execute(
4477 "UPDATE attachment_cutover_state \
4478 SET state = 'incomplete', completed_at = NULL \
4479 WHERE singleton = 1",
4480 [],
4481 )
4482 .unwrap();
4483 writer
4484 .conn_mut()
4485 .execute(
4486 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4487 VALUES ('abandoned-incomplete-v21', ?1, 1)",
4488 [abandoned_ref],
4489 )
4490 .unwrap();
4491 }
4492
4493 let store = Arc::new(
4494 FsBlobStore::new(dir.path().join("blobs"), 0)
4495 .unwrap()
4496 .with_orphan_sweep_grace(Duration::ZERO),
4497 );
4498 let orphan = store.put(b"incomplete V21 orphan".to_vec()).await.unwrap();
4499 let _root_guard = store.write_lock.clone().lock_owned().await;
4500
4501 for dry_run in [true, false] {
4502 let outcome = tokio::time::timeout(
4503 Duration::from_secs(1),
4504 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
4505 )
4506 .await
4507 .expect("incomplete V21 must refuse before waiting for the root lock");
4508 assert!(
4509 matches!(outcome, Err(StorageError::Unsupported { .. })),
4510 "incomplete V21 must return typed Unsupported: {outcome:?}"
4511 );
4512 }
4513
4514 let remaining: i64 = backend
4515 .pool()
4516 .reader()
4517 .unwrap()
4518 .conn()
4519 .query_row(
4520 "SELECT COUNT(*) FROM blob_gc_claims \
4521 WHERE root_key = 'abandoned-incomplete-v21' AND content_ref = ?1",
4522 [abandoned_ref],
4523 |row| row.get(0),
4524 )
4525 .unwrap();
4526 assert_eq!(remaining, 1, "refusal must not recover abandoned claims");
4527 assert!(store.exists(&orphan).await.unwrap());
4528 }
4529
4530 #[tokio::test]
4537 async fn both_sweep_apis_refuse_v20_and_incomplete_v21_epochs_in_both_modes() {
4538 for incomplete_v21 in [false, true] {
4539 let dir = tempfile::tempdir().unwrap();
4540 let db_path = dir.path().join("khive.db");
4541 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4542 {
4543 let mut writer = backend.pool().writer().unwrap();
4544 if incomplete_v21 {
4545 prepare_completed_v21_gc_fixture(writer.conn_mut());
4546 writer
4547 .conn_mut()
4548 .execute(
4549 "UPDATE attachment_cutover_state \
4550 SET state = 'incomplete', completed_at = NULL \
4551 WHERE singleton = 1",
4552 [],
4553 )
4554 .unwrap();
4555 } else {
4556 prepare_v20_gc_fixture(writer.conn_mut());
4557 }
4558 }
4559
4560 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
4561 .unwrap()
4562 .with_orphan_sweep_grace(Duration::ZERO);
4563 let orphan = store
4564 .put(format!("both-apis orphan (incomplete_v21={incomplete_v21})").into_bytes())
4565 .await
4566 .unwrap();
4567
4568 let known_claim_ref = "f".repeat(64);
4572 {
4573 let writer = backend.pool().writer().unwrap();
4574 writer
4575 .conn()
4576 .execute(
4577 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4578 VALUES ('both-api-known-claim', ?1, 1)",
4579 [known_claim_ref.as_str()],
4580 )
4581 .unwrap();
4582 }
4583 let assert_known_claim_unchanged = |arm: &str| {
4584 let remaining: i64 = backend
4585 .pool()
4586 .reader()
4587 .unwrap()
4588 .conn()
4589 .query_row(
4590 "SELECT COUNT(*) FROM blob_gc_claims \
4591 WHERE root_key = 'both-api-known-claim' AND content_ref = ?1",
4592 [known_claim_ref.as_str()],
4593 |row| row.get(0),
4594 )
4595 .unwrap();
4596 assert_eq!(
4597 remaining, 1,
4598 "incomplete_v21={incomplete_v21} arm={arm}: refusal must not mutate \
4599 an existing claim"
4600 );
4601 };
4602
4603 for dry_run in [true, false] {
4604 let snapshot_error = store
4605 .orphan_sweep(&BlobOrphanSweepConfig {
4606 live_refs: std::collections::HashSet::new(),
4607 dry_run,
4608 })
4609 .await
4610 .expect_err("orphan_sweep must refuse regardless of epoch");
4611 assert!(
4612 matches!(snapshot_error, StorageError::Unsupported { .. }),
4613 "incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
4614 from orphan_sweep, got {snapshot_error:?}"
4615 );
4616 assert_known_claim_unchanged(&format!("orphan_sweep dry_run={dry_run}"));
4617
4618 let transactional_error = store
4619 .transactional_orphan_sweep(backend.sql().as_ref(), dry_run)
4620 .await
4621 .expect_err("transactional_orphan_sweep must refuse this epoch");
4622 assert!(
4623 matches!(transactional_error, StorageError::Unsupported { .. }),
4624 "incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
4625 from transactional_orphan_sweep, got {transactional_error:?}"
4626 );
4627 assert_known_claim_unchanged(&format!(
4628 "transactional_orphan_sweep dry_run={dry_run}"
4629 ));
4630 }
4631
4632 assert!(
4633 store.exists(&orphan).await.unwrap(),
4634 "incomplete_v21={incomplete_v21}: a refused sweep must not delete anything"
4635 );
4636 }
4637 }
4638
4639 #[tokio::test]
4659 async fn transactional_orphan_sweep_recheck_refuses_before_root_lock_when_epoch_regresses_after_db_ownership(
4660 ) {
4661 let dir = tempfile::tempdir().unwrap();
4662 let db_path = dir.path().join("khive.db");
4663 let backend = Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4664 {
4665 let mut writer = backend.pool().writer().unwrap();
4666 prepare_completed_v21_gc_fixture(writer.conn_mut());
4667 }
4668 assert!(blob_gc_fencing_complete(backend.sql().as_ref())
4669 .await
4670 .unwrap());
4671
4672 let blob_root = dir.path().join("blobs");
4673 let store = Arc::new(
4674 FsBlobStore::new(blob_root.clone(), 0)
4675 .unwrap()
4676 .with_orphan_sweep_grace(Duration::ZERO),
4677 );
4678 let orphan = store
4679 .put(b"epoch regressed after database ownership".to_vec())
4680 .await
4681 .unwrap();
4682
4683 let canonical_root = blob_root.canonicalize().unwrap();
4686 let _root_write_guard = acquire_root_write_lock(&canonical_root).unwrap();
4687
4688 let canonical_db_path = backend.sql().database_path();
4692 let (reached, release) = db_ownership_sync_hook::install(canonical_db_path.as_deref());
4693 let sweep_store = store.clone();
4694 let sweep_backend = backend.clone();
4695 let handle = tokio::spawn(async move {
4696 sweep_store
4697 .transactional_orphan_sweep(sweep_backend.sql().as_ref(), false)
4698 .await
4699 });
4700
4701 let reached_signal = tokio::time::timeout(Duration::from_secs(1), recv_blocking(reached))
4707 .await
4708 .expect("the sweep must reach database ownership before this test's timeout");
4709 assert!(reached_signal, "hook sender was dropped before signaling");
4710 {
4711 let writer = backend.pool().writer().unwrap();
4712 writer
4713 .conn()
4714 .execute("DELETE FROM _schema_migrations WHERE version = 21", [])
4715 .unwrap();
4716 }
4717 release.send(()).unwrap();
4718
4719 let outcome = tokio::time::timeout(Duration::from_secs(1), handle)
4720 .await
4721 .expect(
4722 "the recheck must refuse before ever waiting on the externally held root lock -- \
4723 under the old (pre-fix) ordering this join times out instead, because the \
4724 sweep blocks acquiring the OS-level root lock held above",
4725 )
4726 .unwrap();
4727 assert!(
4728 matches!(outcome, Err(StorageError::Unsupported { .. })),
4729 "expected the regressed epoch to be caught immediately after database ownership: \
4730 {outcome:?}"
4731 );
4732 assert!(store.exists(&orphan).await.unwrap());
4733 }
4734
4735 #[tokio::test]
4740 async fn transactional_orphan_sweep_accepts_completed_v21_attachment_liveness() {
4741 let dir = tempfile::tempdir().unwrap();
4742 let db_path = dir.path().join("khive.db");
4743 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4744 {
4745 let mut writer = backend.pool().writer().unwrap();
4746 prepare_completed_v21_gc_fixture(writer.conn_mut());
4747 }
4748
4749 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
4750 .unwrap()
4751 .with_orphan_sweep_grace(Duration::ZERO);
4752 let bundle = store.put(b"V21 model bundle".to_vec()).await.unwrap();
4753 let network = store.put(b"V21 FANN network".to_vec()).await.unwrap();
4754 let orphan = store.put(b"V21 true orphan".to_vec()).await.unwrap();
4755 {
4756 let writer = backend.pool().writer().unwrap();
4757 writer
4758 .conn()
4759 .execute(
4760 "INSERT INTO entities \
4761 (id, namespace, kind, entity_type, name, tags, created_at, updated_at) \
4762 VALUES ('model', 'local', 'artifact', 'moodboard_model', \
4763 'model', '[]', 1, 1)",
4764 [],
4765 )
4766 .unwrap();
4767 writer
4768 .conn()
4769 .execute(
4770 "INSERT INTO attachments \
4771 (record_uuid, substrate, role, content_ref, created_at) \
4772 VALUES ('model', 'entity', 'content', ?1, 1), \
4773 ('model', 'entity', 'fann-network', ?2, 1)",
4774 rusqlite::params![bundle.as_str(), network.as_str()],
4775 )
4776 .unwrap();
4777 }
4778
4779 let dry_run = store
4780 .transactional_orphan_sweep(backend.sql().as_ref(), true)
4781 .await
4782 .expect("completed V21 dry run must be supported");
4783 assert_eq!(dry_run.would_delete, 1);
4784 assert_eq!(dry_run.deleted, 0);
4785
4786 let result = store
4787 .transactional_orphan_sweep(backend.sql().as_ref(), false)
4788 .await
4789 .expect("completed V21 destructive sweep must be supported");
4790 assert_eq!(result.deleted, 1);
4791 assert!(store.exists(&bundle).await.unwrap());
4792 assert!(store.exists(&network).await.unwrap());
4793 assert!(!store.exists(&orphan).await.unwrap());
4794 }
4795
4796 #[tokio::test]
4805 async fn transactional_orphan_sweep_refuses_without_the_blob_gc_claims_migration() {
4806 let dir = tempfile::tempdir().unwrap();
4807 let db_path = dir.path().join("khive.db");
4808 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4809 backend.entities().unwrap();
4810 {
4811 let reader = backend.pool().reader().unwrap();
4812 let present: bool = reader
4813 .conn()
4814 .query_row(
4815 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = 'table' \
4816 AND name = 'blob_gc_claims'",
4817 [],
4818 |row| row.get(0),
4819 )
4820 .unwrap();
4821 assert!(
4822 !present,
4823 "this test's premise requires blob_gc_claims to be absent"
4824 );
4825 }
4826
4827 let root = dir.path().join("blobs");
4828 let store = std::sync::Arc::new(
4829 FsBlobStore::new(root.clone(), 0)
4830 .unwrap()
4831 .with_orphan_sweep_grace(Duration::ZERO),
4832 );
4833 let orphan = store.put(b"direct-backend orphan".to_vec()).await.unwrap();
4834
4835 let sql = backend.sql();
4836 let error = store
4837 .transactional_orphan_sweep(sql.as_ref(), false)
4838 .await
4839 .expect_err("sweep must refuse a backend without the blob_gc_claims fencing set");
4840 assert!(
4841 matches!(error, StorageError::Unsupported { .. }),
4842 "expected StorageError::Unsupported, got {error:?}"
4843 );
4844 assert!(
4845 store.exists(&orphan).await.unwrap(),
4846 "a refused sweep must not have deleted anything"
4847 );
4848 }
4849
4850 #[tokio::test]
4851 async fn transactional_orphan_sweep_refuses_an_incomplete_cutover_marker() {
4852 let dir = tempfile::tempdir().unwrap();
4853 let db_path = dir.path().join("khive.db");
4854 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4855 {
4856 let mut writer = backend.pool().writer().unwrap();
4857 prepare_completed_v21_gc_fixture(writer.conn_mut());
4858 writer
4859 .conn_mut()
4860 .execute_batch(
4861 "UPDATE attachment_cutover_state \
4862 SET state = 'incomplete', completed_at = NULL WHERE singleton = 1; \
4863 DELETE FROM _schema_migrations WHERE version = 21;",
4864 )
4865 .unwrap();
4866 }
4867
4868 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
4869 .unwrap()
4870 .with_orphan_sweep_grace(Duration::ZERO);
4871 let orphan = store
4872 .put(b"incomplete-cutover orphan".to_vec())
4873 .await
4874 .unwrap();
4875 let error = store
4876 .transactional_orphan_sweep(backend.sql().as_ref(), false)
4877 .await
4878 .expect_err("sweep must refuse every durable incomplete marker");
4879 assert!(matches!(error, StorageError::Unsupported { .. }));
4880 assert!(
4881 store.exists(&orphan).await.unwrap(),
4882 "refused incomplete-state sweep must preserve every blob"
4883 );
4884 }
4885
4886 #[tokio::test]
4891 async fn transactional_orphan_sweep_refuses_with_incomplete_fencing_triggers() {
4892 let dir = tempfile::tempdir().unwrap();
4893 let db_path = dir.path().join("khive.db");
4894 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4895 {
4896 let mut writer = backend.pool().writer().unwrap();
4897 prepare_completed_v21_gc_fixture(writer.conn_mut());
4898 writer
4899 .conn_mut()
4900 .execute_batch("DROP TRIGGER attachments_reject_claimed_blob_update")
4901 .unwrap();
4902 writer
4903 .conn_mut()
4904 .execute(
4905 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4906 VALUES ('abandoned-partial-fence', \
4907 'dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd', \
4908 1)",
4909 [],
4910 )
4911 .unwrap();
4912 }
4913
4914 let root = dir.path().join("blobs");
4915 let store = std::sync::Arc::new(
4916 FsBlobStore::new(root.clone(), 0)
4917 .unwrap()
4918 .with_orphan_sweep_grace(Duration::ZERO),
4919 );
4920 let orphan = store.put(b"partial-fence orphan".to_vec()).await.unwrap();
4921
4922 let sql = backend.sql();
4923 let _root_guard = store.write_lock.clone().lock_owned().await;
4924 let error = tokio::time::timeout(
4925 Duration::from_secs(1),
4926 store.transactional_orphan_sweep(sql.as_ref(), false),
4927 )
4928 .await
4929 .expect("an incomplete V21 fence must refuse before the root wait")
4930 .expect_err("sweep must refuse when any V21 fencing trigger is missing");
4931 assert!(
4932 matches!(error, StorageError::Unsupported { .. }),
4933 "expected StorageError::Unsupported, got {error:?}"
4934 );
4935 assert!(
4936 store.exists(&orphan).await.unwrap(),
4937 "a refused sweep must not have deleted anything"
4938 );
4939 let remaining: i64 = backend
4940 .pool()
4941 .reader()
4942 .unwrap()
4943 .conn()
4944 .query_row(
4945 "SELECT COUNT(*) FROM blob_gc_claims \
4946 WHERE root_key = 'abandoned-partial-fence'",
4947 [],
4948 |row| row.get(0),
4949 )
4950 .unwrap();
4951 assert_eq!(remaining, 1, "a refused sweep must not recover claims");
4952 }
4953
4954 #[tokio::test]
4960 async fn transactional_orphan_sweep_refuses_same_named_noop_fencing_triggers() {
4961 let dir = tempfile::tempdir().unwrap();
4962 let db_path = dir.path().join("khive.db");
4963 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4964 {
4965 let mut writer = backend.pool().writer().unwrap();
4966 prepare_completed_v21_gc_fixture(writer.conn_mut());
4967 writer
4968 .conn_mut()
4969 .execute_batch(
4970 "DROP TRIGGER attachments_reject_claimed_blob_insert; \
4971 DROP TRIGGER attachments_reject_claimed_blob_update; \
4972 CREATE TRIGGER attachments_reject_claimed_blob_insert \
4973 BEFORE INSERT ON attachments BEGIN SELECT 0; END; \
4974 CREATE TRIGGER attachments_reject_claimed_blob_update \
4975 BEFORE UPDATE OF content_ref ON attachments \
4976 BEGIN SELECT 0; END;",
4977 )
4978 .unwrap();
4979 }
4980
4981 let root = dir.path().join("blobs");
4982 let store = std::sync::Arc::new(
4983 FsBlobStore::new(root.clone(), 0)
4984 .unwrap()
4985 .with_orphan_sweep_grace(Duration::ZERO),
4986 );
4987 let orphan = store.put(b"noop-trigger orphan".to_vec()).await.unwrap();
4988
4989 let sql = backend.sql();
4990 let error = store
4991 .transactional_orphan_sweep(sql.as_ref(), false)
4992 .await
4993 .expect_err("sweep must refuse when the fencing triggers are same-named no-ops");
4994 assert!(
4995 matches!(error, StorageError::Unsupported { .. }),
4996 "expected StorageError::Unsupported, got {error:?}"
4997 );
4998 assert!(
4999 store.exists(&orphan).await.unwrap(),
5000 "a refused sweep must not have deleted anything"
5001 );
5002
5003 let reader = backend.pool().reader().unwrap();
5005 let leftovers: i64 = reader
5006 .conn()
5007 .query_row(
5008 "SELECT (SELECT COUNT(*) FROM blob_gc_claims \
5009 WHERE root_key GLOB '__fence_probe-*') \
5010 + (SELECT COUNT(*) FROM attachments \
5011 WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
5012 [],
5013 |row| row.get(0),
5014 )
5015 .unwrap();
5016 assert_eq!(leftovers, 0, "fence probe rows must not survive the probe");
5017 }
5018
5019 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5020 async fn fence_probe_refuses_id_collision_and_preserves_the_colliding_attachment() {
5021 let dir = tempfile::tempdir().unwrap();
5022 let db_path = dir.path().join("khive.db");
5023 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5024 {
5025 let mut writer = backend.pool().writer().unwrap();
5026 prepare_completed_v21_gc_fixture(writer.conn_mut());
5027 writer
5028 .conn_mut()
5029 .execute(
5030 "INSERT INTO attachments \
5031 (record_uuid, substrate, role, content_ref, media_type, created_at) \
5032 VALUES ('victim-id', 'entity', 'content', \
5033 '2222222222222222222222222222222222222222222222222222222222222222', \
5034 'application/test', 7)",
5035 [],
5036 )
5037 .unwrap();
5038 }
5039
5040 let sql = backend.sql();
5041 let error = super::blob_gc_fence_probe_with_ids(
5042 sql.as_ref(),
5043 "victim-id".to_string(),
5044 "victim-update-id".to_string(),
5045 "victim-insert2-id".to_string(),
5046 "victim-update2-id".to_string(),
5047 "victim-claim-key".to_string(),
5048 )
5049 .await
5050 .expect_err("the probe must refuse when an id it would delete already names a row");
5051 assert!(
5052 matches!(error, StorageError::Unsupported { .. }),
5053 "expected StorageError::Unsupported, got {error:?}"
5054 );
5055
5056 let reader = backend.pool().reader().unwrap();
5057 let (media_type, created_at): (String, i64) = reader
5058 .conn()
5059 .query_row(
5060 "SELECT media_type, created_at FROM attachments \
5061 WHERE record_uuid = 'victim-id' AND role = 'content'",
5062 [],
5063 |row| Ok((row.get(0)?, row.get(1)?)),
5064 )
5065 .expect("the colliding attachment must survive the refused probe untouched");
5066 assert_eq!(media_type, "application/test");
5067 assert_eq!(created_at, 7);
5068 }
5069
5070 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5071 async fn fence_probe_does_not_touch_an_unrelated_retained_entity_sequence() {
5072 let dir = tempfile::tempdir().unwrap();
5073 let db_path = dir.path().join("khive.db");
5074 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5075 {
5076 let mut writer = backend.pool().writer().unwrap();
5077 prepare_completed_v21_gc_fixture(writer.conn_mut());
5078 writer
5082 .conn_mut()
5083 .execute(
5084 "INSERT INTO entities \
5085 (id, namespace, kind, name, tags, created_at, updated_at) \
5086 VALUES ('retained-id', 'local', 'document', 'gone entity', '[]', 7, 7)",
5087 [],
5088 )
5089 .unwrap();
5090 writer
5091 .conn_mut()
5092 .execute("DELETE FROM entities WHERE id = 'retained-id'", [])
5093 .unwrap();
5094 let retained: i64 = writer
5095 .conn_mut()
5096 .query_row(
5097 "SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
5098 [],
5099 |row| row.get(0),
5100 )
5101 .unwrap();
5102 assert_eq!(retained, 1, "fixture requires a retained-only ledger row");
5103 }
5104
5105 let sql = backend.sql();
5106 super::blob_gc_fence_probe_with_ids(
5107 sql.as_ref(),
5108 "retained-id".to_string(),
5109 "retained-update-id".to_string(),
5110 "retained-insert2-id".to_string(),
5111 "retained-update2-id".to_string(),
5112 "retained-claim-key".to_string(),
5113 )
5114 .await
5115 .expect("attachment probe has no reason to mutate an entity sequence row");
5116
5117 let reader = backend.pool().reader().unwrap();
5118 let survivors: i64 = reader
5119 .conn()
5120 .query_row(
5121 "SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
5122 [],
5123 |row| row.get(0),
5124 )
5125 .unwrap();
5126 assert_eq!(
5127 survivors, 1,
5128 "the retained entity ledger row must survive the attachment probe"
5129 );
5130 }
5131
5132 fn nul_embedded_canonical_ref() -> String {
5133 let mut polluted = "a".repeat(64);
5134 polluted.push('\0');
5135 polluted.push_str("zz");
5136 polluted
5137 }
5138
5139 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5144 async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref() {
5145 let dir = tempfile::tempdir().unwrap();
5146 let db_path = dir.path().join("khive.db");
5147 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5148 {
5149 let mut writer = backend.pool().writer().unwrap();
5150 prepare_completed_v21_gc_fixture(writer.conn_mut());
5151 writer
5152 .conn_mut()
5153 .execute(
5154 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5155 VALUES ('nul-claim-key', ?1, 0)",
5156 rusqlite::params![nul_embedded_canonical_ref()],
5157 )
5158 .unwrap();
5159 }
5160
5161 let sql = backend.sql();
5162 let error = super::validate_blob_gc_evidence(sql.as_ref())
5163 .await
5164 .expect_err("a NUL-embedded claim ref must refuse the sweep");
5165 assert!(
5166 error.to_string().contains("blob_gc_claims"),
5167 "expected the claims-table refusal, got {error:?}"
5168 );
5169 }
5170
5171 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5176 async fn blob_gc_evidence_rejects_a_nul_embedded_attachment_ref() {
5177 let dir = tempfile::tempdir().unwrap();
5178 let db_path = dir.path().join("khive.db");
5179 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5180 {
5181 let mut writer = backend.pool().writer().unwrap();
5182 prepare_completed_v21_gc_fixture(writer.conn_mut());
5183 writer
5188 .conn_mut()
5189 .execute_batch("PRAGMA ignore_check_constraints = ON")
5190 .unwrap();
5191 writer
5192 .conn_mut()
5193 .execute(
5194 "INSERT INTO attachments \
5195 (record_uuid, substrate, role, content_ref, created_at) \
5196 VALUES ('nul-attachment-id', 'entity', 'content', ?1, 0)",
5197 rusqlite::params![nul_embedded_canonical_ref()],
5198 )
5199 .expect("ignore_check_constraints must allow the corrupt row to insert");
5200 writer
5201 .conn_mut()
5202 .execute_batch("PRAGMA ignore_check_constraints = OFF")
5203 .unwrap();
5204 }
5205
5206 let sql = backend.sql();
5207 let error = super::validate_blob_gc_evidence(sql.as_ref())
5208 .await
5209 .expect_err("a NUL-embedded attachment ref must refuse the sweep");
5210 assert!(
5211 error.to_string().contains("attachments"),
5212 "expected the attachments-table refusal, got {error:?}"
5213 );
5214 }
5215
5216 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5220 async fn fence_probe_refuses_a_digest_restricted_trigger_rewrite() {
5221 let dir = tempfile::tempdir().unwrap();
5222 let db_path = dir.path().join("khive.db");
5223 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5224 {
5225 let mut writer = backend.pool().writer().unwrap();
5226 prepare_completed_v21_gc_fixture(writer.conn_mut());
5227 }
5228
5229 let sql = backend.sql();
5230 super::blob_gc_fence_probe(sql.as_ref())
5231 .await
5232 .expect("the healthy fence must pass all four probe arms");
5233
5234 {
5235 let mut writer = backend.pool().writer().unwrap();
5236 writer
5237 .conn_mut()
5238 .execute_batch(
5239 "DROP TRIGGER attachments_reject_claimed_blob_insert; \
5240 DROP TRIGGER attachments_reject_claimed_blob_update; \
5241 CREATE TRIGGER attachments_reject_claimed_blob_insert \
5242 BEFORE INSERT ON attachments \
5243 WHEN NEW.content_ref = \
5244 '0000000000000000000000000000000000000000000000000000000000000000' \
5245 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5246 WHERE content_ref = NEW.content_ref) \
5247 BEGIN \
5248 SELECT RAISE(ABORT, \
5249 'content_ref is reserved by an active blob sweep'); \
5250 END; \
5251 CREATE TRIGGER attachments_reject_claimed_blob_update \
5252 BEFORE UPDATE OF content_ref ON attachments \
5253 WHEN NEW.content_ref = \
5254 '0000000000000000000000000000000000000000000000000000000000000000' \
5255 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5256 WHERE content_ref = NEW.content_ref) \
5257 BEGIN \
5258 SELECT RAISE(ABORT, \
5259 'content_ref is reserved by an active blob sweep'); \
5260 END;",
5261 )
5262 .unwrap();
5263 }
5264
5265 let error = super::blob_gc_fence_probe(sql.as_ref())
5266 .await
5267 .expect_err("a sentinel-only fence must fail the second-digest arms");
5268 assert!(
5269 matches!(error, StorageError::Unsupported { .. }),
5270 "expected StorageError::Unsupported, got {error:?}"
5271 );
5272 assert!(
5273 error.to_string().contains("second-digest"),
5274 "the refusal must name a second-digest arm, got {error}"
5275 );
5276 }
5277
5278 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5281 async fn fence_probe_refuses_a_shape_restricted_trigger_rewrite() {
5282 let dir = tempfile::tempdir().unwrap();
5283 let db_path = dir.path().join("khive.db");
5284 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5285 {
5286 let mut writer = backend.pool().writer().unwrap();
5287 prepare_completed_v21_gc_fixture(writer.conn_mut());
5288 writer
5289 .conn_mut()
5290 .execute_batch(
5291 "DROP TRIGGER attachments_reject_claimed_blob_insert; \
5292 DROP TRIGGER attachments_reject_claimed_blob_update; \
5293 CREATE TRIGGER attachments_reject_claimed_blob_insert \
5294 BEFORE INSERT ON attachments \
5295 WHEN NEW.substrate = 'entity' \
5296 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5297 WHERE content_ref = NEW.content_ref) \
5298 BEGIN \
5299 SELECT RAISE(ABORT, \
5300 'content_ref is reserved by an active blob sweep'); \
5301 END; \
5302 CREATE TRIGGER attachments_reject_claimed_blob_update \
5303 BEFORE UPDATE OF content_ref ON attachments \
5304 WHEN NEW.substrate = 'entity' \
5305 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5306 WHERE content_ref = NEW.content_ref) \
5307 BEGIN \
5308 SELECT RAISE(ABORT, \
5309 'content_ref is reserved by an active blob sweep'); \
5310 END;",
5311 )
5312 .unwrap();
5313 }
5314
5315 let sql = backend.sql();
5316 let error = super::blob_gc_fence_probe(sql.as_ref())
5317 .await
5318 .expect_err("an entity-shape-only fence must fail the note-shaped arms");
5319 assert!(
5320 matches!(error, StorageError::Unsupported { .. }),
5321 "expected StorageError::Unsupported, got {error:?}"
5322 );
5323 assert!(
5324 error.to_string().contains("second-digest"),
5325 "the refusal must name a second-digest arm, got {error}"
5326 );
5327 }
5328
5329 fn initialize_utf16le_database(db_path: &std::path::Path) {
5333 let conn = rusqlite::Connection::open(db_path).unwrap();
5334 conn.execute_batch(
5335 "PRAGMA encoding = 'UTF-16le'; \
5336 CREATE TABLE __encoding_pin (x INTEGER); \
5337 DROP TABLE __encoding_pin;",
5338 )
5339 .unwrap();
5340 }
5341
5342 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5346 async fn blob_gc_evidence_accepts_valid_refs_in_a_utf16le_database() {
5347 let dir = tempfile::tempdir().unwrap();
5348 let db_path = dir.path().join("khive.db");
5349 initialize_utf16le_database(&db_path);
5350 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5351 {
5352 let mut writer = backend.pool().writer().unwrap();
5353 let encoding: String = writer
5354 .conn_mut()
5355 .query_row("PRAGMA encoding", [], |row| row.get(0))
5356 .unwrap();
5357 assert_eq!(
5358 encoding, "UTF-16le",
5359 "the fixture database must actually be UTF-16le"
5360 );
5361 prepare_completed_v21_gc_fixture(writer.conn_mut());
5362 writer
5367 .conn_mut()
5368 .execute(
5369 "INSERT INTO attachments \
5370 (record_uuid, substrate, role, content_ref, created_at) \
5371 VALUES ('utf16-valid-attachment', 'entity', 'content', ?1, 0)",
5372 rusqlite::params!["a".repeat(64)],
5373 )
5374 .unwrap();
5375 writer
5376 .conn_mut()
5377 .execute(
5378 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5379 VALUES ('utf16-valid-claim-key', ?1, 0)",
5380 rusqlite::params!["b".repeat(64)],
5381 )
5382 .unwrap();
5383 }
5384
5385 let sql = backend.sql();
5386 super::validate_blob_gc_evidence(sql.as_ref())
5387 .await
5388 .expect("valid canonical refs must pass in a UTF-16LE database");
5389 }
5390
5391 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5394 async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref_in_a_utf16le_database() {
5395 let dir = tempfile::tempdir().unwrap();
5396 let db_path = dir.path().join("khive.db");
5397 initialize_utf16le_database(&db_path);
5398 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5399 {
5400 let mut writer = backend.pool().writer().unwrap();
5401 let encoding: String = writer
5402 .conn_mut()
5403 .query_row("PRAGMA encoding", [], |row| row.get(0))
5404 .unwrap();
5405 assert_eq!(
5406 encoding, "UTF-16le",
5407 "the fixture database must actually be UTF-16le"
5408 );
5409 prepare_completed_v21_gc_fixture(writer.conn_mut());
5410 writer
5411 .conn_mut()
5412 .execute(
5413 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5414 VALUES ('nul-claim-key-utf16', ?1, 0)",
5415 rusqlite::params![nul_embedded_canonical_ref()],
5416 )
5417 .unwrap();
5418 }
5419
5420 let sql = backend.sql();
5421 let error = super::validate_blob_gc_evidence(sql.as_ref())
5422 .await
5423 .expect_err("a NUL-embedded claim ref must refuse the sweep in UTF-16LE too");
5424 assert!(
5425 error.to_string().contains("blob_gc_claims"),
5426 "expected the claims-table refusal, got {error:?}"
5427 );
5428 }
5429
5430 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5431 async fn transactional_orphan_sweep_preserves_put_started_after_liveness_mark() {
5432 let dir = tempfile::tempdir().unwrap();
5433 let db_path = dir.path().join("khive.db");
5434 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5435 {
5436 let mut writer = backend.pool().writer().unwrap();
5437 prepare_completed_v21_gc_fixture(writer.conn_mut());
5438 }
5439 let root = dir.path().join("blobs");
5440 let store = std::sync::Arc::new(
5441 FsBlobStore::new(root.clone(), 0)
5442 .unwrap()
5443 .with_orphan_sweep_grace(Duration::ZERO),
5444 );
5445 let orphan = store.put(b"old orphan".to_vec()).await.unwrap();
5446 let canonical_root = root.canonicalize().unwrap();
5447 let (marked, release, _done) = sync_hook::install(&canonical_root);
5448
5449 let sweep = {
5450 let store = store.clone();
5451 let sql = backend.sql();
5452 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5453 };
5454 assert!(
5455 recv_blocking(marked).await,
5456 "sweep must finish its liveness mark"
5457 );
5458
5459 assert!(
5460 store.write_lock.try_lock().is_err(),
5461 "the sweep must hold the same root lock used by blob writers"
5462 );
5463 let (started_tx, started_rx) = std::sync::mpsc::channel();
5464 let new_ref = {
5465 let root = root.clone();
5466 tokio::task::spawn_blocking(move || {
5467 let _ = started_tx.send(());
5468 put_blocking(&root, 0, b"new concurrent blob".to_vec())
5469 })
5470 };
5471 assert!(recv_blocking(started_rx).await, "blob put must start");
5472
5473 release.send(()).unwrap();
5474 let sweep_result = sweep.await.unwrap().unwrap();
5475 let new_ref = new_ref.await.unwrap().unwrap();
5476
5477 assert_eq!(sweep_result.deleted, 1);
5478 assert!(!store.exists(&orphan).await.unwrap());
5479 assert!(
5480 store.exists(&new_ref).await.unwrap(),
5481 "a blob put started between the liveness mark and physical sweep must survive"
5482 );
5483 }
5484
5485 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5486 async fn transactional_orphan_sweep_releases_sqlite_writer_before_physical_delete() {
5487 let dir = tempfile::tempdir().unwrap();
5488 let db_path = dir.path().join("khive.db");
5489 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5490 {
5491 let mut writer = backend.pool().writer().unwrap();
5492 prepare_completed_v21_gc_fixture(writer.conn_mut());
5493 }
5494 let root = dir.path().join("blobs");
5495 let store = std::sync::Arc::new(
5496 FsBlobStore::new(root.clone(), 0)
5497 .unwrap()
5498 .with_orphan_sweep_grace(Duration::ZERO),
5499 );
5500 let orphan = store
5501 .put(b"claim then delete outside sqlite".to_vec())
5502 .await
5503 .unwrap();
5504 let canonical_root = root.canonicalize().unwrap();
5505 let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
5506
5507 let sweep = {
5508 let store = store.clone();
5509 let sql = backend.sql();
5510 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5511 };
5512 assert!(
5513 recv_blocking(claimed).await,
5514 "sweep must durably claim the orphan before physical deletion"
5515 );
5516 assert!(
5517 store.exists(&orphan).await.unwrap(),
5518 "the test seam must pause before the physical delete"
5519 );
5520 let external_database_lock = fs::OpenOptions::new()
5521 .read(true)
5522 .write(true)
5523 .open(database_gc_lock_path(&db_path))
5524 .unwrap();
5525 assert!(
5526 matches!(
5527 fs4::FileExt::try_lock(&external_database_lock),
5528 Err(fs4::TryLockError::WouldBlock)
5529 ),
5530 "the sweep must retain cross-process database ownership while SQLite's writer is free"
5531 );
5532
5533 let unrelated = rusqlite::Connection::open(&db_path).unwrap();
5538 unrelated.busy_timeout(Duration::from_millis(100)).unwrap();
5539 unrelated
5540 .execute(
5541 "INSERT INTO entities \
5542 (id, namespace, kind, name, tags, created_at, updated_at) \
5543 VALUES ('unrelated-writer', 'local', 'concept', 'unrelated', '[]', 1, 1)",
5544 [],
5545 )
5546 .expect("external filesystem work must not retain SQLite's writer lock");
5547
5548 let claimed_err = unrelated
5552 .execute(
5553 "INSERT INTO attachments \
5554 (record_uuid, substrate, role, content_ref, created_at) \
5555 VALUES ('racing-reference', 'entity', 'content', ?1, 1)",
5556 [orphan.as_str()],
5557 )
5558 .expect_err("a claimed content_ref must fail closed before deletion");
5559 assert!(
5560 claimed_err.to_string().contains("active blob sweep"),
5561 "unexpected claim error: {claimed_err}"
5562 );
5563
5564 release_delete.send(()).unwrap();
5565 let result = sweep.await.unwrap().unwrap();
5566 assert_eq!(result.deleted, 1);
5567 assert!(!store.exists(&orphan).await.unwrap());
5568
5569 let remaining_claims: i64 = unrelated
5570 .query_row(
5571 "SELECT COUNT(*) FROM blob_gc_claims WHERE content_ref = ?1",
5572 [orphan.as_str()],
5573 |row| row.get(0),
5574 )
5575 .unwrap();
5576 assert_eq!(
5577 remaining_claims, 0,
5578 "successful deletion releases the claim"
5579 );
5580 }
5581
5582 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5583 async fn cancelling_sweep_during_delete_keeps_owner_locks_until_blocking_work_finishes() {
5584 let dir = tempfile::tempdir().unwrap();
5585 let db_path = dir.path().join("khive.db");
5586 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5587 {
5588 let mut writer = backend.pool().writer().unwrap();
5589 prepare_completed_v21_gc_fixture(writer.conn_mut());
5590 }
5591 let root = dir.path().join("blobs");
5592 let store = std::sync::Arc::new(
5593 FsBlobStore::new(root.clone(), 0)
5594 .unwrap()
5595 .with_orphan_sweep_grace(Duration::ZERO),
5596 );
5597 let orphan = store.put(b"cancelled sweep orphan".to_vec()).await.unwrap();
5598 let canonical_root = root.canonicalize().unwrap();
5599 let (claimed, release_delete, done) = sync_hook::install(&canonical_root);
5600
5601 let sweep = {
5602 let store = store.clone();
5603 let sql = backend.sql();
5604 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5605 };
5606 assert!(recv_blocking(claimed).await);
5607 sweep.abort();
5608 assert!(sweep.await.unwrap_err().is_cancelled());
5609
5610 let external_root_lock = fs::OpenOptions::new()
5611 .read(true)
5612 .write(true)
5613 .open(root.join(ROOT_WRITE_LOCK_FILE))
5614 .unwrap();
5615 let external_database_lock = fs::OpenOptions::new()
5616 .read(true)
5617 .write(true)
5618 .open(database_gc_lock_path(&db_path))
5619 .unwrap();
5620 assert!(matches!(
5621 fs4::FileExt::try_lock(&external_root_lock),
5622 Err(fs4::TryLockError::WouldBlock)
5623 ));
5624 assert!(matches!(
5625 fs4::FileExt::try_lock(&external_database_lock),
5626 Err(fs4::TryLockError::WouldBlock)
5627 ));
5628
5629 release_delete.send(()).unwrap();
5630 let done_disconnected = tokio::task::spawn_blocking(move || done.recv().is_err())
5631 .await
5632 .unwrap();
5633 assert!(
5634 done_disconnected,
5635 "the cancelled outer task cannot send done"
5636 );
5637 assert!(fs4::FileExt::try_lock(&external_root_lock).is_ok());
5638 assert!(fs4::FileExt::try_lock(&external_database_lock).is_ok());
5639 drop(external_root_lock);
5640 drop(external_database_lock);
5641 assert!(!store.exists(&orphan).await.unwrap());
5642 let stranded_claims: i64 = rusqlite::Connection::open(&db_path)
5643 .unwrap()
5644 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5645 .unwrap();
5646 assert_eq!(
5647 stranded_claims, 1,
5648 "cancellation leaves a fail-closed claim"
5649 );
5650
5651 let recovered = store
5652 .transactional_orphan_sweep(backend.sql().as_ref(), false)
5653 .await
5654 .unwrap();
5655 assert_eq!(recovered.deleted, 0);
5656 let remaining: i64 = rusqlite::Connection::open(&db_path)
5657 .unwrap()
5658 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5659 .unwrap();
5660 assert_eq!(remaining, 0, "the next exclusive owner recovers the claim");
5661 }
5662
5663 #[tokio::test]
5664 async fn transactional_orphan_sweep_recovers_stale_claims_fail_closed() {
5665 let dir = tempfile::tempdir().unwrap();
5666 let db_path = dir.path().join("khive.db");
5667 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5668 {
5669 let mut writer = backend.pool().writer().unwrap();
5670 prepare_completed_v21_gc_fixture(writer.conn_mut());
5671 }
5672 let root = dir.path().join("blobs");
5673 let store = FsBlobStore::new(root.clone(), 0)
5674 .unwrap()
5675 .with_orphan_sweep_grace(Duration::from_secs(60));
5676 let bytes = b"republished after a crashed claim".to_vec();
5677 let content_ref = store.put(bytes.clone()).await.unwrap();
5678 let canonical_root = root.canonicalize().unwrap();
5679 let root_key = blob_root_key(&canonical_root);
5680 let absent_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
5681 let former_probe_seed = "1111111111111111111111111111111111111111111111111111111111111111";
5682 {
5683 let writer = backend.pool().writer().unwrap();
5684 writer
5685 .conn()
5686 .execute(
5687 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5688 VALUES (?1, ?2, 1), (?1, ?3, 1), (?1, ?4, 1)",
5689 rusqlite::params![
5690 root_key,
5691 content_ref.as_str(),
5692 absent_ref,
5693 former_probe_seed
5694 ],
5695 )
5696 .unwrap();
5697 }
5698
5699 assert_eq!(store.put(bytes).await.unwrap(), content_ref);
5705 let result = store
5706 .transactional_orphan_sweep(backend.sql().as_ref(), false)
5707 .await
5708 .unwrap();
5709 assert_eq!(result.deleted, 0);
5710 assert_eq!(result.grace_period_skipped, 1);
5711 assert!(store.exists(&content_ref).await.unwrap());
5712
5713 let remaining: i64 = backend
5714 .pool()
5715 .writer()
5716 .unwrap()
5717 .conn()
5718 .query_row(
5719 "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?1",
5720 [blob_root_key(&canonical_root)],
5721 |row| row.get(0),
5722 )
5723 .unwrap();
5724 assert_eq!(remaining, 0, "the next sweep recovers stale claims");
5725 }
5726
5727 #[tokio::test]
5728 async fn transactional_orphan_sweep_recovers_claims_after_root_relocation() {
5729 let dir = tempfile::tempdir().unwrap();
5730 let db_path = dir.path().join("khive.db");
5731 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5732 {
5733 let mut writer = backend.pool().writer().unwrap();
5734 prepare_completed_v21_gc_fixture(writer.conn_mut());
5735 }
5736
5737 let old_root = dir.path().join("old-blobs");
5738 let bytes = b"claim must follow a relocated blob root".to_vec();
5739 let content_ref = {
5740 let old_store = FsBlobStore::new(old_root.clone(), 0)
5741 .unwrap()
5742 .with_orphan_sweep_grace(Duration::from_secs(60));
5743 old_store.put(bytes).await.unwrap()
5744 };
5745 let old_root_key = blob_root_key(&old_root.canonicalize().unwrap());
5746 backend
5747 .pool()
5748 .writer()
5749 .unwrap()
5750 .conn()
5751 .execute(
5752 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5753 VALUES (?1, ?2, 1)",
5754 rusqlite::params![old_root_key, content_ref.as_str()],
5755 )
5756 .unwrap();
5757
5758 let new_root = dir.path().join("relocated-blobs");
5759 std::fs::rename(&old_root, &new_root).unwrap();
5760 let relocated_store = FsBlobStore::new(new_root, 0)
5761 .unwrap()
5762 .with_orphan_sweep_grace(Duration::from_secs(60));
5763 let result = relocated_store
5764 .transactional_orphan_sweep(backend.sql().as_ref(), false)
5765 .await
5766 .unwrap();
5767
5768 assert_eq!(
5769 result.deleted, 0,
5770 "a fresh relocated blob remains protected"
5771 );
5772 assert_eq!(result.grace_period_skipped, 1);
5773 assert!(relocated_store.exists(&content_ref).await.unwrap());
5774 let remaining: i64 = backend
5775 .pool()
5776 .writer()
5777 .unwrap()
5778 .conn()
5779 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5780 .unwrap();
5781 assert_eq!(
5782 remaining, 0,
5783 "exclusive database sweep ownership makes every pre-existing claim abandoned, \
5784 even when its old path-derived root key no longer matches"
5785 );
5786 }
5787
5788 #[tokio::test]
5789 async fn transactional_orphan_sweep_recovers_claims_copied_by_database_restore() {
5790 let dir = tempfile::tempdir().unwrap();
5791 let source_path = dir.path().join("source.db");
5792 let restored_path = dir.path().join("restored.db");
5793 let bytes = b"claim copied in an online database backup".to_vec();
5794 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
5795 {
5796 let source = crate::StorageBackend::sqlite(&source_path).unwrap();
5797 let mut writer = source.pool().writer().unwrap();
5798 prepare_completed_v21_gc_fixture(writer.conn_mut());
5799 writer
5800 .conn()
5801 .execute(
5802 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5803 VALUES ('source-root-before-backup', ?1, 1)",
5804 [content_ref.as_str()],
5805 )
5806 .unwrap();
5807 writer
5808 .conn()
5809 .execute_batch("PRAGMA wal_checkpoint(TRUNCATE)")
5810 .unwrap();
5811 }
5812 std::fs::copy(&source_path, &restored_path).unwrap();
5813
5814 let restored = crate::StorageBackend::sqlite(&restored_path).unwrap();
5815 let restored_root = dir.path().join("restored-blobs");
5816 let store = FsBlobStore::new(restored_root, 0)
5817 .unwrap()
5818 .with_orphan_sweep_grace(Duration::from_secs(60));
5819 assert_eq!(store.put(bytes).await.unwrap(), content_ref);
5820 let result = store
5821 .transactional_orphan_sweep(restored.sql().as_ref(), false)
5822 .await
5823 .unwrap();
5824
5825 assert_eq!(result.deleted, 0);
5826 assert_eq!(result.grace_period_skipped, 1);
5827 let remaining: i64 = restored
5828 .pool()
5829 .writer()
5830 .unwrap()
5831 .conn()
5832 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5833 .unwrap();
5834 assert_eq!(remaining, 0, "restored claims are abandoned ownership");
5835 }
5836
5837 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5838 async fn transactional_orphan_sweep_bounds_each_durable_claim_batch() {
5839 let dir = tempfile::tempdir().unwrap();
5840 let db_path = dir.path().join("khive.db");
5841 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5842 {
5843 let mut writer = backend.pool().writer().unwrap();
5844 prepare_completed_v21_gc_fixture(writer.conn_mut());
5845 }
5846 let root = dir.path().join("blobs");
5847 let store = std::sync::Arc::new(
5848 FsBlobStore::new(root.clone(), 0)
5849 .unwrap()
5850 .with_orphan_sweep_grace(Duration::ZERO),
5851 );
5852 let candidate_count = BLOB_GC_CLAIM_BATCH_SIZE * 2 + 1;
5853 for index in 0..candidate_count {
5854 store
5855 .put(format!("bounded claim candidate {index}").into_bytes())
5856 .await
5857 .unwrap();
5858 }
5859 let canonical_root = root.canonicalize().unwrap();
5860 let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
5861
5862 let sweep = {
5863 let store = store.clone();
5864 let sql = backend.sql();
5865 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5866 };
5867 assert!(
5868 recv_blocking(claimed).await,
5869 "the first bounded claim batch must commit before deletion"
5870 );
5871 let active_claims: i64 = rusqlite::Connection::open(&db_path)
5872 .unwrap()
5873 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5874 .unwrap();
5875 assert!(active_claims > 0);
5876 assert!(
5877 active_claims <= BLOB_GC_CLAIM_BATCH_SIZE as i64,
5878 "one transaction may expose at most {BLOB_GC_CLAIM_BATCH_SIZE} claim rows; \
5879 observed {active_claims}"
5880 );
5881
5882 release_delete.send(()).unwrap();
5883 let result = sweep.await.unwrap().unwrap();
5884 assert_eq!(result.deleted, candidate_count as u64);
5885 let remaining: i64 = rusqlite::Connection::open(&db_path)
5886 .unwrap()
5887 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5888 .unwrap();
5889 assert_eq!(remaining, 0);
5890 }
5891
5892 #[tokio::test]
5893 async fn abandoned_claim_recovery_deletes_at_most_one_batch_per_writer_hold() {
5894 let dir = tempfile::tempdir().unwrap();
5895 let db_path = dir.path().join("khive.db");
5896 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5897 {
5898 let mut writer = backend.pool().writer().unwrap();
5899 crate::run_migrations(writer.conn_mut()).unwrap();
5900 let tx = writer.conn_mut().transaction().unwrap();
5901 for index in 0..(BLOB_GC_CLAIM_BATCH_SIZE + 1) {
5902 tx.execute(
5903 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5904 VALUES ('abandoned-root', ?1, 1)",
5905 [format!("{index:064x}")],
5906 )
5907 .unwrap();
5908 }
5909 tx.commit().unwrap();
5910 }
5911
5912 let released = release_abandoned_blob_gc_claim_batch(backend.sql().as_ref())
5913 .await
5914 .unwrap();
5915 assert_eq!(released, BLOB_GC_CLAIM_BATCH_SIZE as u64);
5916 let remaining: i64 = backend
5917 .pool()
5918 .writer()
5919 .unwrap()
5920 .conn()
5921 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5922 .unwrap();
5923 assert_eq!(remaining, 1);
5924 }
5925
5926 #[tokio::test]
5927 async fn transactional_orphan_sweep_refuses_corrupt_liveness_and_claim_evidence() {
5928 let dir = tempfile::tempdir().unwrap();
5929 let db_path = dir.path().join("khive.db");
5930 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5931 {
5932 let mut writer = backend.pool().writer().unwrap();
5933 prepare_completed_v21_gc_fixture(writer.conn_mut());
5934 }
5935 let root = dir.path().join("blobs");
5936 let store = FsBlobStore::new(root.clone(), 0)
5937 .unwrap()
5938 .with_orphan_sweep_grace(Duration::ZERO);
5939 let orphan = store
5940 .put(b"must survive corrupt evidence".to_vec())
5941 .await
5942 .unwrap();
5943
5944 let conn = rusqlite::Connection::open(&db_path).unwrap();
5945 conn.execute_batch("PRAGMA ignore_check_constraints = ON")
5946 .unwrap();
5947 conn.execute(
5948 "INSERT INTO attachments \
5949 (record_uuid, substrate, role, content_ref, created_at) \
5950 VALUES ('corrupt-live', 'entity', 'content', 'not-a-content-ref', 1)",
5951 [],
5952 )
5953 .unwrap();
5954 conn.execute_batch("PRAGMA ignore_check_constraints = OFF")
5955 .unwrap();
5956 let live_error = store
5957 .transactional_orphan_sweep(backend.sql().as_ref(), false)
5958 .await
5959 .expect_err("corrupt live evidence must fail closed");
5960 assert!(matches!(live_error, StorageError::InvalidInput { .. }));
5961 assert!(
5962 store.exists(&orphan).await.unwrap(),
5963 "no file may be removed after corrupt live evidence"
5964 );
5965
5966 conn.execute(
5967 "DELETE FROM attachments WHERE record_uuid = 'corrupt-live'",
5968 [],
5969 )
5970 .unwrap();
5971 let root_key = blob_root_key(&root.canonicalize().unwrap());
5972 conn.execute(
5973 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5974 VALUES (?1, 'also-not-a-content-ref', 1)",
5975 [root_key.as_str()],
5976 )
5977 .unwrap();
5978 let claim_error = store
5979 .transactional_orphan_sweep(backend.sql().as_ref(), false)
5980 .await
5981 .expect_err("corrupt durable claim evidence must fail closed");
5982 assert!(matches!(claim_error, StorageError::InvalidInput { .. }));
5983 assert!(
5984 store.exists(&orphan).await.unwrap(),
5985 "no file may be removed after corrupt claim evidence"
5986 );
5987 let remaining: i64 = conn
5988 .query_row(
5989 "SELECT COUNT(*) FROM blob_gc_claims \
5990 WHERE root_key = ?1 AND content_ref = 'also-not-a-content-ref'",
5991 [root_key.as_str()],
5992 |row| row.get(0),
5993 )
5994 .unwrap();
5995 assert_eq!(
5996 remaining, 1,
5997 "corrupt claim evidence is not silently erased"
5998 );
5999 let probe_residue: i64 = conn
6000 .query_row(
6001 "SELECT (SELECT COUNT(*) FROM blob_gc_claims \
6002 WHERE root_key GLOB '__fence_probe-*') \
6003 + (SELECT COUNT(*) FROM attachments \
6004 WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
6005 [],
6006 |row| row.get(0),
6007 )
6008 .unwrap();
6009 assert_eq!(
6010 probe_residue, 0,
6011 "invalid evidence must abort before the functional fence probe"
6012 );
6013 }
6014
6015 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6016 async fn transactional_orphan_sweep_republishes_deduplicated_external_put() {
6017 let dir = tempfile::tempdir().unwrap();
6018 let db_path = dir.path().join("khive.db");
6019 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
6020 {
6021 let mut writer = backend.pool().writer().unwrap();
6022 prepare_completed_v21_gc_fixture(writer.conn_mut());
6023 }
6024 let root = dir.path().join("blobs");
6025 let store = std::sync::Arc::new(
6026 FsBlobStore::new(root.clone(), 0)
6027 .unwrap()
6028 .with_orphan_sweep_grace(Duration::ZERO),
6029 );
6030 let payload = b"existing orphan republished during sweep".to_vec();
6031 let orphan = store.put(payload.clone()).await.unwrap();
6032 let canonical_root = root.canonicalize().unwrap();
6033 let (marked, release, _done) = sync_hook::install(&canonical_root);
6034
6035 let sweep = {
6036 let store = store.clone();
6037 let sql = backend.sql();
6038 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
6039 };
6040 assert!(
6041 recv_blocking(marked).await,
6042 "sweep must finish its liveness mark"
6043 );
6044
6045 let external_lock = fs::OpenOptions::new()
6046 .read(true)
6047 .write(true)
6048 .open(root.join(ROOT_WRITE_LOCK_FILE))
6049 .unwrap();
6050 assert!(
6051 matches!(
6052 fs4::FileExt::try_lock(&external_lock),
6053 Err(fs4::TryLockError::WouldBlock)
6054 ),
6055 "the sweep must exclude a publisher using an independently opened root lock"
6056 );
6057
6058 let (started_tx, started_rx) = std::sync::mpsc::channel();
6059 let republished = {
6060 let root = root.clone();
6061 tokio::task::spawn_blocking(move || {
6062 let _ = started_tx.send(());
6063 put_blocking(&root, 0, payload)
6064 })
6065 };
6066 assert!(recv_blocking(started_rx).await, "blob put must start");
6067
6068 release.send(()).unwrap();
6069 let sweep_result = sweep.await.unwrap().unwrap();
6070 let republished = republished.await.unwrap().unwrap();
6071
6072 assert_eq!(sweep_result.deleted, 1);
6073 assert_eq!(republished, orphan);
6074 assert!(
6075 store.exists(&republished).await.unwrap(),
6076 "a deduplicated put concurrent with the sweep must not return a deleted reference"
6077 );
6078 }
6079
6080 #[tokio::test]
6081 async fn transactional_orphan_sweep_uses_all_attachment_refs_as_live() {
6082 let dir = tempfile::tempdir().unwrap();
6083 let db_path = dir.path().join("khive.db");
6084 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6085 {
6086 let mut writer = backend.pool().writer().unwrap();
6087 prepare_completed_v21_gc_fixture(writer.conn_mut());
6088 }
6089 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6090 .unwrap()
6091 .with_orphan_sweep_grace(Duration::ZERO);
6092 let live = store.put(b"live".to_vec()).await.unwrap();
6093 let soft_deleted = store.put(b"soft deleted".to_vec()).await.unwrap();
6094 let orphan = store.put(b"orphan".to_vec()).await.unwrap();
6095 {
6096 let writer = backend.pool().writer().unwrap();
6097 writer
6098 .conn()
6099 .execute_batch(
6100 "INSERT INTO entities \
6101 (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6102 VALUES ('live', 'local', 'document', 'live', '[]', 1, 1, NULL), \
6103 ('deleted', 'local', 'document', 'deleted', '[]', 1, 1, 2);",
6104 )
6105 .unwrap();
6106 writer
6107 .conn()
6108 .execute(
6109 "INSERT INTO attachments \
6110 (record_uuid, substrate, role, content_ref, created_at) \
6111 VALUES ('live', 'entity', 'content', ?1, 1), \
6112 ('deleted', 'entity', 'content', ?2, 1)",
6113 rusqlite::params![live.as_str(), soft_deleted.as_str()],
6114 )
6115 .unwrap();
6116 }
6117
6118 let dry_run = store
6119 .transactional_orphan_sweep(backend.sql().as_ref(), true)
6120 .await
6121 .unwrap();
6122 assert_eq!(dry_run.would_delete, 1);
6123 assert_eq!(dry_run.deleted, 0);
6124 assert!(store.exists(&soft_deleted).await.unwrap());
6125 assert!(store.exists(&orphan).await.unwrap());
6126
6127 let result = store
6128 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6129 .await
6130 .unwrap();
6131
6132 assert_eq!(result.scanned, 3);
6133 assert_eq!(result.deleted, 1);
6134 assert!(store.exists(&live).await.unwrap());
6135 assert!(
6136 store.exists(&soft_deleted).await.unwrap(),
6137 "soft delete retains attachment rows and their blobs"
6138 );
6139 assert!(!store.exists(&orphan).await.unwrap());
6140 }
6141
6142 #[tokio::test]
6143 async fn transactional_orphan_sweep_protects_a_freshly_published_blob_before_its_reference_commits(
6144 ) {
6145 let dir = tempfile::tempdir().unwrap();
6158 let db_path = dir.path().join("khive.db");
6159 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6160 {
6161 let mut writer = backend.pool().writer().unwrap();
6162 prepare_completed_v21_gc_fixture(writer.conn_mut());
6163 }
6164 let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
6167
6168 let blob = store
6171 .put(b"published, reference not yet committed".to_vec())
6172 .await
6173 .unwrap();
6174
6175 let result = store
6177 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6178 .await
6179 .unwrap();
6180
6181 assert_eq!(result.deleted, 0, "the blob must survive: {result:?}");
6182 assert_eq!(
6183 result.would_delete, 0,
6184 "not treated as a deletable orphan: {result:?}"
6185 );
6186 assert_eq!(
6187 result.grace_period_skipped, 1,
6188 "must be reported as grace-protected rather than silently ignored: {result:?}"
6189 );
6190 assert!(
6191 store.exists(&blob).await.unwrap(),
6192 "a blob still inside its publish grace period must survive the sweep"
6193 );
6194
6195 {
6198 let writer = backend.pool().writer().unwrap();
6199 writer
6200 .conn()
6201 .execute(
6202 "INSERT INTO entities \
6203 (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6204 VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
6205 [],
6206 )
6207 .unwrap();
6208 writer
6209 .conn()
6210 .execute(
6211 "INSERT INTO attachments \
6212 (record_uuid, substrate, role, content_ref, created_at) \
6213 VALUES ('e1', 'entity', 'content', ?1, 1)",
6214 [blob.as_str()],
6215 )
6216 .unwrap();
6217 }
6218
6219 let result = store
6222 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6223 .await
6224 .unwrap();
6225 assert_eq!(result.deleted, 0);
6226 assert!(store.exists(&blob).await.unwrap());
6227 }
6228
6229 #[tokio::test]
6230 async fn put_republishing_an_aged_orphan_restarts_its_grace_clock_before_the_reference_commits()
6231 {
6232 let dir = tempfile::tempdir().unwrap();
6240 let db_path = dir.path().join("khive.db");
6241 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6242 {
6243 let mut writer = backend.pool().writer().unwrap();
6244 prepare_completed_v21_gc_fixture(writer.conn_mut());
6245 }
6246 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6247 .unwrap()
6248 .with_orphan_sweep_grace(Duration::from_secs(60));
6249
6250 let bytes = b"old orphan re-published".to_vec();
6251 let first = store.put(bytes.clone()).await.unwrap();
6252
6253 let path = shard_path(store.root(), &first);
6256 let old_mtime = SystemTime::now() - Duration::from_secs(3600);
6257 fs::OpenOptions::new()
6258 .write(true)
6259 .open(&path)
6260 .unwrap()
6261 .set_modified(old_mtime)
6262 .unwrap();
6263
6264 let second = store.put(bytes).await.unwrap();
6267 assert_eq!(first, second);
6268
6269 let result = store
6272 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6273 .await
6274 .unwrap();
6275 assert_eq!(
6276 result.deleted, 0,
6277 "a dedup-republished blob must survive a sweep landing before its reference \
6278 commits: {result:?}"
6279 );
6280 assert_eq!(
6281 result.grace_period_skipped, 1,
6282 "must be reported as grace-protected, not silently ignored: {result:?}"
6283 );
6284 assert!(store.exists(&first).await.unwrap());
6285
6286 {
6288 let writer = backend.pool().writer().unwrap();
6289 writer
6290 .conn()
6291 .execute(
6292 "INSERT INTO entities \
6293 (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6294 VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
6295 [],
6296 )
6297 .unwrap();
6298 writer
6299 .conn()
6300 .execute(
6301 "INSERT INTO attachments \
6302 (record_uuid, substrate, role, content_ref, created_at) \
6303 VALUES ('e1', 'entity', 'content', ?1, 1)",
6304 [first.as_str()],
6305 )
6306 .unwrap();
6307 }
6308
6309 let result = store
6310 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6311 .await
6312 .unwrap();
6313 assert_eq!(result.deleted, 0);
6314 assert!(
6315 store.exists(&first).await.unwrap(),
6316 "the blob must stay live once its reference has committed"
6317 );
6318 }
6319
6320 #[tokio::test]
6321 async fn put_dedup_mtime_refresh_has_no_observable_effect_under_zero_grace_period() {
6322 let dir = tempfile::tempdir().unwrap();
6329 let db_path = dir.path().join("khive.db");
6330 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6331 {
6332 let mut writer = backend.pool().writer().unwrap();
6333 prepare_completed_v21_gc_fixture(writer.conn_mut());
6334 }
6335 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6336 .unwrap()
6337 .with_orphan_sweep_grace(Duration::ZERO);
6338 let bytes = b"zero grace dedup refresh".to_vec();
6339 let first = store.put(bytes.clone()).await.unwrap();
6340 let second = store.put(bytes.clone()).await.unwrap();
6341 assert_eq!(first, second);
6342
6343 let result = store
6344 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6345 .await
6346 .unwrap();
6347 assert_eq!(
6348 result.deleted, 1,
6349 "a zero grace period must still delete an unreferenced blob even after a dedup \
6350 put refreshed its mtime: {result:?}"
6351 );
6352 assert!(!store.exists(&first).await.unwrap());
6353 }
6354
6355 #[tokio::test]
6356 async fn transactional_orphan_sweep_still_removes_orphans_older_than_the_grace_period() {
6357 let dir = tempfile::tempdir().unwrap();
6362 let db_path = dir.path().join("khive.db");
6363 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6364 {
6365 let mut writer = backend.pool().writer().unwrap();
6366 prepare_completed_v21_gc_fixture(writer.conn_mut());
6367 }
6368 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6369 .unwrap()
6370 .with_orphan_sweep_grace(Duration::from_secs(60));
6371
6372 let orphan = store
6373 .put(b"actually orphaned, published long ago".to_vec())
6374 .await
6375 .unwrap();
6376 let path = shard_path(store.root(), &orphan);
6379 let old_mtime = SystemTime::now() - Duration::from_secs(3600);
6380 fs::OpenOptions::new()
6381 .write(true)
6382 .open(&path)
6383 .unwrap()
6384 .set_modified(old_mtime)
6385 .unwrap();
6386
6387 let result = store
6388 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6389 .await
6390 .unwrap();
6391
6392 assert_eq!(
6393 result.deleted, 1,
6394 "an orphan older than the grace period must still be swept: {result:?}"
6395 );
6396 assert_eq!(result.grace_period_skipped, 0);
6397 assert!(!store.exists(&orphan).await.unwrap());
6398 }
6399
6400 #[test]
6401 fn resolve_blob_root_prefers_env_var() {
6402 let _guard = ENV_LOCK.lock().unwrap();
6403 std::env::set_var("KHIVE_BLOB_ROOT", "/tmp/env-override-root");
6404 let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
6405 std::env::remove_var("KHIVE_BLOB_ROOT");
6406 assert_eq!(resolved.unwrap(), PathBuf::from("/tmp/env-override-root"));
6407 }
6408
6409 #[test]
6410 fn resolve_blob_root_prefers_config_over_default() {
6411 let _guard = ENV_LOCK.lock().unwrap();
6412 std::env::remove_var("KHIVE_BLOB_ROOT");
6413 let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
6414 assert_eq!(resolved.unwrap(), PathBuf::from("/cfg/root"));
6415 }
6416
6417 #[test]
6418 fn resolve_blob_root_defaults_beside_db_dir() {
6419 let _guard = ENV_LOCK.lock().unwrap();
6420 std::env::remove_var("KHIVE_BLOB_ROOT");
6421 let resolved = resolve_blob_root(Some(Path::new("/db/dir")), None);
6422 assert_eq!(resolved.unwrap(), PathBuf::from("/db/dir/blobs"));
6423 }
6424
6425 #[test]
6426 fn resolve_blob_root_errors_with_no_env_config_or_db_dir() {
6427 let _guard = ENV_LOCK.lock().unwrap();
6428 std::env::remove_var("KHIVE_BLOB_ROOT");
6429 let resolved = resolve_blob_root(None, None);
6430 assert!(resolved.is_err());
6431 }
6432
6433 static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
6437
6438 #[cfg(unix)]
6439 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6440 async fn transactional_orphan_sweep_walk_ignores_a_leaf_swapped_for_an_outside_symlink_mid_scan(
6441 ) {
6442 let dir = tempfile::tempdir().unwrap();
6458 let db_path = dir.path().join("khive.db");
6459 let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
6460 {
6461 let mut writer = backend.pool().writer().unwrap();
6462 prepare_completed_v21_gc_fixture(writer.conn_mut());
6463 }
6464
6465 let root = dir.path().join("blobs");
6466 let store = std::sync::Arc::new(
6467 FsBlobStore::new(root.clone(), 0)
6468 .unwrap()
6469 .with_orphan_sweep_grace(Duration::from_secs(3600)),
6470 );
6471
6472 let real = store
6476 .put(b"real freshly-published blob".to_vec())
6477 .await
6478 .unwrap();
6479 let real_path = shard_path(&root, &real);
6480
6481 let outside_dir = dir.path().join("outside");
6486 fs::create_dir_all(&outside_dir).unwrap();
6487 let decoy_path = outside_dir.join(real.as_str());
6488 fs::write(&decoy_path, b"outside decoy, must never be observed").unwrap();
6489 let ancient = SystemTime::now() - Duration::from_secs(7200);
6490 fs::OpenOptions::new()
6491 .write(true)
6492 .open(&decoy_path)
6493 .unwrap()
6494 .set_modified(ancient)
6495 .unwrap();
6496
6497 let (reached, release) = walk_leaf_sync_hook::install(&root);
6498 let sweep = {
6499 let store = store.clone();
6500 let sql = backend.sql();
6501 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), true).await })
6502 };
6503 assert!(
6504 recv_blocking(reached).await,
6505 "sweep walk must reach the leaf classification pause"
6506 );
6507
6508 fs::remove_file(&real_path).unwrap();
6511 std::os::unix::fs::symlink(&decoy_path, &real_path).unwrap();
6512
6513 release.send(()).unwrap();
6514 let result = sweep.await.unwrap().unwrap();
6515
6516 assert_eq!(
6517 result.would_delete, 0,
6518 "an outside decoy's stale mtime must never make an in-root candidate \
6519 eligible for deletion: {result:?}"
6520 );
6521 assert_eq!(
6522 result.grace_period_skipped, 0,
6523 "the swapped leaf is a symlink; `openat(..., O_NOFOLLOW)` refuses it, so it \
6524 must be dropped from candidates entirely rather than counted (real or \
6525 outside) at all: {result:?}"
6526 );
6527 assert_eq!(
6528 result.scanned, 0,
6529 "the symlinked leaf must never be scanned as a candidate: {result:?}"
6530 );
6531
6532 fs::remove_file(&real_path).unwrap();
6536 fs::write(&real_path, b"real freshly-published blob").unwrap();
6537 let control = store
6538 .transactional_orphan_sweep(backend.sql().as_ref(), true)
6539 .await
6540 .unwrap();
6541 assert_eq!(
6542 control.grace_period_skipped, 1,
6543 "control: the un-replaced root must still classify the real candidate as \
6544 grace-protected: {control:?}"
6545 );
6546 assert_eq!(control.would_delete, 0);
6547 }
6548
6549 #[cfg(unix)]
6550 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6551 async fn transactional_orphan_sweep_walk_still_finds_a_real_orphan_past_its_grace_period() {
6552 let dir = tempfile::tempdir().unwrap();
6558 let db_path = dir.path().join("khive.db");
6559 let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6560 {
6561 let mut writer = backend.pool().writer().unwrap();
6562 prepare_completed_v21_gc_fixture(writer.conn_mut());
6563 }
6564 let root = dir.path().join("blobs");
6565 let store = FsBlobStore::new(root.clone(), 0)
6566 .unwrap()
6567 .with_orphan_sweep_grace(Duration::from_secs(60));
6568
6569 let orphan = store.put(b"aged real orphan".to_vec()).await.unwrap();
6570 let path = shard_path(&root, &orphan);
6571 let ancient = SystemTime::now() - Duration::from_secs(3600);
6572 fs::OpenOptions::new()
6573 .write(true)
6574 .open(&path)
6575 .unwrap()
6576 .set_modified(ancient)
6577 .unwrap();
6578
6579 let result = store
6580 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6581 .await
6582 .unwrap();
6583 assert_eq!(
6584 result.deleted, 1,
6585 "a real orphan older than the grace period must still be swept: {result:?}"
6586 );
6587 assert!(!store.exists(&orphan).await.unwrap());
6588 }
6589}