1use std::collections::HashMap;
18use std::fs;
19use std::io::{Read, Write};
20use std::path::{Path, PathBuf};
21use std::sync::{Arc, Mutex as StdMutex, OnceLock};
22use std::time::{Duration, SystemTime};
23
24use async_trait::async_trait;
25
26use khive_storage::blob::{
27 BlobOrphanSweepConfig, BlobOrphanSweepResult, BlobStore, ContentRef, UploadId,
28 MAX_BLOB_WHOLE_BYTES,
29};
30use khive_storage::error::StorageError;
31use khive_storage::types::{SqlRow, SqlStatement, SqlValue, StorageResult};
32use khive_storage::{AtomicUnitOp, SqlAccess, StorageCapability};
33
34use crate::error::SqliteError;
35use uuid::Uuid;
36
37#[path = "blob_uploads.rs"]
38mod uploads;
39
40const ROOT_WRITE_LOCK_FILE: &str = ".khive-blob-write.lock";
41const DATABASE_GC_LOCK_SUFFIX: &str = ".khive-blob-gc.lock";
42const BLOB_GC_CLAIM_BATCH_SIZE: usize = 128;
46
47fn map_io_err(e: std::io::Error, op: &'static str) -> StorageError {
48 StorageError::driver(StorageCapability::Blob, op, e)
49}
50
51#[cfg(any(test, not(unix)))]
52fn shard_path(root: &Path, content_ref: &ContentRef) -> PathBuf {
53 let hex = content_ref.as_str();
54 root.join(&hex[0..2]).join(&hex[2..4]).join(hex)
55}
56
57#[cfg(unix)]
58fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
59 open_dir_no_follow(root)
60}
61
62#[cfg(windows)]
63fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
64 use std::fs::OpenOptions;
65 use std::os::windows::fs::OpenOptionsExt;
66
67 const FILE_SHARE_READ: u32 = 0x1;
68 const FILE_SHARE_WRITE: u32 = 0x2;
69 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
70 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
71
72 let handle = OpenOptions::new()
73 .read(true)
74 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
78 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
79 .open(root)?;
80 let file_type = handle.metadata()?.file_type();
81 if file_type.is_symlink() || !file_type.is_dir() {
82 return Err(std::io::Error::new(
83 std::io::ErrorKind::InvalidInput,
84 format!(
85 "blob store root is not a directory or is a reparse point: {}",
86 root.display()
87 ),
88 ));
89 }
90 Ok(handle)
91}
92
93#[cfg(not(any(unix, windows)))]
94fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
95 let handle = std::fs::File::open(root)?;
96 if !handle.metadata()?.is_dir() {
97 return Err(std::io::Error::new(
98 std::io::ErrorKind::InvalidInput,
99 format!("blob store root is not a directory: {}", root.display()),
100 ));
101 }
102 Ok(handle)
103}
104
105#[cfg(unix)]
111fn verify_blob_root_identity(root: &Path, root_handle: &std::fs::File) -> std::io::Result<()> {
112 use std::os::unix::fs::MetadataExt;
113
114 let current = open_dir_no_follow(root).map_err(|error| {
115 std::io::Error::new(
116 std::io::ErrorKind::InvalidInput,
117 format!(
118 "blob store root is no longer reachable as its initialization-time directory ({}): {error}",
119 root.display()
120 ),
121 )
122 })?;
123 let expected = root_handle.metadata()?;
124 let current = current.metadata()?;
125 if expected.dev() != current.dev() || expected.ino() != current.ino() {
126 return Err(std::io::Error::new(
127 std::io::ErrorKind::InvalidInput,
128 format!(
129 "blob store root no longer names its initialization-time directory: {}",
130 root.display()
131 ),
132 ));
133 }
134 Ok(())
135}
136
137#[cfg(not(unix))]
138fn verify_blob_root_identity(root: &Path, _root_handle: &std::fs::File) -> std::io::Result<()> {
139 let current = root.canonicalize().map_err(|error| {
145 std::io::Error::new(
146 std::io::ErrorKind::InvalidInput,
147 format!(
148 "blob store root is no longer reachable as its initialization-time directory ({}): {error}",
149 root.display()
150 ),
151 )
152 })?;
153 if current != root {
154 return Err(std::io::Error::new(
155 std::io::ErrorKind::InvalidInput,
156 format!(
157 "blob store root no longer names its initialization-time directory: {}",
158 root.display()
159 ),
160 ));
161 }
162 Ok(())
163}
164
165#[cfg(unix)]
183fn unlink_blob_shard_file_no_follow(
184 root: &Path,
185 root_handle: &std::fs::File,
186 content_ref: &ContentRef,
187) -> std::io::Result<()> {
188 use std::os::unix::io::AsRawFd;
189
190 verify_blob_root_identity(root, root_handle)?;
191 let hex = content_ref.as_str();
192 let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
193 let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
194 unlink_entry_at(shard2_dir.as_raw_fd(), hex)
197}
198
199#[cfg(windows)]
200fn unlink_blob_shard_file_no_follow(
201 root: &Path,
202 root_handle: &std::fs::File,
203 content_ref: &ContentRef,
204) -> std::io::Result<()> {
205 use std::fs::OpenOptions;
248 use std::os::windows::ffi::OsStringExt;
249 use std::os::windows::fs::OpenOptionsExt;
250 use std::os::windows::io::AsRawHandle;
251 use windows_sys::Win32::Storage::FileSystem::{
252 FileDispositionInfo, GetFinalPathNameByHandleW, SetFileInformationByHandle,
253 FILE_DISPOSITION_INFO,
254 };
255
256 const FILE_SHARE_READ: u32 = 0x1;
257 const FILE_SHARE_WRITE: u32 = 0x2;
258 const DELETE: u32 = 0x0001_0000;
259 const FILE_READ_ATTRIBUTES: u32 = 0x80;
260 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
261 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
262 const FINAL_PATH_FLAGS: u32 = 0x0;
266
267 fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
268 let dir = OpenOptions::new()
269 .read(true)
270 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
271 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
272 .open(path)?;
273 let file_type = dir.metadata()?.file_type();
274 if file_type.is_symlink() || !file_type.is_dir() {
275 return Err(std::io::Error::new(
276 std::io::ErrorKind::InvalidInput,
277 format!(
278 "refusing to unlink blob shard file through non-directory or \
279 reparse-point path component: {}",
280 path.display()
281 ),
282 ));
283 }
284 Ok(dir)
285 }
286
287 fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<std::path::PathBuf> {
291 let handle = file.as_raw_handle();
292 let mut buf: Vec<u16> = vec![0; 512];
293 loop {
294 let len = unsafe {
295 GetFinalPathNameByHandleW(
296 handle as _,
297 buf.as_mut_ptr(),
298 buf.len() as u32,
299 FINAL_PATH_FLAGS,
300 )
301 };
302 if len == 0 {
303 return Err(std::io::Error::last_os_error());
304 }
305 let len = len as usize;
306 if len <= buf.len() {
307 buf.truncate(len);
308 return Ok(std::path::PathBuf::from(std::ffi::OsString::from_wide(
309 &buf,
310 )));
311 }
312 buf.resize(len, 0);
315 }
316 }
317
318 verify_blob_root_identity(root, root_handle)?;
319 let hex = content_ref.as_str();
320 let shard1 = root.join(&hex[0..2]);
321 let shard2 = shard1.join(&hex[2..4]);
322 let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
323 let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
324
325 let expected = final_path_by_handle(root_handle)?
326 .join(&hex[0..2])
327 .join(&hex[2..4])
328 .join(hex);
329
330 let target = OpenOptions::new()
336 .access_mode(DELETE | FILE_READ_ATTRIBUTES)
337 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
338 .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
339 .open(shard2.join(hex))?;
340
341 let resolved = final_path_by_handle(&target)?;
342 if resolved != expected {
343 return Err(std::io::Error::new(
344 std::io::ErrorKind::InvalidInput,
345 format!(
346 "refusing blob delete: handle resolved outside the verified blob root \
347 (expected {}, resolved {})",
348 expected.display(),
349 resolved.display()
350 ),
351 ));
352 }
353
354 let disposition = FILE_DISPOSITION_INFO { DeleteFile: true };
355 let ok = unsafe {
356 SetFileInformationByHandle(
357 target.as_raw_handle() as _,
358 FileDispositionInfo,
359 std::ptr::from_ref(&disposition).cast(),
360 std::mem::size_of::<FILE_DISPOSITION_INFO>() as u32,
361 )
362 };
363 if ok == 0 {
364 return Err(std::io::Error::last_os_error());
365 }
366 Ok(())
367}
368
369#[cfg(not(any(unix, windows)))]
370fn unlink_blob_shard_file_no_follow(
371 root: &Path,
372 root_handle: &std::fs::File,
373 content_ref: &ContentRef,
374) -> std::io::Result<()> {
375 verify_blob_root_identity(root, root_handle)?;
389 let hex = content_ref.as_str();
390 let shard1 = root.join(&hex[0..2]);
391 let shard2 = shard1.join(&hex[2..4]);
392 for component in [root, shard1.as_path(), shard2.as_path()] {
393 let metadata = fs::symlink_metadata(component)?;
394 if metadata.file_type().is_symlink() {
395 return Err(std::io::Error::new(
396 std::io::ErrorKind::InvalidInput,
397 format!(
398 "refusing to unlink blob shard file through symlinked path component: {}",
399 component.display()
400 ),
401 ));
402 }
403 }
404 fs::remove_file(shard2.join(hex))
405}
406
407#[cfg(unix)]
411fn open_blob_shard_file_no_follow(
412 root: &Path,
413 root_handle: &std::fs::File,
414 content_ref: &ContentRef,
415) -> std::io::Result<std::fs::File> {
416 verify_blob_root_identity(root, root_handle)?;
417 open_blob_shard_file_at_no_follow(root_handle, content_ref, libc::O_RDONLY)
418}
419
420#[cfg(windows)]
421fn open_blob_shard_file_no_follow(
422 root: &Path,
423 root_handle: &std::fs::File,
424 content_ref: &ContentRef,
425) -> std::io::Result<std::fs::File> {
426 use std::fs::OpenOptions;
427 use std::os::windows::ffi::OsStringExt;
428 use std::os::windows::fs::OpenOptionsExt;
429 use std::os::windows::io::AsRawHandle;
430 use windows_sys::Win32::Storage::FileSystem::GetFinalPathNameByHandleW;
431
432 const FILE_SHARE_READ: u32 = 0x1;
433 const FILE_SHARE_WRITE: u32 = 0x2;
434 const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
435 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
436 const FINAL_PATH_FLAGS: u32 = 0x0;
437
438 fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
439 let dir = OpenOptions::new()
440 .read(true)
441 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
444 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
445 .open(path)?;
446 let file_type = dir.metadata()?.file_type();
447 if file_type.is_symlink() || !file_type.is_dir() {
448 return Err(std::io::Error::new(
449 std::io::ErrorKind::InvalidInput,
450 format!(
451 "refusing to read a blob through non-directory or reparse-point component: {}",
452 path.display()
453 ),
454 ));
455 }
456 Ok(dir)
457 }
458
459 fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<PathBuf> {
460 let mut buf: Vec<u16> = vec![0; 512];
461 loop {
462 let len = unsafe {
463 GetFinalPathNameByHandleW(
464 file.as_raw_handle() as _,
465 buf.as_mut_ptr(),
466 buf.len() as u32,
467 FINAL_PATH_FLAGS,
468 )
469 };
470 if len == 0 {
471 return Err(std::io::Error::last_os_error());
472 }
473 let len = len as usize;
474 if len <= buf.len() {
475 buf.truncate(len);
476 return Ok(PathBuf::from(std::ffi::OsString::from_wide(&buf)));
477 }
478 buf.resize(len, 0);
479 }
480 }
481
482 verify_blob_root_identity(root, root_handle)?;
483 let hex = content_ref.as_str();
484 let shard1 = root.join(&hex[0..2]);
485 let shard2 = shard1.join(&hex[2..4]);
486 let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
487 let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
488 let expected = final_path_by_handle(root_handle)?
489 .join(&hex[0..2])
490 .join(&hex[2..4])
491 .join(hex);
492
493 let target = OpenOptions::new()
494 .read(true)
495 .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
496 .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
497 .open(shard2.join(hex))?;
498 let file_type = target.metadata()?.file_type();
499 if file_type.is_symlink() || !file_type.is_file() {
500 return Err(std::io::Error::new(
501 std::io::ErrorKind::InvalidInput,
502 "refusing to read a blob leaf that is not a regular file",
503 ));
504 }
505 let resolved = final_path_by_handle(&target)?;
506 if resolved != expected {
507 return Err(std::io::Error::new(
508 std::io::ErrorKind::InvalidInput,
509 format!(
510 "refusing blob read: handle resolved outside the verified blob root (expected {}, resolved {})",
511 expected.display(),
512 resolved.display()
513 ),
514 ));
515 }
516 Ok(target)
517}
518
519#[cfg(not(any(unix, windows)))]
520fn open_blob_shard_file_no_follow(
521 root: &Path,
522 root_handle: &std::fs::File,
523 _content_ref: &ContentRef,
524) -> std::io::Result<std::fs::File> {
525 verify_blob_root_identity(root, root_handle)?;
526 Err(std::io::Error::new(
527 std::io::ErrorKind::Unsupported,
528 "bounded verified blob reads require handle-relative no-follow file APIs",
529 ))
530}
531
532#[cfg(unix)]
533fn open_dir_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
534 use std::os::unix::ffi::OsStrExt;
535 use std::os::unix::io::FromRawFd;
536
537 let c_path = std::ffi::CString::new(path.as_os_str().as_bytes())
538 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
539 let fd = unsafe {
542 libc::open(
543 c_path.as_ptr(),
544 libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
545 )
546 };
547 if fd < 0 {
548 return Err(std::io::Error::last_os_error());
549 }
550 Ok(unsafe { std::fs::File::from_raw_fd(fd) })
553}
554
555#[cfg(unix)]
556fn openat_dir_no_follow(
557 parent_fd: std::os::unix::io::RawFd,
558 name: &str,
559) -> std::io::Result<std::fs::File> {
560 use std::os::unix::io::FromRawFd;
561
562 let c_name = std::ffi::CString::new(name)
563 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
564 let fd = unsafe {
568 libc::openat(
569 parent_fd,
570 c_name.as_ptr(),
571 libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
572 )
573 };
574 if fd < 0 {
575 return Err(std::io::Error::last_os_error());
576 }
577 Ok(unsafe { std::fs::File::from_raw_fd(fd) })
580}
581
582#[cfg(unix)]
583fn openat_regular_file_no_follow(
584 parent_fd: std::os::unix::io::RawFd,
585 name: &str,
586 access_flags: libc::c_int,
587) -> std::io::Result<std::fs::File> {
588 use std::os::unix::io::FromRawFd;
589
590 let c_name = std::ffi::CString::new(name)
591 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
592 let fd = unsafe {
596 libc::openat(
597 parent_fd,
598 c_name.as_ptr(),
599 access_flags | libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK,
600 )
601 };
602 if fd < 0 {
603 return Err(std::io::Error::last_os_error());
604 }
605 let file = unsafe { std::fs::File::from_raw_fd(fd) };
607 if !file.metadata()?.file_type().is_file() {
608 return Err(std::io::Error::new(
609 std::io::ErrorKind::InvalidInput,
610 format!("refusing non-regular blob store entry: {name}"),
611 ));
612 }
613 Ok(file)
614}
615
616#[cfg(unix)]
617fn open_blob_shard_file_at_no_follow(
618 root_handle: &std::fs::File,
619 content_ref: &ContentRef,
620 access_flags: libc::c_int,
621) -> std::io::Result<std::fs::File> {
622 use std::os::unix::io::AsRawFd;
623
624 let hex = content_ref.as_str();
625 let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
626 let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
627 openat_regular_file_no_follow(shard2_dir.as_raw_fd(), hex, access_flags)
628}
629
630#[cfg(unix)]
631fn open_or_create_dir_at_no_follow(
632 parent_fd: std::os::unix::io::RawFd,
633 name: &str,
634) -> std::io::Result<std::fs::File> {
635 let c_name = std::ffi::CString::new(name)
636 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
637 match openat_dir_no_follow(parent_fd, name) {
638 Ok(dir) => Ok(dir),
639 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
640 let rc = unsafe { libc::mkdirat(parent_fd, c_name.as_ptr(), 0o777) };
642 if rc != 0 {
643 let error = std::io::Error::last_os_error();
644 if error.kind() != std::io::ErrorKind::AlreadyExists {
645 return Err(error);
646 }
647 }
648 openat_dir_no_follow(parent_fd, name)
652 }
653 Err(error) => Err(error),
654 }
655}
656
657#[cfg(unix)]
658fn create_regular_file_at_no_follow(
659 parent_fd: std::os::unix::io::RawFd,
660 name: &str,
661 mode: libc::mode_t,
662) -> std::io::Result<std::fs::File> {
663 use std::os::unix::io::FromRawFd;
664
665 let c_name = std::ffi::CString::new(name)
666 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
667 let fd = unsafe {
670 libc::openat(
671 parent_fd,
672 c_name.as_ptr(),
673 libc::O_RDWR | libc::O_CREAT | libc::O_EXCL | libc::O_NOFOLLOW | libc::O_CLOEXEC,
674 mode as libc::c_uint,
675 )
676 };
677 if fd < 0 {
678 return Err(std::io::Error::last_os_error());
679 }
680 Ok(unsafe { std::fs::File::from_raw_fd(fd) })
682}
683
684#[cfg(unix)]
685fn unlink_entry_at(parent_fd: std::os::unix::io::RawFd, name: &str) -> std::io::Result<()> {
686 let c_name = std::ffi::CString::new(name)
687 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
688 let rc = unsafe { libc::unlinkat(parent_fd, c_name.as_ptr(), 0) };
690 if rc != 0 {
691 return Err(std::io::Error::last_os_error());
692 }
693 Ok(())
694}
695
696#[cfg(unix)]
697fn rename_entry_at(
698 source_fd: std::os::unix::io::RawFd,
699 from: &str,
700 destination_fd: std::os::unix::io::RawFd,
701 to: &str,
702) -> std::io::Result<()> {
703 let c_from = std::ffi::CString::new(from)
704 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
705 let c_to = std::ffi::CString::new(to)
706 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
707 let rc = unsafe { libc::renameat(source_fd, c_from.as_ptr(), destination_fd, c_to.as_ptr()) };
710 if rc != 0 {
711 return Err(std::io::Error::last_os_error());
712 }
713 Ok(())
714}
715
716#[cfg(unix)]
717fn available_space_at(root_handle: &std::fs::File) -> std::io::Result<u64> {
718 use std::os::unix::io::AsRawFd;
719
720 let mut stat: libc::statvfs = unsafe { std::mem::zeroed() };
721 let rc = unsafe { libc::fstatvfs(root_handle.as_raw_fd(), &mut stat) };
723 if rc != 0 {
724 return Err(std::io::Error::last_os_error());
725 }
726 #[allow(clippy::useless_conversion)]
730 Ok(stat.f_frsize.saturating_mul(u64::from(stat.f_bavail)))
731}
732
733#[cfg(unix)]
734fn acquire_root_write_lock_at(root_handle: &std::fs::File) -> StorageResult<std::fs::File> {
735 use std::os::unix::io::{AsRawFd, FromRawFd};
736
737 let c_name = std::ffi::CString::new(ROOT_WRITE_LOCK_FILE)
738 .expect("the static blob root lock name contains no NUL");
739 let fd = unsafe {
742 libc::openat(
743 root_handle.as_raw_fd(),
744 c_name.as_ptr(),
745 libc::O_RDWR | libc::O_CREAT | libc::O_NOFOLLOW | libc::O_CLOEXEC,
746 0o666,
747 )
748 };
749 if fd < 0 {
750 return Err(map_io_err(
751 std::io::Error::last_os_error(),
752 "root_write_lock_open",
753 ));
754 }
755 let lock_file = unsafe { std::fs::File::from_raw_fd(fd) };
757 if !lock_file
758 .metadata()
759 .map_err(|e| map_io_err(e, "root_write_lock_metadata"))?
760 .file_type()
761 .is_file()
762 {
763 return Err(map_io_err(
764 std::io::Error::new(
765 std::io::ErrorKind::InvalidInput,
766 "blob root write lock is not a regular file",
767 ),
768 "root_write_lock_open",
769 ));
770 }
771 fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
772 Ok(lock_file)
773}
774
775fn blob_root_key(root: &Path) -> String {
783 #[cfg(unix)]
784 let bytes = {
785 use std::os::unix::ffi::OsStrExt;
786 root.as_os_str().as_bytes().to_vec()
787 };
788 #[cfg(windows)]
789 let bytes = {
790 use std::os::windows::ffi::OsStrExt;
791 root.as_os_str()
792 .encode_wide()
793 .flat_map(u16::to_le_bytes)
794 .collect::<Vec<_>>()
795 };
796 #[cfg(not(any(unix, windows)))]
797 let bytes = root.to_string_lossy().as_bytes().to_vec();
798 blake3::hash(&bytes).to_hex().to_string()
799}
800
801pub fn resolve_blob_root(
810 db_dir: Option<&Path>,
811 config_root: Option<&Path>,
812) -> Result<PathBuf, SqliteError> {
813 if let Ok(env_root) = std::env::var("KHIVE_BLOB_ROOT") {
814 if !env_root.trim().is_empty() {
815 return Ok(PathBuf::from(env_root));
816 }
817 }
818 if let Some(root) = config_root {
819 return Ok(root.to_path_buf());
820 }
821 if let Some(dir) = db_dir {
822 return Ok(dir.join("blobs"));
823 }
824 Err(SqliteError::InvalidData(
825 "cannot resolve a blob store root: no KHIVE_BLOB_ROOT env var, no configured \
826 root, and the database has no on-disk directory to default beside (in-memory \
827 backend)"
828 .to_string(),
829 ))
830}
831
832fn crosses_floor(available: u64, required_write_bytes: u64, floor_bytes: u64) -> bool {
846 available.saturating_sub(required_write_bytes) < floor_bytes
847}
848
849#[cfg(any(test, not(unix)))]
850fn put_blocking_with_space_probe<F>(
851 root: &Path,
852 floor_bytes: u64,
853 bytes: Vec<u8>,
854 available_space: F,
855) -> StorageResult<ContentRef>
856where
857 F: FnOnce(&Path) -> std::io::Result<u64>,
858{
859 let digest = blake3::hash(&bytes);
860 let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
861 let target = shard_path(root, &content_ref);
862
863 if target.exists() {
874 let file = fs::OpenOptions::new()
875 .write(true)
876 .open(&target)
877 .map_err(|e| map_io_err(e, "put_touch_open"))?;
878 file.set_modified(SystemTime::now())
879 .map_err(|e| map_io_err(e, "put_touch_mtime"))?;
880 return Ok(content_ref);
881 }
882
883 let required_write_bytes = bytes.len() as u64;
884 let available = available_space(root).map_err(|e| map_io_err(e, "put_check_space"))?;
885 if crosses_floor(available, required_write_bytes, floor_bytes) {
886 return Err(StorageError::CapacityFloor {
887 capability: StorageCapability::Blob,
888 volume: root.display().to_string(),
889 available_bytes: available,
890 floor_bytes,
891 });
892 }
893
894 let shard_dir = target
895 .parent()
896 .expect("shard_path always nests under two directory levels");
897 fs::create_dir_all(shard_dir).map_err(|e| map_io_err(e, "put_mkdir"))?;
898
899 let mut tmp = tempfile::Builder::new()
900 .prefix(".tmp-")
901 .tempfile_in(shard_dir)
902 .map_err(|e| map_io_err(e, "put_tempfile"))?;
903 tmp.write_all(&bytes)
904 .map_err(|e| map_io_err(e, "put_write"))?;
905 tmp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
906 tmp.as_file()
907 .sync_all()
908 .map_err(|e| map_io_err(e, "put_fsync"))?;
909
910 let written_len = tmp
911 .as_file()
912 .metadata()
913 .map_err(|e| map_io_err(e, "put_verify"))?
914 .len();
915 if written_len != bytes.len() as u64 {
916 return Err(map_io_err(
917 std::io::Error::other(format!(
918 "temp file length {written_len} does not match {} written bytes",
919 bytes.len()
920 )),
921 "put_verify",
922 ));
923 }
924
925 let temporary = tmp.into_temp_path();
926 publish_blob_path(&temporary, &target)?;
927
928 Ok(content_ref)
929}
930
931#[cfg(any(test, not(unix)))]
932fn publish_blob_path(source: &Path, target: &Path) -> StorageResult<()> {
933 fs::rename(source, target).map_err(|error| map_io_err(error, "put_persist"))
934}
935
936#[cfg(any(test, not(unix)))]
937fn put_blocking(root: &Path, floor_bytes: u64, bytes: Vec<u8>) -> StorageResult<ContentRef> {
938 let _root_write_guard = acquire_root_write_lock(root)?;
939 put_blocking_with_space_probe(root, floor_bytes, bytes, |path| fs4::available_space(path))
940}
941
942fn acquire_root_write_lock_anchored(
943 root: &Path,
944 root_handle: &std::fs::File,
945) -> StorageResult<std::fs::File> {
946 verify_blob_root_identity(root, root_handle)
947 .map_err(|e| map_io_err(e, "root_write_lock_identity"))?;
948 #[cfg(unix)]
949 {
950 acquire_root_write_lock_at(root_handle)
951 }
952 #[cfg(not(unix))]
953 {
954 acquire_root_write_lock(root)
955 }
956}
957
958#[cfg(unix)]
959struct BlobPublication {
960 #[cfg(test)]
961 hook: Option<sync_hook::Publication>,
962}
963
964#[cfg(unix)]
965impl BlobPublication {
966 fn step<T>(
967 &self,
968 operation: &'static str,
969 action: impl FnOnce() -> std::io::Result<T>,
970 ) -> StorageResult<T> {
971 self.io_step(operation, action)
972 .map_err(|error| map_io_err(error, operation))
973 }
974
975 fn io_step<T>(
976 &self,
977 _operation: &'static str,
978 action: impl FnOnce() -> std::io::Result<T>,
979 ) -> std::io::Result<T> {
980 #[cfg(test)]
981 if let Some(hook) = &self.hook {
982 hook.before(_operation)?;
983 }
984 let result = action()?;
985 #[cfg(test)]
986 if let Some(hook) = &self.hook {
987 hook.completed(_operation);
988 }
989 Ok(result)
990 }
991
992 fn sync_directories(
993 &self,
994 root: &fs::File,
995 shard1: &fs::File,
996 shard2: &fs::File,
997 ) -> StorageResult<()> {
998 for (operation, directory) in [
1001 ("put_sync_shard", shard2),
1002 ("put_sync_parent", shard1),
1003 ("put_sync_root", root),
1004 ] {
1005 self.sync_directory(operation, directory)
1006 .map_err(|error| map_io_err(error, operation))?;
1007 }
1008 Ok(())
1009 }
1010
1011 fn sync_directory(&self, operation: &'static str, directory: &fs::File) -> std::io::Result<()> {
1012 self.io_step(operation, || {
1013 sync_directory(directory)?;
1014 #[cfg(test)]
1015 if let Some(hook) = &self.hook {
1016 hook.directory_synced(operation, directory)?;
1017 }
1018 Ok(())
1019 })
1020 }
1021}
1022
1023#[cfg(unix)]
1024fn sync_directory(directory: &fs::File) -> std::io::Result<()> {
1025 use std::os::fd::AsRawFd;
1026
1027 loop {
1028 if unsafe { libc::fsync(directory.as_raw_fd()) } == 0 {
1032 return Ok(());
1033 }
1034 let error = std::io::Error::last_os_error();
1035 if error.kind() != std::io::ErrorKind::Interrupted {
1036 return Err(error);
1037 }
1038 }
1039}
1040
1041#[cfg(unix)]
1042fn publish_blob_at(
1043 root: &fs::File,
1044 shard1: &fs::File,
1045 shard2: &fs::File,
1046 source: &fs::File,
1047 temp_name: &str,
1048 content_ref: &ContentRef,
1049 publication: &BlobPublication,
1050) -> StorageResult<()> {
1051 use std::os::fd::AsRawFd;
1052
1053 if let Err(error) = publication.step("put_persist", || {
1054 rename_entry_at(
1055 source.as_raw_fd(),
1056 temp_name,
1057 shard2.as_raw_fd(),
1058 content_ref.as_str(),
1059 )
1060 }) {
1061 let _ = unlink_entry_at(source.as_raw_fd(), temp_name);
1062 return Err(error);
1063 }
1064 publication.sync_directories(root, shard1, shard2)
1067}
1068
1069#[cfg(unix)]
1070fn put_blocking_from_root_handle(
1071 root: &Path,
1072 root_handle: &std::fs::File,
1073 floor_bytes: u64,
1074 bytes: Vec<u8>,
1075 publication: &BlobPublication,
1076) -> StorageResult<ContentRef> {
1077 use std::os::unix::io::AsRawFd;
1078
1079 let _root_write_guard = acquire_root_write_lock_anchored(root, root_handle)?;
1080 let digest = blake3::hash(&bytes);
1081 let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
1082
1083 let hex = content_ref.as_str();
1088 let existing = (|| -> std::io::Result<_> {
1089 let shard1 = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
1090 let shard2 = openat_dir_no_follow(shard1.as_raw_fd(), &hex[2..4])?;
1091 let file = openat_regular_file_no_follow(shard2.as_raw_fd(), hex, libc::O_WRONLY)?;
1092 Ok((file, shard1, shard2))
1093 })();
1094 match existing {
1095 Ok((file, shard1, shard2)) => {
1096 file.set_modified(SystemTime::now())
1097 .map_err(|e| map_io_err(e, "put_touch_mtime"))?;
1098 publication.step("put_fsync", || file.sync_all())?;
1099 publication.sync_directories(root_handle, &shard1, &shard2)?;
1100 return Ok(content_ref);
1101 }
1102 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
1103 Err(error) => return Err(map_io_err(error, "put_touch_open")),
1104 }
1105
1106 let required_write_bytes = bytes.len() as u64;
1107 let available =
1108 available_space_at(root_handle).map_err(|e| map_io_err(e, "put_check_space"))?;
1109 if crosses_floor(available, required_write_bytes, floor_bytes) {
1110 return Err(StorageError::CapacityFloor {
1111 capability: StorageCapability::Blob,
1112 volume: root.display().to_string(),
1113 available_bytes: available,
1114 floor_bytes,
1115 });
1116 }
1117
1118 let shard1_dir = open_or_create_dir_at_no_follow(root_handle.as_raw_fd(), &hex[0..2])
1119 .map_err(|e| map_io_err(e, "put_mkdir"))?;
1120 let shard2_dir = open_or_create_dir_at_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])
1121 .map_err(|e| map_io_err(e, "put_mkdir"))?;
1122
1123 let temp_name = format!(".tmp-{}", Uuid::new_v4());
1124 let mut temp = create_regular_file_at_no_follow(shard2_dir.as_raw_fd(), &temp_name, 0o600)
1125 .map_err(|e| map_io_err(e, "put_tempfile"))?;
1126 let write_result = (|| -> StorageResult<()> {
1127 temp.write_all(&bytes)
1128 .map_err(|e| map_io_err(e, "put_write"))?;
1129 temp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
1130 publication.step("put_fsync", || temp.sync_all())?;
1131
1132 let written_len = temp
1133 .metadata()
1134 .map_err(|e| map_io_err(e, "put_verify"))?
1135 .len();
1136 if written_len != bytes.len() as u64 {
1137 return Err(map_io_err(
1138 std::io::Error::other(format!(
1139 "temp file length {written_len} does not match {} written bytes",
1140 bytes.len()
1141 )),
1142 "put_verify",
1143 ));
1144 }
1145 Ok(())
1146 })();
1147 drop(temp);
1148 if let Err(error) = write_result {
1149 let _ = unlink_entry_at(shard2_dir.as_raw_fd(), &temp_name);
1150 return Err(error);
1151 }
1152
1153 publish_blob_at(
1154 root_handle,
1155 &shard1_dir,
1156 &shard2_dir,
1157 &shard2_dir,
1158 &temp_name,
1159 &content_ref,
1160 publication,
1161 )?;
1162 Ok(content_ref)
1163}
1164
1165#[cfg(not(unix))]
1166fn put_blocking_from_root_handle(
1167 root: &Path,
1168 root_handle: &std::fs::File,
1169 floor_bytes: u64,
1170 bytes: Vec<u8>,
1171) -> StorageResult<ContentRef> {
1172 verify_blob_root_identity(root, root_handle).map_err(|e| map_io_err(e, "put_root_identity"))?;
1175 put_blocking(root, floor_bytes, bytes)
1176}
1177
1178#[cfg(any(test, not(unix)))]
1179fn acquire_root_write_lock(root: &Path) -> StorageResult<fs::File> {
1180 let lock_file = fs::OpenOptions::new()
1181 .read(true)
1182 .write(true)
1183 .create(true)
1184 .truncate(false)
1185 .open(root.join(ROOT_WRITE_LOCK_FILE))
1186 .map_err(|e| map_io_err(e, "root_write_lock_open"))?;
1187 fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
1188 Ok(lock_file)
1189}
1190
1191fn database_gc_lock_path(database_path: &Path) -> PathBuf {
1192 let mut lock_path = database_path.as_os_str().to_os_string();
1193 lock_path.push(DATABASE_GC_LOCK_SUFFIX);
1194 PathBuf::from(lock_path)
1195}
1196
1197fn acquire_database_gc_lock(database_path: Option<&Path>) -> StorageResult<Option<fs::File>> {
1198 let Some(database_path) = database_path else {
1199 return Ok(None);
1202 };
1203 let lock_path = database_gc_lock_path(database_path);
1204 let lock_file = fs::OpenOptions::new()
1205 .read(true)
1206 .write(true)
1207 .create(true)
1208 .truncate(false)
1209 .open(&lock_path)
1210 .map_err(|e| map_io_err(e, "database_gc_lock_open"))?;
1211 fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "database_gc_lock_acquire"))?;
1212 Ok(Some(lock_file))
1213}
1214
1215#[cfg(all(unix, target_os = "macos"))]
1216fn errno_location() -> *mut libc::c_int {
1217 unsafe { libc::__error() }
1221}
1222
1223#[cfg(all(unix, not(target_os = "macos")))]
1224fn errno_location() -> *mut libc::c_int {
1225 unsafe { libc::__errno_location() }
1228}
1229
1230#[cfg(unix)]
1234fn clear_errno() {
1235 unsafe { *errno_location() = 0 };
1238}
1239
1240#[cfg(unix)]
1241fn current_errno() -> libc::c_int {
1242 unsafe { *errno_location() }
1244}
1245
1246#[cfg(unix)]
1271fn read_dir_names_no_follow(dir_fd: std::os::unix::io::RawFd) -> std::io::Result<Vec<String>> {
1272 use std::os::unix::io::IntoRawFd;
1273
1274 let reopened = openat_dir_no_follow(dir_fd, ".")?;
1275 let owned_fd = reopened.into_raw_fd();
1279 let dirp = unsafe { libc::fdopendir(owned_fd) };
1282 if dirp.is_null() {
1283 let err = std::io::Error::last_os_error();
1284 unsafe { libc::close(owned_fd) };
1286 return Err(err);
1287 }
1288 let mut names = Vec::new();
1289 loop {
1290 clear_errno();
1293 let entry = unsafe { libc::readdir(dirp) };
1295 if entry.is_null() {
1296 if current_errno() != 0 {
1297 let err = std::io::Error::last_os_error();
1298 unsafe { libc::closedir(dirp) };
1302 return Err(err);
1303 }
1304 break;
1305 }
1306 let first = unsafe { *(*entry).d_name.as_ptr() };
1311 if first == b'.' as libc::c_char {
1312 continue;
1313 }
1314 let name = unsafe { std::ffi::CStr::from_ptr((*entry).d_name.as_ptr()) }
1317 .to_string_lossy()
1318 .into_owned();
1319 names.push(name);
1320 }
1321 unsafe { libc::closedir(dirp) };
1323 Ok(names)
1324}
1325
1326#[cfg(unix)]
1344fn walk_blob_files_from_root_handle(
1345 root_handle: &std::fs::File,
1346 _root: &Path,
1350) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
1351 use std::os::unix::io::AsRawFd;
1352
1353 let mut out = Vec::new();
1354 for l1_name in read_dir_names_no_follow(root_handle.as_raw_fd())? {
1355 let l1_dir = match openat_dir_no_follow(root_handle.as_raw_fd(), &l1_name) {
1356 Ok(dir) => dir,
1357 Err(_) => continue,
1358 };
1359 for l2_name in read_dir_names_no_follow(l1_dir.as_raw_fd())? {
1360 let l2_dir = match openat_dir_no_follow(l1_dir.as_raw_fd(), &l2_name) {
1361 Ok(dir) => dir,
1362 Err(_) => continue,
1363 };
1364 for leaf_name in read_dir_names_no_follow(l2_dir.as_raw_fd())? {
1365 let Ok(content_ref) = ContentRef::from_hex(leaf_name.clone()) else {
1368 continue;
1369 };
1370 #[cfg(all(test, unix))]
1371 if let Some(hook) = walk_leaf_sync_hook::take(_root) {
1372 let _ = hook.reached.send(());
1373 let _ = hook.release.recv();
1374 }
1375 let file = match openat_regular_file_no_follow(
1376 l2_dir.as_raw_fd(),
1377 &leaf_name,
1378 libc::O_RDONLY,
1379 ) {
1380 Ok(file) => file,
1381 Err(_) => continue,
1382 };
1383 let mtime = file.metadata().ok().and_then(|meta| meta.modified().ok());
1384 out.push((content_ref, mtime));
1385 }
1386 }
1387 }
1388 Ok(out)
1389}
1390
1391#[cfg(not(unix))]
1399fn walk_blob_files_from_root_handle(
1400 _root_handle: &std::fs::File,
1401 _root: &Path,
1402) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
1403 Err(std::io::Error::new(
1404 std::io::ErrorKind::Unsupported,
1405 "orphan sweep candidate enumeration requires descriptor-relative directory \
1406 reads, available only on unix in this release; refusing to classify via \
1407 path-based reads",
1408 ))
1409}
1410
1411fn within_publish_grace(
1423 mtime: Option<SystemTime>,
1424 now: SystemTime,
1425 grace_period: Duration,
1426) -> bool {
1427 let age = mtime.and_then(|mtime| now.duration_since(mtime).ok());
1428 match age {
1429 Some(age) => age < grace_period,
1430 None => true,
1431 }
1432}
1433
1434#[derive(Debug)]
1435struct PreparedTransactionalSweep {
1436 result: BlobOrphanSweepResult,
1437 candidates: Vec<(ContentRef, bool)>,
1438}
1439
1440fn prepare_transactional_sweep(
1443 files: Vec<(ContentRef, Option<SystemTime>)>,
1444 grace_period: Duration,
1445) -> PreparedTransactionalSweep {
1446 let now = SystemTime::now();
1447 let mut result = BlobOrphanSweepResult::default();
1448 let mut candidates = Vec::with_capacity(files.len());
1449 for (content_ref, mtime) in files {
1450 result.scanned += 1;
1451 let within_grace = within_publish_grace(mtime, now, grace_period);
1452 candidates.push((content_ref, within_grace));
1453 }
1454 PreparedTransactionalSweep { result, candidates }
1455}
1456
1457#[derive(Debug)]
1458struct BlobGcBatchRows {
1459 grace_period_skipped: u64,
1460 would_delete: u64,
1461 claimed_rows: Vec<SqlRow>,
1462}
1463
1464fn required_nonnegative_count(
1465 value: Option<SqlValue>,
1466 operation: &'static str,
1467) -> StorageResult<u64> {
1468 match value {
1469 Some(SqlValue::Integer(value)) if value >= 0 => Ok(value as u64),
1470 other => Err(StorageError::Internal(format!(
1471 "{operation} returned an invalid count: {other:?}"
1472 ))),
1473 }
1474}
1475
1476fn invalid_content_ref(message: String) -> StorageError {
1477 StorageError::InvalidInput {
1478 capability: StorageCapability::Blob,
1479 operation: "transactional_orphan_sweep".into(),
1480 message,
1481 }
1482}
1483
1484async fn blob_gc_fencing_complete(sql: &dyn SqlAccess) -> StorageResult<bool> {
1498 let mut reader = sql.reader().await?;
1499 let present = required_nonnegative_count(
1500 reader
1501 .query_scalar(SqlStatement {
1502 sql: "SELECT COUNT(*) FROM sqlite_master \
1503 WHERE (type = 'table' AND name IN ( \
1504 'blob_gc_claims', 'attachments', \
1505 'attachment_cutover_state')) \
1506 OR (type = 'index' AND name IN ( \
1507 'idx_blob_gc_claims_content_ref', \
1508 'idx_attachments_content_ref')) \
1509 OR (type = 'trigger' AND name IN ( \
1510 'attachments_reject_claimed_blob_insert', \
1511 'attachments_reject_claimed_blob_update'))"
1512 .to_string(),
1513 params: vec![],
1514 label: Some("blob_gc_fencing_complete".to_string()),
1515 })
1516 .await?,
1517 "blob_gc_fencing_complete",
1518 )?;
1519 if present != 7 {
1520 return Ok(false);
1521 }
1522
1523 let legacy_objects = required_nonnegative_count(
1524 reader
1525 .query_scalar(SqlStatement {
1526 sql: "SELECT \
1527 (SELECT COUNT(*) FROM pragma_table_info('entities') \
1528 WHERE name = 'content_ref') \
1529 + (SELECT COUNT(*) FROM sqlite_master \
1530 WHERE (type = 'index' AND name = 'idx_entities_content_ref') \
1531 OR (type = 'trigger' AND name IN ( \
1532 'entities_reject_claimed_blob_insert', \
1533 'entities_reject_claimed_blob_update')))"
1534 .to_string(),
1535 params: vec![],
1536 label: Some("blob_gc_legacy_fencing_absent".to_string()),
1537 })
1538 .await?,
1539 "blob_gc_legacy_fencing_absent",
1540 )?;
1541 if legacy_objects != 0 {
1542 return Ok(false);
1543 }
1544
1545 let complete = required_nonnegative_count(
1570 reader
1571 .query_scalar(SqlStatement {
1572 sql: "SELECT COUNT(*) FROM attachment_cutover_state AS cutover \
1573 WHERE cutover.singleton = 1 \
1574 AND cutover.state = 'complete' \
1575 AND cutover.completed_at IS NOT NULL \
1576 AND (SELECT COUNT(*) FROM _schema_migrations \
1577 WHERE version = ?1 \
1578 AND name = 'attachments_first_class') = 1 \
1579 AND (SELECT COUNT(*) FROM _schema_migrations) = ?1 \
1580 AND (SELECT MIN(version) FROM _schema_migrations) = 1 \
1581 AND (SELECT MAX(version) FROM _schema_migrations) = ?1"
1582 .to_string(),
1583 params: vec![SqlValue::Integer(i64::from(
1584 crate::migrations::ATTACHMENT_CUTOVER_VERSION,
1585 ))],
1586 label: Some("blob_gc_cutover_complete".to_string()),
1587 })
1588 .await?,
1589 "blob_gc_cutover_complete",
1590 )?;
1591 Ok(complete == 1)
1592}
1593
1594fn unsupported_blob_gc_epoch() -> StorageError {
1595 StorageError::Unsupported {
1596 capability: StorageCapability::Blob,
1597 operation: "transactional_orphan_sweep".into(),
1598 message: "transactional blob GC requires a complete V21 attachment cutover with \
1599 the attachment claim-fencing set; refusing both report-only and \
1600 destructive sweep in this database epoch"
1601 .into(),
1602 }
1603}
1604
1605const BLOB_GC_FENCE_PROBE_REF: &str =
1609 "0000000000000000000000000000000000000000000000000000000000000000";
1610
1611const BLOB_GC_FENCE_TRIGGER_MESSAGE: &str = "content_ref is reserved by an active blob sweep";
1614
1615async fn blob_gc_fence_probe(sql: &dyn SqlAccess) -> StorageResult<()> {
1630 let run = Uuid::new_v4().simple().to_string();
1631 blob_gc_fence_probe_with_ids(
1632 sql,
1633 format!("__blob-gc-fence-probe-insert-{run}__"),
1634 format!("__blob-gc-fence-probe-update-{run}__"),
1635 format!("__blob-gc-fence-probe-insert2-{run}__"),
1636 format!("__blob-gc-fence-probe-update2-{run}__"),
1637 format!("__fence_probe-{run}__"),
1638 )
1639 .await
1640}
1641
1642async fn blob_gc_fence_probe_with_ids(
1647 sql: &dyn SqlAccess,
1648 insert_id: String,
1649 update_id: String,
1650 insert2_id: String,
1651 update2_id: String,
1652 claim_key: String,
1653) -> StorageResult<()> {
1654 fn fence_rejection(result: Result<u64, StorageError>) -> Result<bool, String> {
1655 match result {
1656 Ok(_) => Ok(false),
1657 Err(error) => {
1658 let text = error.to_string();
1659 if text.contains(BLOB_GC_FENCE_TRIGGER_MESSAGE) {
1660 Ok(true)
1661 } else {
1662 Err(text)
1663 }
1664 }
1665 }
1666 }
1667
1668 fn required_seed(value: Option<SqlValue>) -> StorageResult<String> {
1669 match value {
1670 Some(SqlValue::Text(seed)) => Ok(seed),
1671 _ => Err(StorageError::Unsupported {
1672 capability: StorageCapability::Blob,
1673 operation: "transactional_orphan_sweep".into(),
1674 message: "the blob GC fence probe could not select an unclaimed \
1675 canonical seed; refusing deletion so a later sweep can retry"
1676 .into(),
1677 }),
1678 }
1679 }
1680
1681 let op: AtomicUnitOp = Box::new(move |writer| {
1682 Box::pin(async move {
1683 let preexisting = writer
1687 .query_row(SqlStatement {
1688 sql: "SELECT (SELECT COUNT(*) FROM attachments \
1689 WHERE record_uuid IN (?1, ?2, ?3, ?4)) \
1690 + (SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?5)"
1691 .to_string(),
1692 params: vec![
1693 SqlValue::Text(insert_id.clone()),
1694 SqlValue::Text(update_id.clone()),
1695 SqlValue::Text(insert2_id.clone()),
1696 SqlValue::Text(update2_id.clone()),
1697 SqlValue::Text(claim_key.clone()),
1698 ],
1699 label: Some("blob_gc_fence_probe_ownership_guard".to_string()),
1700 })
1701 .await?
1702 .and_then(|row| row.columns.first().map(|c| c.value.clone()));
1703 match preexisting {
1704 Some(SqlValue::Integer(0)) => {}
1705 Some(SqlValue::Integer(_)) => {
1706 return Err(StorageError::Unsupported {
1707 capability: StorageCapability::Blob,
1708 operation: "transactional_orphan_sweep".into(),
1709 message: "the blob GC fence probe's row ids collide with existing \
1710 rows; refusing to probe rather than delete data the \
1711 probe does not own"
1712 .into(),
1713 });
1714 }
1715 _ => {
1716 return Err(StorageError::Internal(
1717 "blob GC fence probe ownership guard returned no count".into(),
1718 ));
1719 }
1720 }
1721
1722 let seed_ref = writer
1728 .query_row(SqlStatement {
1729 sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1730 SELECT 1, lower(hex(randomblob(32))) \
1731 UNION ALL \
1732 SELECT attempt + 1, lower(hex(randomblob(32))) \
1733 FROM candidates WHERE attempt < 8 \
1734 ) \
1735 SELECT candidate.content_ref FROM candidates AS candidate \
1736 WHERE candidate.content_ref <> ?1 \
1737 AND NOT EXISTS ( \
1738 SELECT 1 FROM blob_gc_claims \
1739 WHERE content_ref = candidate.content_ref \
1740 ) \
1741 LIMIT 1"
1742 .to_string(),
1743 params: vec![SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string())],
1744 label: Some("blob_gc_fence_probe_select_seed".to_string()),
1745 })
1746 .await?
1747 .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1748 let seed_ref = required_seed(seed_ref)?;
1749
1750 let seed2_ref = writer
1751 .query_row(SqlStatement {
1752 sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1753 SELECT 1, lower(hex(randomblob(32))) \
1754 UNION ALL \
1755 SELECT attempt + 1, lower(hex(randomblob(32))) \
1756 FROM candidates WHERE attempt < 8 \
1757 ) \
1758 SELECT candidate.content_ref FROM candidates AS candidate \
1759 WHERE candidate.content_ref NOT IN (?1, ?2) \
1760 AND NOT EXISTS ( \
1761 SELECT 1 FROM blob_gc_claims \
1762 WHERE content_ref = candidate.content_ref \
1763 ) \
1764 LIMIT 1"
1765 .to_string(),
1766 params: vec![
1767 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1768 SqlValue::Text(seed_ref.clone()),
1769 ],
1770 label: Some("blob_gc_fence_probe_select_seed2".to_string()),
1771 })
1772 .await?
1773 .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1774 let seed2_ref = required_seed(seed2_ref)?;
1775
1776 let probe2_ref = writer
1780 .query_row(SqlStatement {
1781 sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1782 SELECT 1, lower(hex(randomblob(32))) \
1783 UNION ALL \
1784 SELECT attempt + 1, lower(hex(randomblob(32))) \
1785 FROM candidates WHERE attempt < 8 \
1786 ) \
1787 SELECT candidate.content_ref FROM candidates AS candidate \
1788 WHERE candidate.content_ref NOT IN (?1, ?2, ?3) \
1789 AND NOT EXISTS ( \
1790 SELECT 1 FROM blob_gc_claims \
1791 WHERE content_ref = candidate.content_ref \
1792 ) \
1793 LIMIT 1"
1794 .to_string(),
1795 params: vec![
1796 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1797 SqlValue::Text(seed_ref.clone()),
1798 SqlValue::Text(seed2_ref.clone()),
1799 ],
1800 label: Some("blob_gc_fence_probe_select_probe2".to_string()),
1801 })
1802 .await?
1803 .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1804 let probe2_ref = required_seed(probe2_ref)?;
1805
1806 writer
1807 .execute(SqlStatement {
1808 sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
1809 VALUES (?1, ?2, 0), (?1, ?3, 0)"
1810 .to_string(),
1811 params: vec![
1812 SqlValue::Text(claim_key.clone()),
1813 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1814 SqlValue::Text(probe2_ref.clone()),
1815 ],
1816 label: Some("blob_gc_fence_probe_claim".to_string()),
1817 })
1818 .await?;
1819
1820 let insert_attempt = writer
1821 .execute(SqlStatement {
1822 sql: "INSERT INTO attachments \
1823 (record_uuid, substrate, role, content_ref, created_at) \
1824 VALUES (?1, 'entity', 'content', ?2, 0)"
1825 .to_string(),
1826 params: vec![
1827 SqlValue::Text(insert_id.clone()),
1828 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1829 ],
1830 label: Some("blob_gc_fence_probe_insert_arm".to_string()),
1831 })
1832 .await;
1833 let insert_fenced = fence_rejection(insert_attempt);
1834
1835 writer
1836 .execute(SqlStatement {
1837 sql: "INSERT INTO attachments \
1838 (record_uuid, substrate, role, content_ref, created_at) \
1839 VALUES (?1, 'entity', 'content', ?2, 0)"
1840 .to_string(),
1841 params: vec![SqlValue::Text(update_id.clone()), SqlValue::Text(seed_ref)],
1842 label: Some("blob_gc_fence_probe_update_arm_seed".to_string()),
1843 })
1844 .await?;
1845 let update_attempt = writer
1846 .execute(SqlStatement {
1847 sql: "UPDATE attachments SET content_ref = ?1 \
1848 WHERE record_uuid = ?2 AND role = 'content'"
1849 .to_string(),
1850 params: vec![
1851 SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1852 SqlValue::Text(update_id.clone()),
1853 ],
1854 label: Some("blob_gc_fence_probe_update_arm".to_string()),
1855 })
1856 .await;
1857 let update_fenced = fence_rejection(update_attempt);
1858
1859 let insert2_attempt = writer
1860 .execute(SqlStatement {
1861 sql: "INSERT INTO attachments \
1862 (record_uuid, substrate, role, content_ref, created_at) \
1863 VALUES (?1, 'note', 'evidence', ?2, 0)"
1864 .to_string(),
1865 params: vec![
1866 SqlValue::Text(insert2_id.clone()),
1867 SqlValue::Text(probe2_ref.clone()),
1868 ],
1869 label: Some("blob_gc_fence_probe_insert2_arm".to_string()),
1870 })
1871 .await;
1872 let insert2_fenced = fence_rejection(insert2_attempt);
1873
1874 writer
1875 .execute(SqlStatement {
1876 sql: "INSERT INTO attachments \
1877 (record_uuid, substrate, role, content_ref, created_at) \
1878 VALUES (?1, 'note', 'evidence', ?2, 0)"
1879 .to_string(),
1880 params: vec![
1881 SqlValue::Text(update2_id.clone()),
1882 SqlValue::Text(seed2_ref),
1883 ],
1884 label: Some("blob_gc_fence_probe_update2_arm_seed".to_string()),
1885 })
1886 .await?;
1887 let update2_attempt = writer
1888 .execute(SqlStatement {
1889 sql: "UPDATE attachments SET content_ref = ?1 \
1890 WHERE record_uuid = ?2 AND role = 'evidence'"
1891 .to_string(),
1892 params: vec![
1893 SqlValue::Text(probe2_ref),
1894 SqlValue::Text(update2_id.clone()),
1895 ],
1896 label: Some("blob_gc_fence_probe_update2_arm".to_string()),
1897 })
1898 .await;
1899 let update2_fenced = fence_rejection(update2_attempt);
1900
1901 writer
1904 .execute(SqlStatement {
1905 sql: "DELETE FROM attachments WHERE record_uuid IN (?1, ?2, ?3, ?4)"
1906 .to_string(),
1907 params: vec![
1908 SqlValue::Text(insert_id.clone()),
1909 SqlValue::Text(update_id.clone()),
1910 SqlValue::Text(insert2_id.clone()),
1911 SqlValue::Text(update2_id.clone()),
1912 ],
1913 label: Some("blob_gc_fence_probe_cleanup_attachments".to_string()),
1914 })
1915 .await?;
1916 writer
1917 .execute(SqlStatement {
1918 sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
1919 params: vec![SqlValue::Text(claim_key)],
1920 label: Some("blob_gc_fence_probe_cleanup_claim".to_string()),
1921 })
1922 .await?;
1923
1924 Ok(
1925 Box::new((insert_fenced, update_fenced, insert2_fenced, update2_fenced))
1926 as Box<dyn std::any::Any + Send>,
1927 )
1928 })
1929 });
1930 let outcome = sql.atomic_unit(op).await?;
1931 let (insert_fenced, update_fenced, insert2_fenced, update2_fenced) = *outcome
1932 .downcast::<(
1933 Result<bool, String>,
1934 Result<bool, String>,
1935 Result<bool, String>,
1936 Result<bool, String>,
1937 )>()
1938 .map_err(|_| {
1939 StorageError::Internal("blob GC fence probe returned an unexpected outcome type".into())
1940 })?;
1941 let arm_verdict = |arm: &str, fenced: Result<bool, String>| -> StorageResult<()> {
1942 match fenced {
1943 Ok(true) => Ok(()),
1944 Ok(false) => Err(StorageError::Unsupported {
1945 capability: StorageCapability::Blob,
1946 operation: "transactional_orphan_sweep".into(),
1947 message: format!(
1948 "the V21 fencing triggers exist by name but did not reject a claimed \
1949 content_ref on the attachment {arm} path; refusing unfenced deletion"
1950 ),
1951 }),
1952 Err(other) => Err(StorageError::Unsupported {
1953 capability: StorageCapability::Blob,
1954 operation: "transactional_orphan_sweep".into(),
1955 message: format!(
1956 "the blob GC fence probe could not verify the attachment {arm} fence \
1957 (unexpected rejection: {other}); refusing unfenced deletion"
1958 ),
1959 }),
1960 }
1961 };
1962 arm_verdict("INSERT", insert_fenced)?;
1963 arm_verdict("UPDATE", update_fenced)?;
1964 arm_verdict("second-digest INSERT", insert2_fenced)?;
1965 arm_verdict("second-digest UPDATE", update2_fenced)
1966}
1967
1968async fn validate_blob_gc_evidence(sql: &dyn SqlAccess) -> StorageResult<()> {
1969 let mut reader = sql.reader().await?;
1974 let canonical_bytes = match reader
1983 .query_row(SqlStatement {
1984 sql: "SELECT length(CAST('x' AS BLOB))".to_string(),
1985 params: vec![],
1986 label: Some("blob_gc_validate_encoding_width".to_string()),
1987 })
1988 .await?
1989 .and_then(|row| row.columns.first().map(|column| column.value.clone()))
1990 {
1991 Some(SqlValue::Integer(width)) if (1..=4).contains(&width) => width * 64,
1992 other => {
1993 return Err(invalid_content_ref(format!(
1994 "the text-encoding width probe returned {other:?}; refusing GC validation"
1995 )));
1996 }
1997 };
1998 let invalid_claim = reader
1999 .query_row(SqlStatement {
2000 sql: "SELECT content_ref FROM blob_gc_claims \
2001 WHERE typeof(content_ref) <> 'text' \
2002 OR length(content_ref) <> 64 \
2003 OR length(CAST(content_ref AS BLOB)) <> ?1 \
2004 OR content_ref GLOB '*[^0-9a-f]*' \
2005 LIMIT 1"
2006 .to_string(),
2007 params: vec![SqlValue::Integer(canonical_bytes)],
2008 label: Some("blob_gc_validate_existing_claims".to_string()),
2009 })
2010 .await?;
2011 if invalid_claim.is_some() {
2012 return Err(invalid_content_ref(
2013 "blob_gc_claims.content_ref contained a non-canonical value".into(),
2014 ));
2015 }
2016
2017 let invalid_live = reader
2018 .query_row(SqlStatement {
2019 sql: "SELECT content_ref FROM attachments \
2020 WHERE typeof(content_ref) <> 'text' \
2021 OR length(content_ref) <> 64 \
2022 OR length(CAST(content_ref AS BLOB)) <> ?1 \
2023 OR content_ref GLOB '*[^0-9a-f]*' \
2024 LIMIT 1"
2025 .to_string(),
2026 params: vec![SqlValue::Integer(canonical_bytes)],
2027 label: Some("blob_gc_validate_live_refs".to_string()),
2028 })
2029 .await?;
2030 if invalid_live.is_some() {
2031 return Err(invalid_content_ref(
2032 "attachments.content_ref contained a non-canonical value".into(),
2033 ));
2034 }
2035 Ok(())
2036}
2037
2038async fn release_abandoned_blob_gc_claim_batch(sql: &dyn SqlAccess) -> StorageResult<u64> {
2039 let op: AtomicUnitOp = Box::new(move |writer| {
2040 Box::pin(async move {
2041 let released = writer
2042 .execute(SqlStatement {
2043 sql: "DELETE FROM blob_gc_claims \
2044 WHERE rowid IN ( \
2045 SELECT rowid FROM blob_gc_claims \
2046 ORDER BY rowid LIMIT ?1 \
2047 )"
2048 .to_string(),
2049 params: vec![SqlValue::Integer(BLOB_GC_CLAIM_BATCH_SIZE as i64)],
2050 label: Some("blob_gc_release_abandoned_claim_batch".to_string()),
2051 })
2052 .await?;
2053 Ok(Box::new(released) as Box<dyn std::any::Any + Send>)
2054 })
2055 });
2056 let released = sql.atomic_unit(op).await?;
2057 released.downcast::<u64>().map(|count| *count).map_err(|_| {
2058 StorageError::Internal(
2059 "transactional orphan sweep returned an unexpected recovery count type".into(),
2060 )
2061 })
2062}
2063
2064async fn claim_blob_gc_batch(
2065 sql: &dyn SqlAccess,
2066 root_key: String,
2067 candidates: &[(ContentRef, bool)],
2068 dry_run: bool,
2069) -> StorageResult<BlobGcBatchRows> {
2070 debug_assert!(candidates.len() <= BLOB_GC_CLAIM_BATCH_SIZE);
2071 let eligible_refs = candidates
2072 .iter()
2073 .filter(|(_, within_grace)| !within_grace)
2074 .map(|(content_ref, _)| content_ref.to_string())
2075 .collect::<Vec<_>>();
2076 let grace_refs = candidates
2077 .iter()
2078 .filter(|(_, within_grace)| *within_grace)
2079 .map(|(content_ref, _)| content_ref.to_string())
2080 .collect::<Vec<_>>();
2081 let eligible_json = serde_json::to_string(&eligible_refs).map_err(|error| {
2082 StorageError::Internal(format!(
2083 "failed to prepare blob GC eligible candidate batch: {error}"
2084 ))
2085 })?;
2086 let grace_json = serde_json::to_string(&grace_refs).map_err(|error| {
2087 StorageError::Internal(format!(
2088 "failed to prepare blob GC grace candidate batch: {error}"
2089 ))
2090 })?;
2091 let claimed_at = chrono::Utc::now().timestamp_micros();
2092 let op: AtomicUnitOp = Box::new(move |writer| {
2093 Box::pin(async move {
2094 let grace_period_skipped = required_nonnegative_count(
2095 writer
2096 .query_scalar(SqlStatement {
2097 sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
2098 WHERE NOT EXISTS ( \
2099 SELECT 1 FROM attachments \
2100 WHERE content_ref = candidate.value \
2101 )"
2102 .to_string(),
2103 params: vec![SqlValue::Text(grace_json)],
2104 label: Some("blob_gc_count_grace_candidates_batch".to_string()),
2105 })
2106 .await?,
2107 "blob_gc_count_grace_candidates_batch",
2108 )?;
2109
2110 if dry_run {
2111 let would_delete = required_nonnegative_count(
2112 writer
2113 .query_scalar(SqlStatement {
2114 sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
2115 WHERE NOT EXISTS ( \
2116 SELECT 1 FROM attachments \
2117 WHERE content_ref = candidate.value \
2118 )"
2119 .to_string(),
2120 params: vec![SqlValue::Text(eligible_json)],
2121 label: Some("blob_gc_count_dry_run_candidates_batch".to_string()),
2122 })
2123 .await?,
2124 "blob_gc_count_dry_run_candidates_batch",
2125 )?;
2126 return Ok(Box::new(BlobGcBatchRows {
2127 grace_period_skipped,
2128 would_delete,
2129 claimed_rows: Vec::new(),
2130 }) as Box<dyn std::any::Any + Send>);
2131 }
2132
2133 writer
2134 .execute(SqlStatement {
2135 sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
2136 SELECT ?1, candidate.value, ?3 \
2137 FROM json_each(?2) AS candidate \
2138 WHERE NOT EXISTS ( \
2139 SELECT 1 FROM attachments \
2140 WHERE content_ref = candidate.value \
2141 )"
2142 .to_string(),
2143 params: vec![
2144 SqlValue::Text(root_key.clone()),
2145 SqlValue::Text(eligible_json),
2146 SqlValue::Integer(claimed_at),
2147 ],
2148 label: Some("blob_gc_claim_candidate_batch".to_string()),
2149 })
2150 .await?;
2151
2152 let claimed_rows = writer
2153 .query_all(SqlStatement {
2154 sql: "SELECT content_ref FROM blob_gc_claims \
2155 WHERE root_key = ?1 ORDER BY content_ref"
2156 .to_string(),
2157 params: vec![SqlValue::Text(root_key)],
2158 label: Some("blob_gc_claimed_candidate_batch".to_string()),
2159 })
2160 .await?;
2161 Ok(Box::new(BlobGcBatchRows {
2162 grace_period_skipped,
2163 would_delete: claimed_rows.len() as u64,
2164 claimed_rows,
2165 }) as Box<dyn std::any::Any + Send>)
2166 })
2167 });
2168 let rows = sql.atomic_unit(op).await?;
2169 rows.downcast::<BlobGcBatchRows>()
2170 .map(|rows| *rows)
2171 .map_err(|_| {
2172 StorageError::Internal(
2173 "transactional orphan sweep returned an unexpected batch-row type".into(),
2174 )
2175 })
2176}
2177
2178fn parse_blob_gc_claim_rows(rows: Vec<SqlRow>) -> StorageResult<Vec<ContentRef>> {
2179 let mut claimed = Vec::with_capacity(rows.len());
2180 for row in rows {
2181 let raw = match row.get("content_ref") {
2182 Some(SqlValue::Text(raw)) => raw.clone(),
2183 _ => {
2184 return Err(invalid_content_ref(
2185 "blob_gc_claims.content_ref contained a non-text value".into(),
2186 ));
2187 }
2188 };
2189 claimed.push(ContentRef::from_hex(raw).map_err(invalid_content_ref)?);
2190 }
2191 Ok(claimed)
2192}
2193
2194async fn release_blob_gc_batch(sql: &dyn SqlAccess, root_key: String) -> StorageResult<()> {
2195 let cleanup: AtomicUnitOp = Box::new(move |writer| {
2196 Box::pin(async move {
2197 writer
2198 .execute(SqlStatement {
2199 sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
2200 params: vec![SqlValue::Text(root_key)],
2201 label: Some("blob_gc_release_claim_batch".to_string()),
2202 })
2203 .await?;
2204 Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
2205 })
2206 });
2207 sql.atomic_unit(cleanup).await?;
2208 Ok(())
2209}
2210
2211type SweepLockMap = HashMap<Option<PathBuf>, Arc<DatabaseGcProcessLock>>;
2218
2219#[derive(Debug, Default)]
2220struct DatabaseGcProcessLock {
2221 held: StdMutex<bool>,
2222 released: std::sync::Condvar,
2223 #[cfg(test)]
2224 waiters: std::sync::atomic::AtomicUsize,
2225}
2226
2227impl DatabaseGcProcessLock {
2228 fn acquire(self: &Arc<Self>) -> DatabaseGcProcessGuard {
2229 let mut held = self
2230 .held
2231 .lock()
2232 .unwrap_or_else(std::sync::PoisonError::into_inner);
2233 while *held {
2234 #[cfg(test)]
2235 self.waiters
2236 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2237 held = self
2238 .released
2239 .wait(held)
2240 .unwrap_or_else(std::sync::PoisonError::into_inner);
2241 #[cfg(test)]
2242 self.waiters
2243 .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
2244 }
2245 *held = true;
2246 DatabaseGcProcessGuard {
2247 lock: Arc::clone(self),
2248 }
2249 }
2250
2251 fn try_acquire(self: &Arc<Self>) -> Option<DatabaseGcProcessGuard> {
2252 let mut held = self
2253 .held
2254 .lock()
2255 .unwrap_or_else(std::sync::PoisonError::into_inner);
2256 if *held {
2257 return None;
2258 }
2259 *held = true;
2260 Some(DatabaseGcProcessGuard {
2261 lock: Arc::clone(self),
2262 })
2263 }
2264}
2265
2266#[derive(Debug)]
2267struct DatabaseGcProcessGuard {
2268 lock: Arc<DatabaseGcProcessLock>,
2269}
2270
2271impl Drop for DatabaseGcProcessGuard {
2272 fn drop(&mut self) {
2273 let mut held = self
2274 .lock
2275 .held
2276 .lock()
2277 .unwrap_or_else(std::sync::PoisonError::into_inner);
2278 debug_assert!(*held, "database GC process owner released twice");
2279 *held = false;
2280 self.lock.released.notify_one();
2281 }
2282}
2283
2284fn database_sweep_locks() -> &'static StdMutex<SweepLockMap> {
2285 static REGISTRY: OnceLock<StdMutex<SweepLockMap>> = OnceLock::new();
2286 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2287}
2288
2289fn sweep_lock_for_database(database_path: Option<&Path>) -> Arc<DatabaseGcProcessLock> {
2290 let key = database_path.map(Path::to_path_buf);
2291 let mut locks = database_sweep_locks()
2292 .lock()
2293 .unwrap_or_else(std::sync::PoisonError::into_inner);
2294 locks
2295 .entry(key)
2296 .or_insert_with(|| Arc::new(DatabaseGcProcessLock::default()))
2297 .clone()
2298}
2299
2300#[cfg(test)]
2301pub(crate) fn database_gc_waiter_count(database_path: Option<&Path>) -> usize {
2302 sweep_lock_for_database(database_path)
2303 .waiters
2304 .load(std::sync::atomic::Ordering::SeqCst)
2305}
2306
2307pub struct DatabaseGcOwnerGuard {
2315 _process_guard: DatabaseGcProcessGuard,
2316 _advisory_guard: Option<fs::File>,
2317 database_path: Option<PathBuf>,
2318}
2319
2320pub(crate) fn acquire_database_gc_owner_for_path_blocking(
2321 database_path: Option<PathBuf>,
2322) -> StorageResult<DatabaseGcOwnerGuard> {
2323 let process_guard = sweep_lock_for_database(database_path.as_deref()).acquire();
2324 let advisory_guard = acquire_database_gc_lock(database_path.as_deref())?;
2325 Ok(DatabaseGcOwnerGuard {
2326 _process_guard: process_guard,
2327 _advisory_guard: advisory_guard,
2328 database_path,
2329 })
2330}
2331
2332pub(crate) fn try_acquire_database_gc_owner_for_path(
2340 database_path: PathBuf,
2341) -> StorageResult<DatabaseGcOwnerGuard> {
2342 let process_guard = sweep_lock_for_database(Some(&database_path))
2343 .try_acquire()
2344 .ok_or_else(|| {
2345 StorageError::Internal(format!(
2346 "database GC owner for {} is already held; retry schema migration through the \
2347 coordinated backend boot path",
2348 database_path.display()
2349 ))
2350 })?;
2351 let lock_path = database_gc_lock_path(&database_path);
2352 let advisory_guard = fs::OpenOptions::new()
2353 .read(true)
2354 .write(true)
2355 .create(true)
2356 .truncate(false)
2357 .open(&lock_path)
2358 .map_err(|error| map_io_err(error, "database_gc_lock_open"))?;
2359 fs4::FileExt::try_lock(&advisory_guard)
2360 .map_err(|error| map_io_err(error.into(), "database_gc_lock_try_acquire"))?;
2361 Ok(DatabaseGcOwnerGuard {
2362 _process_guard: process_guard,
2363 _advisory_guard: Some(advisory_guard),
2364 database_path: Some(database_path),
2365 })
2366}
2367
2368impl std::fmt::Debug for DatabaseGcOwnerGuard {
2369 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2370 formatter
2371 .debug_struct("DatabaseGcOwnerGuard")
2372 .field("database_path", &self.database_path)
2373 .finish_non_exhaustive()
2374 }
2375}
2376
2377impl DatabaseGcOwnerGuard {
2378 pub fn database_path(&self) -> Option<&Path> {
2381 self.database_path.as_deref()
2382 }
2383}
2384
2385pub async fn acquire_database_gc_owner(sql: &dyn SqlAccess) -> StorageResult<DatabaseGcOwnerGuard> {
2389 let database_path = sql.database_path();
2390 tokio::task::spawn_blocking(move || acquire_database_gc_owner_for_path_blocking(database_path))
2391 .await
2392 .map_err(|error| {
2393 StorageError::driver(StorageCapability::Blob, "acquire_database_gc_owner", error)
2394 })?
2395}
2396
2397fn root_write_locks() -> &'static StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>> {
2408 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>>> =
2409 OnceLock::new();
2410 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2411}
2412
2413fn write_lock_for_root(root: &Path) -> std::io::Result<Arc<tokio::sync::Mutex<()>>> {
2422 let canonical = root.canonicalize()?;
2423 let mut locks = root_write_locks()
2424 .lock()
2425 .unwrap_or_else(std::sync::PoisonError::into_inner);
2426 Ok(locks
2427 .entry(canonical)
2428 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
2429 .clone())
2430}
2431
2432#[derive(Debug)]
2434pub struct FsBlobStore {
2435 root: PathBuf,
2436 root_handle: Arc<fs::File>,
2440 floor_bytes: u64,
2441 write_lock: Arc<tokio::sync::Mutex<()>>,
2458 orphan_sweep_grace: Duration,
2465}
2466
2467impl FsBlobStore {
2468 pub const DEFAULT_FLOOR_BYTES: u64 = 100_000_000_000;
2471
2472 pub const DEFAULT_ORPHAN_SWEEP_GRACE: Duration = Duration::from_secs(3600);
2477
2478 pub fn new(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
2484 match fs::create_dir(&root) {
2485 Ok(()) => {}
2486 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
2487 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
2488 let parent = root
2489 .parent()
2490 .filter(|path| !path.as_os_str().is_empty())
2491 .unwrap_or_else(|| Path::new("."));
2492 return Err(std::io::Error::new(
2493 std::io::ErrorKind::NotFound,
2494 format!("blob root parent missing: {}", parent.display()),
2495 )
2496 .into());
2497 }
2498 Err(error) => return Err(error.into()),
2499 }
2500 let store = Self::open_existing(root, floor_bytes)?;
2501 #[cfg(unix)]
2502 {
2503 use std::os::fd::AsRawFd;
2504
2505 let publication = BlobPublication {
2506 #[cfg(test)]
2507 hook: sync_hook::take(&store.root).and_then(|hook| hook.publication),
2508 };
2509 let parent = openat_dir_no_follow(store.root_handle.as_raw_fd(), "..")?;
2510 publication.sync_directory("init_sync_root", &store.root_handle)?;
2511 publication.sync_directory("init_sync_parent", &parent)?;
2512 }
2513 Ok(store)
2514 }
2515
2516 pub fn open_existing(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
2521 let root = root.canonicalize()?;
2525 let metadata = fs::metadata(&root)?;
2526 if !metadata.is_dir() {
2527 return Err(SqliteError::InvalidData(format!(
2528 "blob store root is not a directory: {}",
2529 root.display()
2530 )));
2531 }
2532 let root_handle = Arc::new(open_blob_root_handle(&root)?);
2533 let write_lock = write_lock_for_root(&root)?;
2534 verify_blob_root_identity(&root, &root_handle)?;
2535 Ok(Self {
2536 root,
2537 root_handle,
2538 floor_bytes,
2539 write_lock,
2540 orphan_sweep_grace: Self::DEFAULT_ORPHAN_SWEEP_GRACE,
2541 })
2542 }
2543
2544 pub fn with_orphan_sweep_grace(mut self, grace_period: Duration) -> Self {
2547 self.orphan_sweep_grace = grace_period;
2548 self
2549 }
2550
2551 pub fn root(&self) -> &Path {
2553 &self.root
2554 }
2555}
2556
2557#[async_trait]
2558impl BlobStore for FsBlobStore {
2559 async fn begin_upload(&self, declared_size: u64) -> StorageResult<UploadId> {
2560 uploads::begin(self, declared_size).await
2561 }
2562
2563 async fn append_part(&self, id: &UploadId, bytes: Vec<u8>) -> StorageResult<u64> {
2564 uploads::append(self, id.clone(), bytes).await
2565 }
2566
2567 async fn commit_upload(&self, id: &UploadId, content_ref: &ContentRef) -> StorageResult<()> {
2568 uploads::commit(self, id.clone(), content_ref.clone()).await
2569 }
2570
2571 async fn abort_upload(&self, id: &UploadId) -> StorageResult<()> {
2572 uploads::abort(self, id.clone()).await
2573 }
2574
2575 async fn sweep_uploads(&self, idle_for: Duration) -> StorageResult<u64> {
2576 uploads::sweep(self, idle_for).await
2577 }
2578
2579 async fn put(&self, bytes: Vec<u8>) -> StorageResult<ContentRef> {
2580 let owned_guard = self.write_lock.clone().lock_owned().await;
2590 let root = self.root.clone();
2591 let root_handle = Arc::clone(&self.root_handle);
2592 let floor_bytes = self.floor_bytes;
2593 #[cfg(test)]
2600 let hook = sync_hook::take(&root);
2601 tokio::task::spawn_blocking(move || {
2602 #[cfg_attr(not(test), allow(clippy::let_and_return))]
2607 let result = {
2608 let _owned_guard = owned_guard;
2609 #[cfg(test)]
2610 if let Some(h) = &hook {
2611 let _ = h.reached.send(());
2612 let _ = h.release.recv();
2613 }
2614 put_blocking_from_root_handle(
2615 &root,
2616 &root_handle,
2617 floor_bytes,
2618 bytes,
2619 #[cfg(unix)]
2620 &BlobPublication {
2621 #[cfg(test)]
2622 hook: hook.as_ref().and_then(|hook| hook.publication.clone()),
2623 },
2624 )
2625 };
2626 #[cfg(test)]
2627 if let Some(h) = &hook {
2628 let _ = h.done.send(());
2629 }
2630 result
2631 })
2632 .await
2633 .map_err(|e| StorageError::driver(StorageCapability::Blob, "put", e))?
2634 }
2635
2636 async fn get_bounded_verified(
2637 &self,
2638 content_ref: &ContentRef,
2639 max_bytes: u64,
2640 ) -> StorageResult<Vec<u8>> {
2641 if max_bytes > MAX_BLOB_WHOLE_BYTES {
2645 return Err(StorageError::InvalidInput {
2646 capability: StorageCapability::Blob,
2647 operation: "get_bounded_verified".into(),
2648 message: format!(
2649 "max_bytes {max_bytes} exceeds the {MAX_BLOB_WHOLE_BYTES}-byte portable whole-buffer envelope"
2650 ),
2651 });
2652 }
2653
2654 let root = self.root.clone();
2655 let root_handle = Arc::clone(&self.root_handle);
2656 let content_ref = content_ref.clone();
2657 #[cfg(test)]
2658 let read_hook = bounded_read_sync_hook::take(&root);
2659 tokio::task::spawn_blocking(move || {
2660 let mut file = open_blob_shard_file_no_follow(&root, &root_handle, &content_ref)
2664 .map_err(|e| {
2665 if e.kind() == std::io::ErrorKind::NotFound {
2666 StorageError::NotFound {
2667 capability: StorageCapability::Blob,
2668 resource: "blob",
2669 key: content_ref.to_string(),
2670 }
2671 } else if e.kind() == std::io::ErrorKind::Unsupported {
2672 StorageError::Unsupported {
2673 capability: StorageCapability::Blob,
2674 operation: "get_bounded_verified".into(),
2675 message: e.to_string(),
2676 }
2677 } else {
2678 map_io_err(e, "get_bounded_verified.open")
2679 }
2680 })?;
2681
2682 let metadata_bytes = file
2683 .metadata()
2684 .map_err(|e| map_io_err(e, "get_bounded_verified.metadata"))?
2685 .len();
2686 #[cfg(test)]
2687 if let Some(hook) = &read_hook {
2688 let _ = hook.reached.send(());
2689 let _ = hook.release.recv();
2690 }
2691 if metadata_bytes > max_bytes {
2692 return Err(StorageError::BlobTooLarge {
2693 content_ref,
2694 max_bytes,
2695 observed_at_least: metadata_bytes,
2696 });
2697 }
2698
2699 let mut bytes = Vec::with_capacity(metadata_bytes as usize);
2702 (&mut file)
2703 .take(max_bytes + 1)
2704 .read_to_end(&mut bytes)
2705 .map_err(|e| map_io_err(e, "get_bounded_verified.read"))?;
2706 let actual_bytes = bytes.len() as u64;
2707 if actual_bytes > max_bytes {
2708 return Err(StorageError::BlobTooLarge {
2709 content_ref,
2710 max_bytes,
2711 observed_at_least: actual_bytes,
2712 });
2713 }
2714 if metadata_bytes != actual_bytes {
2715 return Err(StorageError::BlobSizeMismatch {
2716 content_ref,
2717 metadata_bytes,
2718 actual_bytes,
2719 });
2720 }
2721
2722 let actual = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
2723 if actual != content_ref {
2724 return Err(StorageError::BlobDigestMismatch {
2725 expected: content_ref,
2726 actual,
2727 });
2728 }
2729 Ok(bytes)
2730 })
2731 .await
2732 .map_err(|e| StorageError::driver(StorageCapability::Blob, "get_bounded_verified", e))?
2733 }
2734
2735 async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool> {
2736 let root = self.root.clone();
2737 let root_handle = Arc::clone(&self.root_handle);
2738 let content_ref = content_ref.clone();
2739 tokio::task::spawn_blocking(move || {
2740 match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2741 Ok(_) => Ok(true),
2742 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
2743 Err(error) => Err(map_io_err(error, "exists")),
2744 }
2745 })
2746 .await
2747 .map_err(|e| StorageError::driver(StorageCapability::Blob, "exists", e))?
2748 }
2749
2750 async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>> {
2751 let root = self.root.clone();
2752 let root_handle = Arc::clone(&self.root_handle);
2753 let content_ref = content_ref.clone();
2754 tokio::task::spawn_blocking(move || {
2755 match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2756 Ok(file) => file
2757 .metadata()
2758 .map(|metadata| Some(metadata.len()))
2759 .map_err(|error| map_io_err(error, "size")),
2760 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
2761 Err(error) => Err(map_io_err(error, "size")),
2762 }
2763 })
2764 .await
2765 .map_err(|e| StorageError::driver(StorageCapability::Blob, "size", e))?
2766 }
2767
2768 async fn delete(&self, content_ref: &ContentRef) -> StorageResult<bool> {
2769 let root = self.root.clone();
2770 let root_handle = Arc::clone(&self.root_handle);
2771 let content_ref = content_ref.clone();
2772 tokio::task::spawn_blocking(move || {
2773 match unlink_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2774 Ok(()) => Ok(true),
2775 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
2776 Err(e) => Err(map_io_err(e, "delete")),
2777 }
2778 })
2779 .await
2780 .map_err(|e| StorageError::driver(StorageCapability::Blob, "delete", e))?
2781 }
2782
2783 async fn orphan_sweep(
2793 &self,
2794 config: &BlobOrphanSweepConfig,
2795 ) -> StorageResult<BlobOrphanSweepResult> {
2796 let _ = config;
2797 Err(StorageError::Unsupported {
2798 capability: StorageCapability::Blob,
2799 operation: "orphan_sweep".into(),
2800 message: "caller-snapshot orphan_sweep is disabled in this compatibility release; \
2801 it cannot prove a completed V21 attachment epoch, use \
2802 transactional_orphan_sweep instead"
2803 .into(),
2804 })
2805 }
2806
2807 async fn transactional_orphan_sweep(
2829 &self,
2830 sql: &dyn SqlAccess,
2831 dry_run: bool,
2832 ) -> StorageResult<BlobOrphanSweepResult> {
2833 if !blob_gc_fencing_complete(sql).await? {
2838 return Err(unsupported_blob_gc_epoch());
2839 }
2840
2841 let database_path = sql.database_path();
2849 let lock_database_path = database_path.clone();
2850 #[cfg(test)]
2851 let hook_database_path = database_path.clone();
2852 let (database_guard, database_file_guard) = tokio::task::spawn_blocking(move || {
2853 let process_guard = sweep_lock_for_database(lock_database_path.as_deref()).acquire();
2854 let file_guard = acquire_database_gc_lock(lock_database_path.as_deref())?;
2855 #[cfg(test)]
2856 if let Some(hook) = db_ownership_sync_hook::take(hook_database_path.as_deref()) {
2857 let _ = hook.reached.send(());
2858 let _ = hook.release.recv();
2859 }
2860 Ok::<_, StorageError>((process_guard, file_guard))
2861 })
2862 .await
2863 .map_err(|e| {
2864 StorageError::driver(
2865 StorageCapability::Blob,
2866 "transactional_orphan_sweep_lock",
2867 e,
2868 )
2869 })??;
2870
2871 if !blob_gc_fencing_complete(sql).await? {
2877 return Err(unsupported_blob_gc_epoch());
2878 }
2879
2880 let root_guard = self.write_lock.clone().lock_owned().await;
2881 let root = self.root.clone();
2882 let root_handle = Arc::clone(&self.root_handle);
2883 let scan_root = root.clone();
2884 let scan_root_handle = Arc::clone(&root_handle);
2885 let grace_period = self.orphan_sweep_grace;
2886 let (write_guards, canonical_root, prepared) = tokio::task::spawn_blocking(move || {
2887 verify_blob_root_identity(&scan_root, &scan_root_handle)
2888 .map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
2889 let canonical_root = scan_root;
2892 let root_write_guard =
2893 acquire_root_write_lock_anchored(&canonical_root, &scan_root_handle)?;
2894 let candidates = walk_blob_files_from_root_handle(&scan_root_handle, &canonical_root)
2895 .map_err(|e| map_io_err(e, "transactional_orphan_sweep_walk"))?;
2896 verify_blob_root_identity(&canonical_root, &scan_root_handle)
2897 .map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
2898 let prepared = prepare_transactional_sweep(candidates, grace_period);
2899 Ok::<_, StorageError>((
2900 (
2901 database_guard,
2902 database_file_guard,
2903 root_guard,
2904 root_write_guard,
2905 ),
2906 canonical_root,
2907 prepared,
2908 ))
2909 })
2910 .await
2911 .map_err(|e| {
2912 StorageError::driver(
2913 StorageCapability::Blob,
2914 "transactional_orphan_sweep_walk",
2915 e,
2916 )
2917 })??;
2918 let root_key = blob_root_key(&canonical_root);
2919 validate_blob_gc_evidence(sql).await?;
2920 blob_gc_fence_probe(sql).await?;
2921 if !dry_run {
2922 loop {
2923 let released = release_abandoned_blob_gc_claim_batch(sql).await?;
2924 if released < BLOB_GC_CLAIM_BATCH_SIZE as u64 {
2925 break;
2926 }
2927 }
2928 }
2929
2930 let mut write_guards = write_guards;
2931 let mut result = prepared.result;
2932 let mut delete_error = None;
2933 #[cfg(test)]
2934 let mut hook: Option<sync_hook::Hook> = None;
2935 #[cfg(not(test))]
2936 let mut hook: Option<()> = None;
2937 #[cfg(test)]
2938 let mut hook_paused = false;
2939
2940 for candidates in prepared.candidates.chunks(BLOB_GC_CLAIM_BATCH_SIZE) {
2947 let batch = claim_blob_gc_batch(sql, root_key.clone(), candidates, dry_run).await?;
2948 result.grace_period_skipped += batch.grace_period_skipped;
2949 result.would_delete += batch.would_delete;
2950 if dry_run {
2951 continue;
2952 }
2953
2954 let claimed_refs = parse_blob_gc_claim_rows(batch.claimed_rows)?;
2955 if claimed_refs.is_empty() {
2956 continue;
2957 }
2958
2959 #[cfg(test)]
2960 if hook.is_none() {
2961 hook = sync_hook::take(&root);
2962 }
2963 #[cfg(test)]
2964 let pause_hook = hook.is_some() && !hook_paused;
2965 #[cfg(test)]
2966 if pause_hook {
2967 hook_paused = true;
2968 }
2969
2970 let delete_root = canonical_root.clone();
2971 let delete_root_handle = Arc::clone(&root_handle);
2972 let (returned_guards, deleted, batch_delete_error, returned_hook) =
2973 tokio::task::spawn_blocking(move || {
2974 #[cfg(test)]
2975 if pause_hook {
2976 if let Some(hook) = &hook {
2977 let _ = hook.reached.send(());
2978 let _ = hook.release.recv();
2979 }
2980 }
2981
2982 let mut deleted = 0_u64;
2983 let mut first_error = None;
2984 for content_ref in claimed_refs {
2985 match unlink_blob_shard_file_no_follow(
2986 &delete_root,
2987 &delete_root_handle,
2988 &content_ref,
2989 ) {
2990 Ok(()) => deleted += 1,
2991 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
2992 Err(error) => {
2993 first_error =
2994 Some(map_io_err(error, "transactional_orphan_sweep_delete"));
2995 break;
2996 }
2997 }
2998 }
2999 (write_guards, deleted, first_error, hook)
3004 })
3005 .await
3006 .map_err(|error| {
3007 StorageError::driver(
3008 StorageCapability::Blob,
3009 "transactional_orphan_sweep_delete",
3010 error,
3011 )
3012 })?;
3013 write_guards = returned_guards;
3014 hook = returned_hook;
3015 result.deleted += deleted;
3016
3017 release_blob_gc_batch(sql, root_key.clone()).await?;
3022 if batch_delete_error.is_some() {
3023 delete_error = batch_delete_error;
3024 break;
3025 }
3026 }
3027
3028 drop(write_guards);
3029 #[cfg(test)]
3030 if let Some(hook) = hook {
3031 let _ = hook.done.send(());
3032 }
3033 #[cfg(not(test))]
3034 let _ = hook;
3035 if let Some(error) = delete_error {
3036 return Err(error);
3037 }
3038 Ok(result)
3039 }
3040}
3041
3042#[cfg(test)]
3057mod sync_hook {
3058 use std::collections::{HashMap, VecDeque};
3059 use std::path::{Path, PathBuf};
3060 use std::sync::mpsc::{Receiver, Sender};
3061 use std::sync::{Mutex as StdMutex, OnceLock};
3062
3063 pub(super) struct Hook {
3064 pub(super) reached: Sender<()>,
3065 pub(super) release: Receiver<()>,
3066 pub(super) done: Sender<()>,
3067 #[cfg(unix)]
3068 pub(super) publication: Option<Publication>,
3069 }
3070
3071 fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
3072 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
3073 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
3074 }
3075
3076 pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>, Receiver<()>) {
3079 let canonical = root
3080 .canonicalize()
3081 .expect("root must exist before installing a sync_hook");
3082 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
3083 let (release_tx, release_rx) = std::sync::mpsc::channel();
3084 let (done_tx, done_rx) = std::sync::mpsc::channel();
3085 registry()
3086 .lock()
3087 .unwrap_or_else(std::sync::PoisonError::into_inner)
3088 .entry(canonical)
3089 .or_default()
3090 .push_back(Hook {
3091 reached: reached_tx,
3092 release: release_rx,
3093 done: done_tx,
3094 #[cfg(unix)]
3095 publication: None,
3096 });
3097 (reached_rx, release_tx, done_rx)
3098 }
3099
3100 pub(super) fn take(root: &Path) -> Option<Hook> {
3106 let canonical = root.canonicalize().ok()?;
3107 registry()
3108 .lock()
3109 .unwrap_or_else(std::sync::PoisonError::into_inner)
3110 .get_mut(&canonical)
3111 .and_then(VecDeque::pop_front)
3112 }
3113
3114 #[cfg(unix)]
3115 type StepAction = (&'static str, Box<dyn FnOnce() + Send>);
3116
3117 #[cfg(unix)]
3118 type DirectorySync = (&'static str, u64, u64);
3119
3120 #[cfg(unix)]
3121 #[derive(Clone)]
3122 pub(super) struct Publication {
3123 pub(super) completed: std::sync::Arc<StdMutex<Vec<&'static str>>>,
3124 pub(super) directories: std::sync::Arc<StdMutex<Vec<DirectorySync>>>,
3125 fail_at: Option<&'static str>,
3126 action: std::sync::Arc<StdMutex<Option<StepAction>>>,
3127 }
3128
3129 #[cfg(unix)]
3130 impl Publication {
3131 pub(super) fn before(&self, operation: &'static str) -> std::io::Result<()> {
3132 if self.fail_at == Some(operation) {
3133 return Err(std::io::Error::other("injected publication failure"));
3134 }
3135 let action = {
3136 let mut slot = self.action.lock().unwrap();
3137 if slot.as_ref().is_some_and(|(at, _)| *at == operation) {
3138 slot.take()
3139 } else {
3140 None
3141 }
3142 };
3143 if let Some((_, action)) = action {
3144 action();
3145 }
3146 Ok(())
3147 }
3148
3149 pub(super) fn on_step(
3150 &self,
3151 operation: &'static str,
3152 action: impl FnOnce() + Send + 'static,
3153 ) {
3154 *self.action.lock().unwrap() = Some((operation, Box::new(action)));
3155 }
3156
3157 pub(super) fn completed(&self, operation: &'static str) {
3158 self.completed.lock().unwrap().push(operation);
3159 }
3160
3161 pub(super) fn directory_synced(
3162 &self,
3163 operation: &'static str,
3164 directory: &std::fs::File,
3165 ) -> std::io::Result<()> {
3166 use std::os::unix::fs::MetadataExt;
3167
3168 let metadata = directory.metadata()?;
3169 self.directories
3170 .lock()
3171 .unwrap()
3172 .push((operation, metadata.dev(), metadata.ino()));
3173 Ok(())
3174 }
3175 }
3176
3177 #[cfg(unix)]
3180 pub(super) fn install_publication(root: &Path, fail_at: Option<&'static str>) -> Publication {
3181 let publication = Publication {
3182 completed: std::sync::Arc::default(),
3183 directories: std::sync::Arc::default(),
3184 fail_at,
3185 action: std::sync::Arc::default(),
3186 };
3187 let (reached, _) = std::sync::mpsc::channel();
3188 let (_, release) = std::sync::mpsc::channel();
3189 let (done, _) = std::sync::mpsc::channel();
3190 registry()
3191 .lock()
3192 .unwrap_or_else(std::sync::PoisonError::into_inner)
3193 .entry(root.canonicalize().unwrap())
3194 .or_default()
3195 .push_back(Hook {
3196 reached,
3197 release,
3198 done,
3199 publication: Some(publication.clone()),
3200 });
3201 publication
3202 }
3203}
3204
3205#[cfg(all(test, unix))]
3221mod walk_leaf_sync_hook {
3222 use std::collections::HashMap;
3223 use std::path::{Path, PathBuf};
3224 use std::sync::mpsc::{Receiver, Sender};
3225 use std::sync::{Mutex as StdMutex, OnceLock};
3226
3227 pub(super) struct Hook {
3228 pub(super) reached: Sender<()>,
3229 pub(super) release: Receiver<()>,
3230 }
3231
3232 fn registry() -> &'static StdMutex<HashMap<PathBuf, Hook>> {
3233 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Hook>>> = OnceLock::new();
3234 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
3235 }
3236
3237 pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
3238 let canonical = root
3239 .canonicalize()
3240 .expect("root must exist before installing a walk_leaf_sync_hook");
3241 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
3242 let (release_tx, release_rx) = std::sync::mpsc::channel();
3243 registry()
3244 .lock()
3245 .unwrap_or_else(std::sync::PoisonError::into_inner)
3246 .insert(
3247 canonical,
3248 Hook {
3249 reached: reached_tx,
3250 release: release_rx,
3251 },
3252 );
3253 (reached_rx, release_tx)
3254 }
3255
3256 pub(super) fn take(root: &Path) -> Option<Hook> {
3259 let canonical = root.canonicalize().ok()?;
3260 registry()
3261 .lock()
3262 .unwrap_or_else(std::sync::PoisonError::into_inner)
3263 .remove(&canonical)
3264 }
3265}
3266
3267#[cfg(test)]
3272mod bounded_read_sync_hook {
3273 use std::collections::{HashMap, VecDeque};
3274 use std::path::{Path, PathBuf};
3275 use std::sync::mpsc::{Receiver, Sender};
3276 use std::sync::{Mutex as StdMutex, OnceLock};
3277
3278 pub(super) struct Hook {
3279 pub(super) reached: Sender<()>,
3280 pub(super) release: Receiver<()>,
3281 }
3282
3283 fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
3284 static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
3285 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
3286 }
3287
3288 pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
3289 let canonical = root
3290 .canonicalize()
3291 .expect("root must exist before installing a bounded-read hook");
3292 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
3293 let (release_tx, release_rx) = std::sync::mpsc::channel();
3294 registry()
3295 .lock()
3296 .unwrap_or_else(std::sync::PoisonError::into_inner)
3297 .entry(canonical)
3298 .or_default()
3299 .push_back(Hook {
3300 reached: reached_tx,
3301 release: release_rx,
3302 });
3303 (reached_rx, release_tx)
3304 }
3305
3306 pub(super) fn take(root: &Path) -> Option<Hook> {
3307 let canonical = root.canonicalize().ok()?;
3308 registry()
3309 .lock()
3310 .unwrap_or_else(std::sync::PoisonError::into_inner)
3311 .get_mut(&canonical)
3312 .and_then(VecDeque::pop_front)
3313 }
3314}
3315
3316#[cfg(test)]
3323mod db_ownership_sync_hook {
3324 use std::collections::{HashMap, VecDeque};
3325 use std::path::{Path, PathBuf};
3326 use std::sync::mpsc::{Receiver, Sender};
3327 use std::sync::{Mutex as StdMutex, OnceLock};
3328
3329 pub(super) struct Hook {
3330 pub(super) reached: Sender<()>,
3331 pub(super) release: Receiver<()>,
3332 }
3333
3334 fn registry() -> &'static StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>> {
3335 static REGISTRY: OnceLock<StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>>> =
3336 OnceLock::new();
3337 REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
3338 }
3339
3340 pub(super) fn install(database_path: Option<&Path>) -> (Receiver<()>, Sender<()>) {
3341 let key = database_path.map(Path::to_path_buf);
3342 let (reached_tx, reached_rx) = std::sync::mpsc::channel();
3343 let (release_tx, release_rx) = std::sync::mpsc::channel();
3344 registry()
3345 .lock()
3346 .unwrap_or_else(std::sync::PoisonError::into_inner)
3347 .entry(key)
3348 .or_default()
3349 .push_back(Hook {
3350 reached: reached_tx,
3351 release: release_rx,
3352 });
3353 (reached_rx, release_tx)
3354 }
3355
3356 pub(super) fn take(database_path: Option<&Path>) -> Option<Hook> {
3357 let key = database_path.map(Path::to_path_buf);
3358 registry()
3359 .lock()
3360 .unwrap_or_else(std::sync::PoisonError::into_inner)
3361 .get_mut(&key)
3362 .and_then(VecDeque::pop_front)
3363 }
3364}
3365
3366#[cfg(test)]
3367mod tests {
3368 use super::*;
3369
3370 fn store(floor_bytes: u64) -> (tempfile::TempDir, FsBlobStore) {
3371 let dir = tempfile::tempdir().unwrap();
3372 let root = dir.path().join("blobs");
3373 let store = FsBlobStore::new(root, floor_bytes)
3377 .unwrap()
3378 .with_orphan_sweep_grace(Duration::ZERO);
3379 (dir, store)
3380 }
3381
3382 #[cfg(unix)]
3383 #[tokio::test]
3384 async fn publication_barriers_cover_fresh_shards_and_dedup() {
3385 use std::os::unix::fs::MetadataExt;
3386
3387 let (_dir, store) = store(0);
3388 let bytes = b"publication barrier order".to_vec();
3389 let hook = sync_hook::install_publication(store.root(), None);
3390 let content_ref = store.put(bytes.clone()).await.unwrap();
3391 assert_eq!(
3392 *hook.completed.lock().unwrap(),
3393 [
3394 "put_fsync",
3395 "put_persist",
3396 "put_sync_shard",
3397 "put_sync_parent",
3398 "put_sync_root"
3399 ]
3400 );
3401 let path = shard_path(store.root(), &content_ref);
3402 let directory_identity = |operation, path: &Path| {
3403 let metadata = fs::metadata(path).unwrap();
3404 (operation, metadata.dev(), metadata.ino())
3405 };
3406 let expected_directories = [
3407 directory_identity("put_sync_shard", path.parent().unwrap()),
3408 directory_identity("put_sync_parent", path.parent().unwrap().parent().unwrap()),
3409 directory_identity("put_sync_root", store.root()),
3410 ];
3411 assert_eq!(*hook.directories.lock().unwrap(), expected_directories);
3412 let inode = fs::metadata(&path).unwrap().ino();
3413 let reopened = FsBlobStore::open_existing(store.root().to_path_buf(), 0).unwrap();
3414 let hook = sync_hook::install_publication(reopened.root(), None);
3415 assert_eq!(reopened.put(bytes.clone()).await.unwrap(), content_ref);
3416 assert_eq!(
3417 *hook.completed.lock().unwrap(),
3418 [
3419 "put_fsync",
3420 "put_sync_shard",
3421 "put_sync_parent",
3422 "put_sync_root"
3423 ]
3424 );
3425 assert_eq!(*hook.directories.lock().unwrap(), expected_directories);
3426 assert_eq!(fs::metadata(&path).unwrap().ino(), inode);
3427 assert_eq!(
3428 reopened
3429 .get_bounded_verified(&content_ref, bytes.len() as u64)
3430 .await
3431 .unwrap(),
3432 bytes
3433 );
3434 }
3435
3436 #[cfg(unix)]
3437 #[test]
3438 fn publication_barriers_initialize_root_but_not_read_only_open() {
3439 use std::os::unix::fs::MetadataExt;
3440
3441 let dir = tempfile::tempdir().unwrap();
3442 let root = dir.path().join("blobs");
3443 fs::create_dir(&root).unwrap();
3444 let hook = sync_hook::install_publication(&root, None);
3445 fs::remove_dir(&root).unwrap();
3446 let store = FsBlobStore::new(root.clone(), 0).unwrap();
3447 assert_eq!(
3448 *hook.completed.lock().unwrap(),
3449 ["init_sync_root", "init_sync_parent"]
3450 );
3451 let root_metadata = fs::metadata(&root).unwrap();
3452 let parent_metadata = fs::metadata(dir.path()).unwrap();
3453 assert_eq!(
3454 *hook.directories.lock().unwrap(),
3455 [
3456 ("init_sync_root", root_metadata.dev(), root_metadata.ino()),
3457 (
3458 "init_sync_parent",
3459 parent_metadata.dev(),
3460 parent_metadata.ino()
3461 ),
3462 ]
3463 );
3464
3465 let hook = sync_hook::install_publication(store.root(), Some("init_sync_root"));
3466 let reopened = FsBlobStore::open_existing(root.clone(), 0).unwrap();
3467 assert!(hook.completed.lock().unwrap().is_empty());
3468 assert!(FsBlobStore::new(root, 0).is_err());
3469 assert!(hook.completed.lock().unwrap().is_empty());
3470 assert_eq!(reopened.root(), store.root());
3471 }
3472
3473 #[cfg(unix)]
3474 #[test]
3475 fn publication_barriers_retry_each_initialization_failure() {
3476 for operation in ["init_sync_root", "init_sync_parent"] {
3477 let dir = tempfile::tempdir().unwrap();
3478 let root = dir.path().join("blobs");
3479 fs::create_dir(&root).unwrap();
3480 let hook = sync_hook::install_publication(&root, Some(operation));
3481 fs::remove_dir(&root).unwrap();
3482 assert!(FsBlobStore::new(root.clone(), 0).is_err(), "{operation}");
3483 assert!(root.is_dir());
3484 assert!(!hook.completed.lock().unwrap().contains(&operation));
3485
3486 let retry = sync_hook::install_publication(&root, None);
3487 FsBlobStore::new(root, 0).unwrap();
3488 assert_eq!(
3489 *retry.completed.lock().unwrap(),
3490 ["init_sync_root", "init_sync_parent"]
3491 );
3492 }
3493 }
3494
3495 #[test]
3496 fn publication_barriers_refuse_missing_root_parent_without_creating_it() {
3497 let dir = tempfile::tempdir().unwrap();
3498 let parent = dir.path().join("missing").join("nested");
3499 let root = parent.join("blobs");
3500 let error = match FsBlobStore::new(root.clone(), 0) {
3501 Ok(_) => panic!("missing parent must not be created"),
3502 Err(error) => error,
3503 };
3504 assert!(
3505 error.to_string().contains("blob root parent missing"),
3506 "{error}"
3507 );
3508 assert!(
3509 error.to_string().contains(&parent.display().to_string()),
3510 "{error}"
3511 );
3512 assert!(!dir.path().join("missing").exists());
3513
3514 fs::create_dir_all(&parent).unwrap();
3515 FsBlobStore::new(root, 0).unwrap();
3516 }
3517
3518 #[cfg(unix)]
3519 #[test]
3520 fn publication_barriers_fsync_propagates_kernel_errors() {
3521 let (socket, _peer) = std::os::unix::net::UnixStream::pair().unwrap();
3522 let handle = fs::File::from(std::os::fd::OwnedFd::from(socket));
3523 assert!(sync_directory(&handle).is_err());
3526 }
3527
3528 #[cfg(unix)]
3529 #[tokio::test]
3530 async fn publication_barriers_fail_before_rename_without_publishing() {
3531 for operation in ["put_fsync", "put_persist"] {
3532 let (_dir, store) = store(0);
3533 let bytes = b"unpublished failure".to_vec();
3534 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
3535 let hook = sync_hook::install_publication(store.root(), Some(operation));
3536 let error = store.put(bytes.clone()).await.unwrap_err();
3537 assert!(error.to_string().contains(operation), "{error}");
3538 let path = shard_path(store.root(), &content_ref);
3539 assert!(!path.exists());
3540 assert_eq!(fs::read_dir(path.parent().unwrap()).unwrap().count(), 0);
3541 assert!(!hook.completed.lock().unwrap().contains(&operation));
3542 assert_eq!(store.put(bytes).await.unwrap(), content_ref);
3543 }
3544 }
3545
3546 #[cfg(unix)]
3547 #[tokio::test]
3548 async fn publication_barriers_retry_after_each_directory_failure() {
3549 use std::os::unix::fs::MetadataExt;
3550
3551 for operation in ["put_sync_shard", "put_sync_parent", "put_sync_root"] {
3552 let (_dir, store) = store(0);
3553 let bytes = b"retry a visible but unacknowledged blob".to_vec();
3554 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
3555 sync_hook::install_publication(store.root(), Some(operation));
3556 let error = store.put(bytes.clone()).await.unwrap_err();
3557 assert!(error.to_string().contains(operation), "{error}");
3558 let path = shard_path(store.root(), &content_ref);
3559 let inode = fs::metadata(&path).unwrap().ino();
3560 assert_eq!(fs::read(&path).unwrap(), bytes);
3561
3562 let reopened = FsBlobStore::open_existing(store.root().to_path_buf(), 0).unwrap();
3563 let failed_retry = sync_hook::install_publication(reopened.root(), Some(operation));
3564 let error = reopened.put(bytes.clone()).await.unwrap_err();
3565 assert!(error.to_string().contains(operation), "{error}");
3566 assert!(!failed_retry
3567 .completed
3568 .lock()
3569 .unwrap()
3570 .contains(&"put_persist"));
3571 let repaired = sync_hook::install_publication(reopened.root(), None);
3572 assert_eq!(reopened.put(bytes.clone()).await.unwrap(), content_ref);
3573 assert_eq!(
3574 *repaired.completed.lock().unwrap(),
3575 [
3576 "put_fsync",
3577 "put_sync_shard",
3578 "put_sync_parent",
3579 "put_sync_root"
3580 ]
3581 );
3582 assert_eq!(fs::metadata(&path).unwrap().ino(), inode);
3583 assert_eq!(fs::read(&path).unwrap(), bytes);
3584 }
3585 }
3586
3587 #[cfg(unix)]
3588 #[tokio::test]
3589 async fn publication_barriers_fault_is_scoped_to_one_root_and_put() {
3590 let (_a, first) = store(0);
3591 let (_b, second) = store(0);
3592 sync_hook::install_publication(first.root(), Some("put_sync_shard"));
3593 second.put(b"second".to_vec()).await.unwrap();
3594 assert!(first.put(b"first".to_vec()).await.is_err());
3595 first.put(b"first".to_vec()).await.unwrap();
3596 }
3597
3598 #[cfg(unix)]
3599 #[tokio::test]
3600 async fn publication_barriers_keep_open_handles_when_root_path_changes() {
3601 let (dir, store) = store(0);
3602 let root = store.root().to_path_buf();
3603 let moved = dir.path().join("moved-root");
3604 let hook = sync_hook::install_publication(&root, None);
3605 let moved_for_hook = moved.clone();
3606 hook.on_step("put_sync_shard", move || {
3607 fs::rename(&root, moved_for_hook).unwrap();
3608 std::os::unix::fs::symlink(&root, &root).unwrap();
3611 });
3612 let bytes = b"pinned publication".to_vec();
3613 let content_ref = store.put(bytes.clone()).await.unwrap();
3614 assert_eq!(fs::read(shard_path(&moved, &content_ref)).unwrap(), bytes);
3615 assert_eq!(
3616 *hook.completed.lock().unwrap(),
3617 [
3618 "put_fsync",
3619 "put_persist",
3620 "put_sync_shard",
3621 "put_sync_parent",
3622 "put_sync_root"
3623 ]
3624 );
3625 assert!(store.put(bytes).await.is_err());
3626 }
3627
3628 #[cfg(unix)]
3629 #[tokio::test]
3630 async fn publication_barriers_reopen_in_a_fresh_process() {
3631 let (_dir, store) = store(0);
3632 let bytes = b"fresh process publication".to_vec();
3633 let content_ref = store.put(bytes).await.unwrap();
3634 let result = std::process::Command::new(std::env::current_exe().unwrap())
3635 .args([
3636 "--exact",
3637 "stores::blob::tests::publication_barriers_process_reader",
3638 "--ignored",
3639 "--nocapture",
3640 ])
3641 .env("KHIVE_TEST_BLOB_PUBLICATION_ROOT", store.root())
3642 .env("KHIVE_TEST_BLOB_PUBLICATION_REF", content_ref.as_str())
3643 .output()
3644 .unwrap();
3645 assert!(result.status.success(), "{result:?}");
3646 assert!(String::from_utf8_lossy(&result.stdout).contains("verified published object"));
3647 }
3648
3649 #[cfg(unix)]
3650 #[test]
3651 #[ignore = "subprocess helper for publication_barriers_reopen_in_a_fresh_process"]
3652 fn publication_barriers_process_reader() {
3653 let root = PathBuf::from(std::env::var_os("KHIVE_TEST_BLOB_PUBLICATION_ROOT").unwrap());
3654 let content_ref =
3655 ContentRef::from_hex(std::env::var("KHIVE_TEST_BLOB_PUBLICATION_REF").unwrap())
3656 .unwrap();
3657 let store = FsBlobStore::open_existing(root, 0).unwrap();
3658 let bytes = tokio::runtime::Runtime::new()
3659 .unwrap()
3660 .block_on(store.get_bounded_verified(&content_ref, 1024))
3661 .unwrap();
3662 assert_eq!(bytes, b"fresh process publication");
3663 println!("verified published object");
3664 }
3665
3666 fn prepare_v20_gc_fixture(conn: &mut rusqlite::Connection) {
3669 conn.execute_batch(include_str!("../../sql/schema-migrations-table.sql"))
3670 .expect("create migration ledger");
3671 for migration in crate::MIGRATIONS
3672 .iter()
3673 .filter(|migration| migration.version <= 20)
3674 {
3675 let tx = conn.transaction().expect("begin historical migration");
3676 tx.execute_batch(migration.up)
3677 .expect("apply historical migration body");
3678 tx.execute(
3679 "INSERT INTO _schema_migrations (version, name, applied_at) \
3680 VALUES (?1, ?2, 0)",
3681 rusqlite::params![migration.version, migration.name],
3682 )
3683 .expect("record historical migration");
3684 tx.commit().expect("commit historical migration");
3685 }
3686 }
3687
3688 fn prepare_completed_v21_gc_fixture(conn: &mut rusqlite::Connection) {
3697 prepare_v20_gc_fixture(conn);
3698 crate::migrations::stage_attachment_cutover(conn).expect("stage canonical completed V21");
3699 crate::migrations::finalize_attachment_cutover(conn)
3700 .expect("finalize canonical completed V21");
3701 let version = crate::migrations::read_schema_version(conn)
3702 .expect("read canonical completed V21 ledger");
3703 assert_eq!(
3704 version,
3705 crate::migrations::ATTACHMENT_CUTOVER_VERSION,
3706 "GC-gate fixtures need a completed-V21 ledger; a later migration chain \
3707 must provide a pinned through-V21 fixture builder for these tests"
3708 );
3709 }
3710
3711 #[tokio::test]
3712 async fn completed_v21_gc_gate_requires_new_indexes_and_absent_legacy_column() {
3713 let dir = tempfile::tempdir().unwrap();
3714 let db_path = dir.path().join("khive.db");
3715 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
3716 {
3717 let mut writer = backend.pool().writer().unwrap();
3718 prepare_completed_v21_gc_fixture(writer.conn_mut());
3719 }
3720 assert!(blob_gc_fencing_complete(backend.sql().as_ref())
3721 .await
3722 .unwrap());
3723
3724 {
3725 let writer = backend.pool().writer().unwrap();
3726 writer
3727 .conn()
3728 .execute_batch("DROP INDEX idx_attachments_content_ref")
3729 .unwrap();
3730 }
3731 assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
3732 .await
3733 .unwrap());
3734
3735 {
3736 let writer = backend.pool().writer().unwrap();
3737 writer
3738 .conn()
3739 .execute_batch(
3740 "CREATE INDEX idx_attachments_content_ref \
3741 ON attachments(content_ref); \
3742 ALTER TABLE entities ADD COLUMN content_ref TEXT",
3743 )
3744 .unwrap();
3745 }
3746 assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
3747 .await
3748 .unwrap());
3749 }
3750
3751 #[tokio::test]
3761 async fn completed_v21_gate_acceptance_matrix_rejects_each_removed_fact_independently() {
3762 let cases: &[(&str, &str)] = &[
3763 (
3764 "blob_gc_claims_content_ref_index_dropped",
3765 "DROP INDEX idx_blob_gc_claims_content_ref",
3766 ),
3767 (
3768 "v21_ledger_row_deleted",
3769 "DELETE FROM _schema_migrations WHERE version = 21",
3770 ),
3771 (
3772 "v21_ledger_row_renamed",
3773 "UPDATE _schema_migrations SET name = 'not_attachments_first_class' \
3774 WHERE version = 21",
3775 ),
3776 (
3786 "below_v21_ledger_rows_deleted",
3787 "DELETE FROM _schema_migrations WHERE version < 21",
3788 ),
3789 (
3790 "marker_row_deleted",
3791 "DELETE FROM attachment_cutover_state WHERE singleton = 1",
3792 ),
3793 (
3794 "marker_completed_at_null_while_state_complete",
3801 "DROP TABLE attachment_cutover_state; \
3802 CREATE TABLE attachment_cutover_state ( \
3803 singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
3804 state TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
3805 started_at INTEGER NOT NULL, \
3806 completed_at INTEGER \
3807 ) STRICT; \
3808 INSERT INTO attachment_cutover_state \
3809 (singleton, state, started_at, completed_at) \
3810 VALUES (1, 'complete', 21, NULL)",
3811 ),
3812 (
3813 "insert_fence_dropped_alone",
3814 "DROP TRIGGER attachments_reject_claimed_blob_insert",
3815 ),
3816 ];
3817
3818 for (case, mutation_sql) in cases {
3819 let dir = tempfile::tempdir().unwrap();
3820 let db_path = dir.path().join("khive.db");
3821 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
3822 {
3823 let mut writer = backend.pool().writer().unwrap();
3824 prepare_completed_v21_gc_fixture(writer.conn_mut());
3825 writer
3826 .conn_mut()
3827 .execute_batch(mutation_sql)
3828 .unwrap_or_else(|e| panic!("case {case}: failed to apply mutation: {e}"));
3829 }
3830
3831 assert!(
3832 !blob_gc_fencing_complete(backend.sql().as_ref())
3833 .await
3834 .unwrap(),
3835 "case {case}: gate must reject with this fact removed"
3836 );
3837
3838 let store = Arc::new(
3839 FsBlobStore::new(dir.path().join("blobs"), 0)
3840 .unwrap()
3841 .with_orphan_sweep_grace(Duration::ZERO),
3842 );
3843 let orphan = store
3844 .put(format!("gate matrix orphan for {case}").into_bytes())
3845 .await
3846 .unwrap();
3847 let abandoned_ref = "c".repeat(64);
3848 {
3849 let writer = backend.pool().writer().unwrap();
3850 writer
3851 .conn()
3852 .execute(
3853 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
3854 VALUES ('gate-matrix-abandoned', ?1, 1)",
3855 [abandoned_ref.as_str()],
3856 )
3857 .unwrap();
3858 }
3859
3860 let _root_guard = store.write_lock.clone().lock_owned().await;
3861 for dry_run in [true, false] {
3862 let outcome = tokio::time::timeout(
3863 Duration::from_secs(1),
3864 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
3865 )
3866 .await
3867 .unwrap_or_else(|_| panic!("case {case}: refusal must precede the root wait"));
3868 assert!(
3869 matches!(outcome, Err(StorageError::Unsupported { .. })),
3870 "case {case} dry_run={dry_run}: expected Unsupported, got {outcome:?}"
3871 );
3872 }
3873
3874 assert!(
3875 store.exists(&orphan).await.unwrap(),
3876 "case {case}: a refused sweep must not delete anything"
3877 );
3878 let reader = backend.pool().reader().unwrap();
3879 let remaining: i64 = reader
3880 .conn()
3881 .query_row(
3882 "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = 'gate-matrix-abandoned'",
3883 [],
3884 |row| row.get(0),
3885 )
3886 .unwrap();
3887 assert_eq!(remaining, 1, "case {case}: refusal must not recover claims");
3888
3889 let probe_claims: i64 = reader
3894 .conn()
3895 .query_row(
3896 "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key GLOB '__fence_probe-*'",
3897 [],
3898 |row| row.get(0),
3899 )
3900 .unwrap();
3901 assert_eq!(
3902 probe_claims, 0,
3903 "case {case}: refusal must leave no fence-probe claim residue"
3904 );
3905 let probe_attachments: i64 = reader
3906 .conn()
3907 .query_row(
3908 "SELECT COUNT(*) FROM attachments \
3909 WHERE record_uuid GLOB '__blob-gc-fence-probe-*'",
3910 [],
3911 |row| row.get(0),
3912 )
3913 .unwrap();
3914 assert_eq!(
3915 probe_attachments, 0,
3916 "case {case}: refusal must leave no fence-probe attachment residue"
3917 );
3918 }
3919 }
3920
3921 #[tokio::test]
3937 async fn gate_rejects_ledger_ahead_of_binary_latest() {
3938 let dir = tempfile::tempdir().unwrap();
3939 let db_path = dir.path().join("khive.db");
3940 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
3941 {
3942 let mut writer = backend.pool().writer().unwrap();
3943 prepare_completed_v21_gc_fixture(writer.conn_mut());
3944 }
3946 assert!(
3947 blob_gc_fencing_complete(backend.sql().as_ref())
3948 .await
3949 .unwrap(),
3950 "control: completed fixture at the binary's latest version must pass"
3951 );
3952
3953 {
3954 let writer = backend.pool().writer().unwrap();
3955 writer
3956 .conn()
3957 .execute(
3958 "INSERT INTO _schema_migrations (version, name, applied_at) \
3959 VALUES (?1, 'post_cutover_feature', unixepoch())",
3960 [i64::from(crate::migrations::latest_schema_version()) + 1],
3961 )
3962 .unwrap();
3963 }
3964 assert!(
3965 !blob_gc_fencing_complete(backend.sql().as_ref())
3966 .await
3967 .unwrap(),
3968 "a ledger ahead of the binary's latest schema version must fail closed"
3969 );
3970 }
3971
3972 #[tokio::test]
3980 async fn gate_rejects_incomplete_ledger_behind_v21() {
3981 let dir = tempfile::tempdir().unwrap();
3982 let db_path = dir.path().join("khive.db");
3983 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
3984 {
3985 let mut writer = backend.pool().writer().unwrap();
3986 prepare_completed_v21_gc_fixture(writer.conn_mut());
3987 }
3988 assert!(
3989 blob_gc_fencing_complete(backend.sql().as_ref())
3990 .await
3991 .unwrap(),
3992 "control: the untouched completed fixture must pass"
3993 );
3994
3995 {
3996 let writer = backend.pool().writer().unwrap();
3997 writer
3998 .conn()
3999 .execute_batch("DELETE FROM _schema_migrations WHERE version < 21")
4000 .unwrap();
4001 let (v21_named, max_version): (i64, i64) = writer
4003 .conn()
4004 .query_row(
4005 "SELECT (SELECT COUNT(*) FROM _schema_migrations \
4006 WHERE version = 21 AND name = 'attachments_first_class'), \
4007 (SELECT MAX(version) FROM _schema_migrations)",
4008 [],
4009 |row| Ok((row.get(0)?, row.get(1)?)),
4010 )
4011 .unwrap();
4012 assert_eq!(
4013 (v21_named, max_version),
4014 (1, 21),
4015 "fixture must keep the named V21 row and MAX(version) = 21 so the \
4016 rejection can only come from the contiguity clause"
4017 );
4018 }
4019 assert!(
4020 !blob_gc_fencing_complete(backend.sql().as_ref())
4021 .await
4022 .unwrap(),
4023 "an incomplete ledger behind V21 must fail closed despite a valid terminal row"
4024 );
4025 }
4026
4027 #[tokio::test]
4034 async fn completed_v21_gate_fails_closed_when_marker_read_errors() {
4035 let dir = tempfile::tempdir().unwrap();
4036 let db_path = dir.path().join("khive.db");
4037 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
4038 {
4039 let mut writer = backend.pool().writer().unwrap();
4040 prepare_completed_v21_gc_fixture(writer.conn_mut());
4041 writer
4046 .conn_mut()
4047 .execute_batch(
4048 "DROP TABLE attachment_cutover_state; \
4049 CREATE TABLE attachment_cutover_state ( \
4050 singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
4051 state TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
4052 started_at INTEGER NOT NULL \
4053 ) STRICT; \
4054 INSERT INTO attachment_cutover_state (singleton, state, started_at) \
4055 VALUES (1, 'complete', 21)",
4056 )
4057 .unwrap();
4058 }
4059
4060 let gate_error = blob_gc_fencing_complete(backend.sql().as_ref())
4061 .await
4062 .expect_err("a marker read error must propagate, not silently resolve to false");
4063 assert!(
4064 !matches!(gate_error, StorageError::Unsupported { .. }),
4065 "a read error is a distinct failure from the typed epoch refusal: {gate_error:?}"
4066 );
4067
4068 let store = Arc::new(
4069 FsBlobStore::new(dir.path().join("blobs"), 0)
4070 .unwrap()
4071 .with_orphan_sweep_grace(Duration::ZERO),
4072 );
4073 let orphan = store
4074 .put(b"marker read error orphan".to_vec())
4075 .await
4076 .unwrap();
4077 let _root_guard = store.write_lock.clone().lock_owned().await;
4078 for dry_run in [true, false] {
4079 let outcome = tokio::time::timeout(
4080 Duration::from_secs(1),
4081 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
4082 )
4083 .await
4084 .expect("a marker read error must fail before waiting on the root lock");
4085 assert!(
4086 outcome.is_err(),
4087 "dry_run={dry_run}: expected the sweep to fail closed, got {outcome:?}"
4088 );
4089 }
4090 assert!(store.exists(&orphan).await.unwrap());
4091 }
4092
4093 #[test]
4094 fn database_sweep_owner_is_keyed_by_database_not_blob_root() {
4095 let dir = tempfile::tempdir().unwrap();
4096 let database = dir.path().join("khive.db");
4097 let same_database_a = sweep_lock_for_database(Some(&database));
4098 let same_database_b = sweep_lock_for_database(Some(&database));
4099 let other_database = sweep_lock_for_database(Some(&dir.path().join("other.db")));
4100 let mut expected_lock_path = database.as_os_str().to_os_string();
4101 expected_lock_path.push(DATABASE_GC_LOCK_SUFFIX);
4102
4103 assert!(Arc::ptr_eq(&same_database_a, &same_database_b));
4104 assert!(!Arc::ptr_eq(&same_database_a, &other_database));
4105 assert_eq!(
4106 database_gc_lock_path(&database),
4107 PathBuf::from(expected_lock_path)
4108 );
4109 }
4110
4111 #[tokio::test]
4112 async fn database_gc_owner_holds_process_and_advisory_fences_until_drop() {
4113 let dir = tempfile::tempdir().unwrap();
4114 let database = dir.path().join("owner.db");
4115 let backend = crate::StorageBackend::sqlite_for_test(&database).unwrap();
4116 let owner = acquire_database_gc_owner(backend.sql().as_ref())
4117 .await
4118 .unwrap();
4119 let canonical_database = owner
4120 .database_path()
4121 .expect("file-backed owner path")
4122 .to_path_buf();
4123
4124 assert!(
4125 sweep_lock_for_database(Some(&canonical_database))
4126 .try_acquire()
4127 .is_none(),
4128 "boot and sweep must share one process-local database owner"
4129 );
4130 let external = fs::OpenOptions::new()
4131 .read(true)
4132 .write(true)
4133 .open(database_gc_lock_path(&canonical_database))
4134 .unwrap();
4135 assert!(
4136 matches!(
4137 fs4::FileExt::try_lock(&external),
4138 Err(fs4::TryLockError::WouldBlock)
4139 ),
4140 "the reusable owner must also retain the cross-process advisory fence"
4141 );
4142
4143 drop(owner);
4144 fs4::FileExt::try_lock(&external).expect("owner drop releases advisory fence");
4145 }
4146
4147 #[cfg(unix)]
4148 #[test]
4149 fn database_gc_lock_path_preserves_non_utf8_identity() {
4150 use std::os::unix::ffi::{OsStrExt, OsStringExt};
4151
4152 let database = PathBuf::from(std::ffi::OsString::from_vec(
4153 b"khive-non-utf8-\xff.db".to_vec(),
4154 ));
4155 let lock_path = database_gc_lock_path(&database);
4156 let mut expected = database.as_os_str().as_bytes().to_vec();
4157 expected.extend_from_slice(DATABASE_GC_LOCK_SUFFIX.as_bytes());
4158 assert_eq!(lock_path.as_os_str().as_bytes(), expected);
4159 }
4160
4161 #[cfg(unix)]
4170 #[tokio::test]
4171 async fn unlink_blob_shard_refuses_symlinked_shard_dir_and_still_sweeps_real_shard() {
4172 let (dir, store) = store(0);
4173 let root = dir.path().join("blobs");
4174
4175 let real = store.put(b"real blob content".to_vec()).await.unwrap();
4178
4179 let outside = tempfile::tempdir().unwrap();
4184 let victim = outside.path().join("victim.txt");
4185 fs::write(&victim, b"do not delete me").unwrap();
4186
4187 let real_prefix = &real.as_str()[0..2];
4188 let attack_prefix = if real_prefix == "aa" { "bb" } else { "aa" };
4189 let fake_ref = ContentRef::from_hex(format!("{attack_prefix}{}", "0".repeat(62))).unwrap();
4190 let fake_hex = fake_ref.as_str().to_string();
4191
4192 let shard1 = root.join(attack_prefix);
4193 std::os::unix::fs::symlink(outside.path(), &shard1).unwrap();
4194 fs::create_dir_all(outside.path().join(&fake_hex[2..4])).unwrap();
4200 fs::write(
4201 outside.path().join(&fake_hex[2..4]).join(&fake_hex),
4202 b"decoy",
4203 )
4204 .unwrap();
4205
4206 let root_handle = open_blob_root_handle(&root).unwrap();
4207 let error = unlink_blob_shard_file_no_follow(&root, &root_handle, &fake_ref).unwrap_err();
4208 assert!(
4214 matches!(
4215 error.raw_os_error(),
4216 Some(libc::ELOOP) | Some(libc::ENOTDIR)
4217 ),
4218 "opening a symlinked shard directory must be refused, not followed; got: {error}"
4219 );
4220 assert!(
4221 victim.exists(),
4222 "the file outside the blob root must never be touched by a refused shard-dir open"
4223 );
4224
4225 unlink_blob_shard_file_no_follow(&root, &root_handle, &real).unwrap();
4231 assert!(!store.exists(&real).await.unwrap());
4232 }
4233
4234 async fn recv_blocking(rx: std::sync::mpsc::Receiver<()>) -> bool {
4240 tokio::task::spawn_blocking(move || rx.recv().is_ok())
4241 .await
4242 .expect("recv_blocking thread panicked")
4243 }
4244
4245 #[tokio::test]
4246 async fn put_bounded_get_roundtrip() {
4247 let (_dir, store) = store(0);
4248 let bytes = b"hello blob store".to_vec();
4249 let content_ref = store.put(bytes.clone()).await.unwrap();
4250 let fetched = store
4251 .get_bounded_verified(&content_ref, bytes.len() as u64)
4252 .await
4253 .unwrap();
4254 assert_eq!(fetched, bytes);
4255 }
4256
4257 #[tokio::test]
4258 async fn bounded_verified_get_accepts_exact_and_portable_maximum_limits() {
4259 let (_dir, store) = store(0);
4260 let bytes = b"bounded fs blob".to_vec();
4261 let content_ref = store.put(bytes.clone()).await.unwrap();
4262
4263 assert_eq!(
4264 store
4265 .get_bounded_verified(&content_ref, bytes.len() as u64)
4266 .await
4267 .unwrap(),
4268 bytes
4269 );
4270 assert_eq!(
4271 store
4272 .get_bounded_verified(&content_ref, MAX_BLOB_WHOLE_BYTES)
4273 .await
4274 .unwrap(),
4275 bytes
4276 );
4277 }
4278
4279 #[cfg(unix)]
4280 #[tokio::test]
4281 async fn bounded_verified_get_resolves_a_configured_symlink_root_once() {
4282 use std::os::unix::fs::symlink;
4283
4284 let dir = tempfile::tempdir().unwrap();
4285 let target = dir.path().join("blob-target");
4286 fs::create_dir(&target).unwrap();
4287 let configured = dir.path().join("blob-configured");
4288 symlink(&target, &configured).unwrap();
4289
4290 let store = FsBlobStore::new(configured, 0).unwrap();
4291 assert_eq!(store.root(), target.canonicalize().unwrap());
4292 let bytes = b"symlink-configured root".to_vec();
4293 let content_ref = store.put(bytes.clone()).await.unwrap();
4294 assert_eq!(
4295 store
4296 .get_bounded_verified(&content_ref, bytes.len() as u64)
4297 .await
4298 .unwrap(),
4299 bytes
4300 );
4301 }
4302
4303 #[cfg(unix)]
4304 #[tokio::test]
4305 async fn fs_blob_store_refuses_root_replacement_before_put() {
4306 use std::os::unix::fs::symlink;
4307
4308 let dir = tempfile::tempdir().unwrap();
4309 let root = dir.path().join("blobs");
4310 let store = FsBlobStore::new(root.clone(), 0).unwrap();
4311 let original_root = dir.path().join("blobs-original");
4312 fs::rename(&root, &original_root).unwrap();
4313
4314 let redirected_root = tempfile::tempdir().unwrap();
4315 symlink(redirected_root.path(), &root).unwrap();
4316
4317 let bytes = b"must not land in a replaced root".to_vec();
4318 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
4319 store
4320 .put(bytes)
4321 .await
4322 .expect_err("a root replaced after construction must be refused");
4323
4324 assert!(
4325 !shard_path(redirected_root.path(), &content_ref).exists(),
4326 "put must not publish into the tree selected by the replacement symlink"
4327 );
4328 assert!(
4329 !shard_path(&original_root, &content_ref).exists(),
4330 "a refused put must not mutate the initialization-time root either"
4331 );
4332 }
4333
4334 #[cfg(unix)]
4335 #[tokio::test]
4336 async fn fs_blob_store_refuses_ancestor_replacement_for_existing_blob_operations() {
4337 use std::os::unix::fs::symlink;
4338
4339 let dir = tempfile::tempdir().unwrap();
4340 let ancestor = dir.path().join("store-parent");
4341 let root = ancestor.join("blobs");
4342 fs::create_dir(&ancestor).unwrap();
4343 let store = FsBlobStore::new(root.clone(), 0).unwrap();
4344 let bytes = b"same bytes in both trees".to_vec();
4345 let content_ref = store.put(bytes.clone()).await.unwrap();
4346
4347 let original_ancestor = dir.path().join("store-parent-original");
4348 fs::rename(&ancestor, &original_ancestor).unwrap();
4349 let original_blob = shard_path(&original_ancestor.join("blobs"), &content_ref);
4350
4351 let redirected_ancestor = tempfile::tempdir().unwrap();
4352 let redirected_blob = shard_path(&redirected_ancestor.path().join("blobs"), &content_ref);
4353 fs::create_dir_all(redirected_blob.parent().unwrap()).unwrap();
4354 fs::write(&redirected_blob, &bytes).unwrap();
4355 symlink(redirected_ancestor.path(), &ancestor).unwrap();
4356
4357 store
4358 .get_bounded_verified(&content_ref, bytes.len() as u64)
4359 .await
4360 .expect_err("a read through a replaced root ancestor must be refused");
4361 store
4362 .exists(&content_ref)
4363 .await
4364 .expect_err("exists through a replaced root ancestor must be refused");
4365 store
4366 .size(&content_ref)
4367 .await
4368 .expect_err("size through a replaced root ancestor must be refused");
4369 store
4370 .delete(&content_ref)
4371 .await
4372 .expect_err("delete through a replaced root ancestor must be refused");
4373
4374 assert!(
4375 original_blob.exists(),
4376 "refusal must preserve the initialization-time blob"
4377 );
4378 assert!(
4379 redirected_blob.exists(),
4380 "refusal must not read as authority or delete the redirected blob"
4381 );
4382 }
4383
4384 #[tokio::test]
4385 async fn bounded_verified_get_rejects_same_size_digest_corruption() {
4386 let (_dir, store) = store(0);
4387 let expected_bytes = b"expected".to_vec();
4388 let actual_bytes = b"mutated!".to_vec();
4389 let expected = store.put(expected_bytes).await.unwrap();
4390 let actual = ContentRef::from_digest_bytes(blake3::hash(&actual_bytes).as_bytes());
4391 fs::write(shard_path(store.root(), &expected), actual_bytes).unwrap();
4392
4393 let err = store.get_bounded_verified(&expected, 8).await.unwrap_err();
4394 assert!(matches!(
4395 err,
4396 StorageError::BlobDigestMismatch {
4397 expected: ref got_expected,
4398 actual: ref got_actual,
4399 } if got_expected == &expected && got_actual == &actual
4400 ));
4401 }
4402
4403 #[tokio::test]
4404 async fn bounded_verified_get_stops_at_max_plus_one_after_file_growth() {
4405 let (_dir, store) = store(0);
4406 let store = Arc::new(store);
4407 let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
4408 let path = shard_path(store.root(), &content_ref);
4409 let (reached, release) = bounded_read_sync_hook::install(store.root());
4410
4411 let read_store = Arc::clone(&store);
4412 let read_ref = content_ref.clone();
4413 let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
4414 assert!(
4415 recv_blocking(reached).await,
4416 "read must reach the metadata seam"
4417 );
4418 let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
4419 writer.write_all(b"efgh-poison-tail").unwrap();
4420 writer.flush().unwrap();
4421 release.send(()).unwrap();
4422
4423 let err = read.await.unwrap().unwrap_err();
4424 assert!(matches!(
4425 err,
4426 StorageError::BlobTooLarge {
4427 content_ref: ref got,
4428 max_bytes: 4,
4429 observed_at_least: 5,
4430 } if got == &content_ref
4431 ));
4432 }
4433
4434 #[tokio::test]
4435 async fn bounded_verified_get_reports_growth_within_limit_as_size_mismatch() {
4436 let (_dir, store) = store(0);
4437 let store = Arc::new(store);
4438 let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
4439 let path = shard_path(store.root(), &content_ref);
4440 let (reached, release) = bounded_read_sync_hook::install(store.root());
4441
4442 let read_store = Arc::clone(&store);
4443 let read_ref = content_ref.clone();
4444 let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 8).await });
4445 assert!(
4446 recv_blocking(reached).await,
4447 "read must reach the metadata seam"
4448 );
4449 let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
4450 writer.write_all(b"ef").unwrap();
4451 writer.flush().unwrap();
4452 release.send(()).unwrap();
4453
4454 let err = read.await.unwrap().unwrap_err();
4455 assert!(matches!(
4456 err,
4457 StorageError::BlobSizeMismatch {
4458 content_ref: ref got,
4459 metadata_bytes: 4,
4460 actual_bytes: 6,
4461 } if got == &content_ref
4462 ));
4463 }
4464
4465 #[tokio::test]
4466 async fn bounded_verified_get_reports_truncation_as_size_mismatch() {
4467 let (_dir, store) = store(0);
4468 let store = Arc::new(store);
4469 let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
4470 let path = shard_path(store.root(), &content_ref);
4471 let (reached, release) = bounded_read_sync_hook::install(store.root());
4472
4473 let read_store = Arc::clone(&store);
4474 let read_ref = content_ref.clone();
4475 let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
4476 assert!(
4477 recv_blocking(reached).await,
4478 "read must reach the metadata seam"
4479 );
4480 let mut writer = fs::OpenOptions::new()
4481 .write(true)
4482 .truncate(true)
4483 .open(&path)
4484 .unwrap();
4485 writer.write_all(b"abc").unwrap();
4486 writer.flush().unwrap();
4487 release.send(()).unwrap();
4488
4489 let err = read.await.unwrap().unwrap_err();
4490 assert!(matches!(
4491 err,
4492 StorageError::BlobSizeMismatch {
4493 content_ref: ref got,
4494 metadata_bytes: 4,
4495 actual_bytes: 3,
4496 } if got == &content_ref
4497 ));
4498 }
4499
4500 #[cfg(unix)]
4501 #[tokio::test]
4502 async fn bounded_verified_get_keeps_the_opened_inode_when_the_path_is_replaced() {
4503 let (_dir, store) = store(0);
4504 let store = Arc::new(store);
4505 let original = b"original".to_vec();
4506 let replacement = b"replaced".to_vec();
4507 let content_ref = store.put(original.clone()).await.unwrap();
4508 let path = shard_path(store.root(), &content_ref);
4509 let moved_path = path.with_extension("opened-inode");
4510 let max_bytes = original.len() as u64;
4511 let (reached, release) = bounded_read_sync_hook::install(store.root());
4512
4513 let read_store = Arc::clone(&store);
4514 let read_ref = content_ref.clone();
4515 let read =
4516 tokio::spawn(
4517 async move { read_store.get_bounded_verified(&read_ref, max_bytes).await },
4518 );
4519 assert!(
4520 recv_blocking(reached).await,
4521 "read must reach the metadata seam"
4522 );
4523 fs::rename(&path, &moved_path).unwrap();
4524 fs::write(&path, replacement).unwrap();
4525 release.send(()).unwrap();
4526
4527 assert_eq!(read.await.unwrap().unwrap(), original);
4528 }
4529
4530 #[cfg(unix)]
4531 #[tokio::test]
4532 async fn bounded_verified_get_refuses_a_symlink_leaf() {
4533 use std::os::unix::fs::symlink;
4534
4535 let dir = tempfile::tempdir().unwrap();
4536 let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
4537 let outside_bytes = b"outside but digest matching".to_vec();
4538 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
4539 let outside = dir.path().join("outside");
4540 fs::write(&outside, &outside_bytes).unwrap();
4541 let leaf = shard_path(store.root(), &content_ref);
4542 fs::create_dir_all(leaf.parent().unwrap()).unwrap();
4543 symlink(&outside, &leaf).unwrap();
4544
4545 let err = store
4546 .get_bounded_verified(&content_ref, outside_bytes.len() as u64)
4547 .await
4548 .unwrap_err();
4549 assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
4550 }
4551
4552 #[cfg(unix)]
4553 #[tokio::test]
4554 async fn bounded_verified_get_refuses_a_symlinked_shard_component() {
4555 use std::os::unix::fs::symlink;
4556
4557 let dir = tempfile::tempdir().unwrap();
4558 let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
4559 let outside_bytes = b"outside through shard link".to_vec();
4560 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
4561 let hex = content_ref.as_str();
4562 let outside_shard1 = dir.path().join("outside-shard1");
4563 let outside_shard2 = outside_shard1.join(&hex[2..4]);
4564 fs::create_dir_all(&outside_shard2).unwrap();
4565 fs::write(outside_shard2.join(hex), &outside_bytes).unwrap();
4566 symlink(&outside_shard1, store.root().join(&hex[0..2])).unwrap();
4567
4568 let err = store
4569 .get_bounded_verified(&content_ref, outside_bytes.len() as u64)
4570 .await
4571 .unwrap_err();
4572 assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
4573 }
4574
4575 #[tokio::test]
4576 async fn put_content_ref_matches_blake3_digest() {
4577 let (_dir, store) = store(0);
4578 let bytes = b"digest check".to_vec();
4579 let content_ref = store.put(bytes.clone()).await.unwrap();
4580 let expected = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
4581 assert_eq!(content_ref, expected);
4582 }
4583
4584 #[tokio::test]
4585 async fn put_dedups_identical_content() {
4586 let (_dir, store) = store(0);
4587 let bytes = b"same bytes twice".to_vec();
4588 let first = store.put(bytes.clone()).await.unwrap();
4589 let second = store.put(bytes.clone()).await.unwrap();
4590 assert_eq!(first, second);
4591 assert_eq!(
4592 store
4593 .get_bounded_verified(&first, bytes.len() as u64)
4594 .await
4595 .unwrap(),
4596 bytes
4597 );
4598 }
4599
4600 #[tokio::test]
4601 async fn exists_reflects_put_and_delete() {
4602 let (_dir, store) = store(0);
4603 let bytes = b"exists check".to_vec();
4604 let content_ref = store.put(bytes).await.unwrap();
4605 assert!(store.exists(&content_ref).await.unwrap());
4606
4607 assert!(store.delete(&content_ref).await.unwrap());
4608 assert!(!store.exists(&content_ref).await.unwrap());
4609 }
4610
4611 #[tokio::test]
4612 async fn delete_missing_content_ref_returns_false() {
4613 let (_dir, store) = store(0);
4614 let missing = ContentRef::from_hex("f".repeat(64)).unwrap();
4615 assert!(!store.delete(&missing).await.unwrap());
4616 }
4617
4618 #[tokio::test]
4619 async fn size_reports_byte_length_for_a_present_object() {
4620 let (_dir, store) = store(0);
4621 let bytes = b"size check".to_vec();
4622 let content_ref = store.put(bytes.clone()).await.unwrap();
4623 assert_eq!(
4624 store.size(&content_ref).await.unwrap(),
4625 Some(bytes.len() as u64)
4626 );
4627 }
4628
4629 #[tokio::test]
4630 async fn size_returns_none_for_an_absent_object() {
4631 let (_dir, store) = store(0);
4632 let missing = ContentRef::from_hex("9".repeat(64)).unwrap();
4633 assert_eq!(store.size(&missing).await.unwrap(), None);
4634 }
4635
4636 #[tokio::test]
4637 async fn bounded_get_missing_content_ref_returns_not_found() {
4638 let (_dir, store) = store(0);
4639 let missing = ContentRef::from_hex("e".repeat(64)).unwrap();
4640 let err = store
4641 .get_bounded_verified(&missing, MAX_BLOB_WHOLE_BYTES)
4642 .await
4643 .unwrap_err();
4644 assert!(matches!(err, StorageError::NotFound { .. }));
4645 }
4646
4647 #[tokio::test]
4648 async fn put_refuses_below_free_space_floor() {
4649 let (_dir, store) = store(u64::MAX);
4652 let err = store.put(b"too big a floor".to_vec()).await.unwrap_err();
4653 match err {
4654 StorageError::CapacityFloor {
4655 floor_bytes,
4656 available_bytes,
4657 ..
4658 } => {
4659 assert_eq!(floor_bytes, u64::MAX);
4660 assert!(available_bytes < u64::MAX);
4661 }
4662 other => panic!("expected CapacityFloor, got {other:?}"),
4663 }
4664 }
4665
4666 #[tokio::test]
4667 async fn capacity_floor_error_names_the_floor_and_volume() {
4668 let (_dir, store) = store(u64::MAX);
4669 let err = store.put(b"x".to_vec()).await.unwrap_err();
4670 let msg = err.to_string();
4671 assert!(
4672 msg.contains(&u64::MAX.to_string()),
4673 "must name the floor: {msg}"
4674 );
4675 assert!(msg.contains("Blob"), "must name the capability: {msg}");
4676 }
4677
4678 #[test]
4679 fn crosses_floor_is_write_size_aware_at_the_exact_boundary() {
4680 assert!(crosses_floor(101, 2, 100));
4685 assert!(!crosses_floor(101, 1, 100));
4686 }
4687
4688 #[test]
4689 fn crosses_floor_accepts_a_write_that_lands_exactly_on_the_floor() {
4690 assert!(!crosses_floor(100, 0, 100));
4691 }
4692
4693 #[test]
4694 fn crosses_floor_rejects_a_write_that_lands_one_byte_under_the_floor() {
4695 assert!(crosses_floor(100, 1, 100));
4696 }
4697
4698 #[test]
4699 fn crosses_floor_saturates_instead_of_underflowing_when_write_exceeds_available() {
4700 assert!(crosses_floor(10, 100, 50));
4701 assert!(!crosses_floor(10, 100, 0));
4707 }
4708
4709 #[test]
4710 fn put_refuses_a_write_that_would_cross_the_floor_even_though_available_alone_clears_it() {
4711 let dir = tempfile::tempdir().unwrap();
4712 let root = dir.path().join("blobs");
4713 fs::create_dir_all(&root).unwrap();
4714
4715 let err = put_blocking_with_space_probe(&root, 100, vec![7u8; 2], |_| Ok(101)).unwrap_err();
4722 assert!(
4723 matches!(err, StorageError::CapacityFloor { .. }),
4724 "a write-size-aware floor check must reject a write that pushes the volume \
4725 below the floor even though available space alone still clears it: {err:?}"
4726 );
4727 }
4728
4729 #[test]
4730 fn a_later_put_checks_a_fresh_capacity_snapshot() {
4731 let dir = tempfile::tempdir().unwrap();
4732 let root = dir.path().join("blobs");
4733 fs::create_dir_all(&root).unwrap();
4734
4735 let first = put_blocking_with_space_probe(&root, 100, vec![1u8; 2], |_| Ok(102));
4743 let second = put_blocking_with_space_probe(&root, 100, vec![2u8; 2], |_| Ok(101));
4744
4745 assert!(
4746 first.is_ok(),
4747 "the first put may land on the floor: {first:?}"
4748 );
4749 assert!(
4750 matches!(second, Err(StorageError::CapacityFloor { .. })),
4751 "the later put must use its lower capacity snapshot: {second:?}"
4752 );
4753 }
4754
4755 #[tokio::test]
4756 async fn concurrent_puts_from_two_independently_constructed_stores_share_the_root_lock() {
4757 let dir = tempfile::tempdir().unwrap();
4810 let root = dir.path().join("blobs");
4811 fs::create_dir_all(&root).unwrap();
4812 let canonical_root = root.canonicalize().unwrap();
4813
4814 let store_a = std::sync::Arc::new(FsBlobStore::new(root.clone(), 0).unwrap());
4817 let store_b = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
4818
4819 let (a_reached, a_release, _a_done) = sync_hook::install(&canonical_root);
4820 let a = {
4821 let store_a = store_a.clone();
4822 tokio::spawn(async move { store_a.put(b"store_a payload".to_vec()).await })
4823 };
4824 assert!(
4825 recv_blocking(a_reached).await,
4826 "store_a's put must reach the sync_hook checkpoint"
4827 );
4828
4829 assert!(
4834 store_b.write_lock.try_lock().is_err(),
4835 "store_b's write_lock was NOT held while store_a's put held its guard -- the two \
4836 independently constructed stores do NOT share one lock"
4837 );
4838
4839 a_release.send(()).unwrap();
4843 let result_a = a.await.unwrap();
4844 assert!(result_a.is_ok(), "store_a's put must succeed: {result_a:?}");
4845
4846 let result_b = store_b.put(b"store_b payload".to_vec()).await;
4849 assert!(result_b.is_ok(), "store_b's put must succeed: {result_b:?}");
4850 }
4851
4852 #[tokio::test]
4853 async fn aborting_the_outer_put_future_does_not_release_the_guard_before_persist_completes() {
4854 let dir = tempfile::tempdir().unwrap();
4874 let root = dir.path().join("blobs");
4875 fs::create_dir_all(&root).unwrap();
4876 let canonical_root = root.canonicalize().unwrap();
4877
4878 let store = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
4879 let (reached, release, done) = sync_hook::install(&canonical_root);
4880 let handle = {
4881 let store = store.clone();
4882 tokio::spawn(async move { store.put(b"cancellation race payload".to_vec()).await })
4883 };
4884
4885 assert!(
4886 recv_blocking(reached).await,
4887 "put must reach the sync_hook checkpoint -- owned guard already moved into the \
4888 closure -- before this test can mean anything"
4889 );
4890
4891 handle.abort();
4892 let abort_result = handle.await;
4893 match &abort_result {
4894 Err(e) if e.is_cancelled() => {}
4895 other => panic!(
4896 "the outer task must actually have been cancelled for this test to be \
4897 meaningful: {other:?}"
4898 ),
4899 }
4900
4901 let shared_lock = write_lock_for_root(&canonical_root).unwrap();
4902 assert!(
4903 shared_lock.try_lock().is_err(),
4904 "the guard must still be held by the detached blocking write immediately after \
4905 the outer future was cancelled -- if this is free, the guard was released with \
4906 the aborted frame instead of moving into the spawn_blocking closure"
4907 );
4908
4909 release.send(()).unwrap();
4914 assert!(
4915 recv_blocking(done).await,
4916 "the detached write must signal completion once it actually persists"
4917 );
4918 assert!(
4919 shared_lock.try_lock().is_ok(),
4920 "the guard must be free once the detached write's completion was observed"
4921 );
4922 }
4923
4924 #[tokio::test]
4932 async fn orphan_sweep_is_disabled_in_both_modes_regardless_of_live_refs() {
4933 let (_dir, store) = store(0);
4934 let blob = store
4935 .put(b"never swept by this API".to_vec())
4936 .await
4937 .unwrap();
4938 let mut live_refs = std::collections::HashSet::new();
4939 live_refs.insert(blob.clone());
4940
4941 for dry_run in [true, false] {
4942 let error = store
4943 .orphan_sweep(&BlobOrphanSweepConfig {
4944 live_refs: live_refs.clone(),
4945 dry_run,
4946 })
4947 .await
4948 .expect_err("caller-snapshot orphan_sweep must be disabled");
4949 assert!(
4950 matches!(error, StorageError::Unsupported { .. }),
4951 "expected typed Unsupported, got {error:?}"
4952 );
4953 }
4954 assert!(store.exists(&blob).await.unwrap());
4955 }
4956
4957 #[tokio::test]
4962 async fn transactional_orphan_sweep_refuses_v20_before_root_or_claim_mutation() {
4963 let dir = tempfile::tempdir().unwrap();
4964 let db_path = dir.path().join("khive.db");
4965 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
4966 {
4967 let mut writer = backend.pool().writer().unwrap();
4968 prepare_v20_gc_fixture(writer.conn_mut());
4969 }
4970
4971 let root = dir.path().join("blobs");
4972 let store = Arc::new(
4973 FsBlobStore::new(root, 0)
4974 .unwrap()
4975 .with_orphan_sweep_grace(Duration::ZERO),
4976 );
4977 let bundle = store.put(b"legacy model bundle".to_vec()).await.unwrap();
4978 let network = store.put(b"legacy FANN network".to_vec()).await.unwrap();
4979 let orphan = store.put(b"ordinary old orphan".to_vec()).await.unwrap();
4980 let abandoned_ref = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
4981 {
4982 let writer = backend.pool().writer().unwrap();
4983 writer
4984 .conn()
4985 .execute(
4986 "INSERT INTO entities \
4987 (id, namespace, kind, entity_type, name, tags, created_at, updated_at, \
4988 content_ref) \
4989 VALUES ('legacy-model', 'local', 'artifact', 'moodboard_model', \
4990 'legacy model', '[]', 1, 1, ?1)",
4991 [bundle.as_str()],
4992 )
4993 .unwrap();
4994 writer
4995 .conn()
4996 .execute(
4997 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4998 VALUES ('abandoned-before-compat', ?1, 1)",
4999 [abandoned_ref],
5000 )
5001 .unwrap();
5002 }
5003
5004 let _root_guard = store.write_lock.clone().lock_owned().await;
5007 for dry_run in [true, false] {
5008 let outcome = tokio::time::timeout(
5009 Duration::from_secs(1),
5010 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
5011 )
5012 .await
5013 .expect("V20 refusal must happen before waiting for the held root lock");
5014 let error = outcome.expect_err("V20 transactional sweep must be disabled");
5015 match error {
5016 StorageError::Unsupported {
5017 capability: StorageCapability::Blob,
5018 operation,
5019 message,
5020 } => {
5021 assert_eq!(operation, "transactional_orphan_sweep");
5022 assert!(
5023 message.contains("complete V21 attachment cutover"),
5024 "unexpected compatibility diagnostic: {message}"
5025 );
5026 }
5027 other => panic!("expected typed Unsupported refusal, got {other:?}"),
5028 }
5029 }
5030
5031 let reader = backend.pool().reader().unwrap();
5032 let abandoned: i64 = reader
5033 .conn()
5034 .query_row(
5035 "SELECT COUNT(*) FROM blob_gc_claims \
5036 WHERE root_key = 'abandoned-before-compat' AND content_ref = ?1",
5037 [abandoned_ref],
5038 |row| row.get(0),
5039 )
5040 .unwrap();
5041 assert_eq!(abandoned, 1, "V20 refusal must not clean abandoned claims");
5042 drop(reader);
5043 assert!(store.exists(&bundle).await.unwrap());
5044 assert!(store.exists(&network).await.unwrap());
5045 assert!(store.exists(&orphan).await.unwrap());
5046 }
5047
5048 #[tokio::test]
5053 async fn transactional_orphan_sweep_refuses_incomplete_v21_marker_without_mutation() {
5054 let dir = tempfile::tempdir().unwrap();
5055 let db_path = dir.path().join("khive.db");
5056 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5057 let abandoned_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
5058 {
5059 let mut writer = backend.pool().writer().unwrap();
5060 prepare_completed_v21_gc_fixture(writer.conn_mut());
5061 writer
5062 .conn_mut()
5063 .execute(
5064 "UPDATE attachment_cutover_state \
5065 SET state = 'incomplete', completed_at = NULL \
5066 WHERE singleton = 1",
5067 [],
5068 )
5069 .unwrap();
5070 writer
5071 .conn_mut()
5072 .execute(
5073 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5074 VALUES ('abandoned-incomplete-v21', ?1, 1)",
5075 [abandoned_ref],
5076 )
5077 .unwrap();
5078 }
5079
5080 let store = Arc::new(
5081 FsBlobStore::new(dir.path().join("blobs"), 0)
5082 .unwrap()
5083 .with_orphan_sweep_grace(Duration::ZERO),
5084 );
5085 let orphan = store.put(b"incomplete V21 orphan".to_vec()).await.unwrap();
5086 let _root_guard = store.write_lock.clone().lock_owned().await;
5087
5088 for dry_run in [true, false] {
5089 let outcome = tokio::time::timeout(
5090 Duration::from_secs(1),
5091 store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
5092 )
5093 .await
5094 .expect("incomplete V21 must refuse before waiting for the root lock");
5095 assert!(
5096 matches!(outcome, Err(StorageError::Unsupported { .. })),
5097 "incomplete V21 must return typed Unsupported: {outcome:?}"
5098 );
5099 }
5100
5101 let remaining: i64 = backend
5102 .pool()
5103 .reader()
5104 .unwrap()
5105 .conn()
5106 .query_row(
5107 "SELECT COUNT(*) FROM blob_gc_claims \
5108 WHERE root_key = 'abandoned-incomplete-v21' AND content_ref = ?1",
5109 [abandoned_ref],
5110 |row| row.get(0),
5111 )
5112 .unwrap();
5113 assert_eq!(remaining, 1, "refusal must not recover abandoned claims");
5114 assert!(store.exists(&orphan).await.unwrap());
5115 }
5116
5117 #[tokio::test]
5124 async fn both_sweep_apis_refuse_v20_and_incomplete_v21_epochs_in_both_modes() {
5125 for incomplete_v21 in [false, true] {
5126 let dir = tempfile::tempdir().unwrap();
5127 let db_path = dir.path().join("khive.db");
5128 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5129 {
5130 let mut writer = backend.pool().writer().unwrap();
5131 if incomplete_v21 {
5132 prepare_completed_v21_gc_fixture(writer.conn_mut());
5133 writer
5134 .conn_mut()
5135 .execute(
5136 "UPDATE attachment_cutover_state \
5137 SET state = 'incomplete', completed_at = NULL \
5138 WHERE singleton = 1",
5139 [],
5140 )
5141 .unwrap();
5142 } else {
5143 prepare_v20_gc_fixture(writer.conn_mut());
5144 }
5145 }
5146
5147 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
5148 .unwrap()
5149 .with_orphan_sweep_grace(Duration::ZERO);
5150 let orphan = store
5151 .put(format!("both-apis orphan (incomplete_v21={incomplete_v21})").into_bytes())
5152 .await
5153 .unwrap();
5154
5155 let known_claim_ref = "f".repeat(64);
5159 {
5160 let writer = backend.pool().writer().unwrap();
5161 writer
5162 .conn()
5163 .execute(
5164 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5165 VALUES ('both-api-known-claim', ?1, 1)",
5166 [known_claim_ref.as_str()],
5167 )
5168 .unwrap();
5169 }
5170 let assert_known_claim_unchanged = |arm: &str| {
5171 let remaining: i64 = backend
5172 .pool()
5173 .reader()
5174 .unwrap()
5175 .conn()
5176 .query_row(
5177 "SELECT COUNT(*) FROM blob_gc_claims \
5178 WHERE root_key = 'both-api-known-claim' AND content_ref = ?1",
5179 [known_claim_ref.as_str()],
5180 |row| row.get(0),
5181 )
5182 .unwrap();
5183 assert_eq!(
5184 remaining, 1,
5185 "incomplete_v21={incomplete_v21} arm={arm}: refusal must not mutate \
5186 an existing claim"
5187 );
5188 };
5189
5190 for dry_run in [true, false] {
5191 let snapshot_error = store
5192 .orphan_sweep(&BlobOrphanSweepConfig {
5193 live_refs: std::collections::HashSet::new(),
5194 dry_run,
5195 })
5196 .await
5197 .expect_err("orphan_sweep must refuse regardless of epoch");
5198 assert!(
5199 matches!(snapshot_error, StorageError::Unsupported { .. }),
5200 "incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
5201 from orphan_sweep, got {snapshot_error:?}"
5202 );
5203 assert_known_claim_unchanged(&format!("orphan_sweep dry_run={dry_run}"));
5204
5205 let transactional_error = store
5206 .transactional_orphan_sweep(backend.sql().as_ref(), dry_run)
5207 .await
5208 .expect_err("transactional_orphan_sweep must refuse this epoch");
5209 assert!(
5210 matches!(transactional_error, StorageError::Unsupported { .. }),
5211 "incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
5212 from transactional_orphan_sweep, got {transactional_error:?}"
5213 );
5214 assert_known_claim_unchanged(&format!(
5215 "transactional_orphan_sweep dry_run={dry_run}"
5216 ));
5217 }
5218
5219 assert!(
5220 store.exists(&orphan).await.unwrap(),
5221 "incomplete_v21={incomplete_v21}: a refused sweep must not delete anything"
5222 );
5223 }
5224 }
5225
5226 #[tokio::test]
5246 async fn transactional_orphan_sweep_recheck_refuses_before_root_lock_when_epoch_regresses_after_db_ownership(
5247 ) {
5248 let dir = tempfile::tempdir().unwrap();
5249 let db_path = dir.path().join("khive.db");
5250 let backend = Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
5251 {
5252 let mut writer = backend.pool().writer().unwrap();
5253 prepare_completed_v21_gc_fixture(writer.conn_mut());
5254 }
5255 assert!(blob_gc_fencing_complete(backend.sql().as_ref())
5256 .await
5257 .unwrap());
5258
5259 let blob_root = dir.path().join("blobs");
5260 let store = Arc::new(
5261 FsBlobStore::new(blob_root.clone(), 0)
5262 .unwrap()
5263 .with_orphan_sweep_grace(Duration::ZERO),
5264 );
5265 let orphan = store
5266 .put(b"epoch regressed after database ownership".to_vec())
5267 .await
5268 .unwrap();
5269
5270 let canonical_root = blob_root.canonicalize().unwrap();
5273 let _root_write_guard = acquire_root_write_lock(&canonical_root).unwrap();
5274
5275 let canonical_db_path = backend.sql().database_path();
5279 let (reached, release) = db_ownership_sync_hook::install(canonical_db_path.as_deref());
5280 let sweep_store = store.clone();
5281 let sweep_backend = backend.clone();
5282 let handle = tokio::spawn(async move {
5283 sweep_store
5284 .transactional_orphan_sweep(sweep_backend.sql().as_ref(), false)
5285 .await
5286 });
5287
5288 let reached_signal = tokio::time::timeout(Duration::from_secs(1), recv_blocking(reached))
5294 .await
5295 .expect("the sweep must reach database ownership before this test's timeout");
5296 assert!(reached_signal, "hook sender was dropped before signaling");
5297 {
5298 let writer = backend.pool().writer().unwrap();
5299 writer
5300 .conn()
5301 .execute("DELETE FROM _schema_migrations WHERE version = 21", [])
5302 .unwrap();
5303 }
5304 release.send(()).unwrap();
5305
5306 let outcome = tokio::time::timeout(Duration::from_secs(1), handle)
5307 .await
5308 .expect(
5309 "the recheck must refuse before ever waiting on the externally held root lock -- \
5310 under the old (pre-fix) ordering this join times out instead, because the \
5311 sweep blocks acquiring the OS-level root lock held above",
5312 )
5313 .unwrap();
5314 assert!(
5315 matches!(outcome, Err(StorageError::Unsupported { .. })),
5316 "expected the regressed epoch to be caught immediately after database ownership: \
5317 {outcome:?}"
5318 );
5319 assert!(store.exists(&orphan).await.unwrap());
5320 }
5321
5322 #[tokio::test]
5327 async fn transactional_orphan_sweep_accepts_completed_v21_attachment_liveness() {
5328 let dir = tempfile::tempdir().unwrap();
5329 let db_path = dir.path().join("khive.db");
5330 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5331 {
5332 let mut writer = backend.pool().writer().unwrap();
5333 prepare_completed_v21_gc_fixture(writer.conn_mut());
5334 }
5335
5336 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
5337 .unwrap()
5338 .with_orphan_sweep_grace(Duration::ZERO);
5339 let bundle = store.put(b"V21 model bundle".to_vec()).await.unwrap();
5340 let network = store.put(b"V21 FANN network".to_vec()).await.unwrap();
5341 let orphan = store.put(b"V21 true orphan".to_vec()).await.unwrap();
5342 {
5343 let writer = backend.pool().writer().unwrap();
5344 writer
5345 .conn()
5346 .execute(
5347 "INSERT INTO entities \
5348 (id, namespace, kind, entity_type, name, tags, created_at, updated_at) \
5349 VALUES ('model', 'local', 'artifact', 'moodboard_model', \
5350 'model', '[]', 1, 1)",
5351 [],
5352 )
5353 .unwrap();
5354 writer
5355 .conn()
5356 .execute(
5357 "INSERT INTO attachments \
5358 (record_uuid, substrate, role, content_ref, created_at) \
5359 VALUES ('model', 'entity', 'content', ?1, 1), \
5360 ('model', 'entity', 'fann-network', ?2, 1)",
5361 rusqlite::params![bundle.as_str(), network.as_str()],
5362 )
5363 .unwrap();
5364 }
5365
5366 let dry_run = store
5367 .transactional_orphan_sweep(backend.sql().as_ref(), true)
5368 .await
5369 .expect("completed V21 dry run must be supported");
5370 assert_eq!(dry_run.would_delete, 1);
5371 assert_eq!(dry_run.deleted, 0);
5372
5373 let result = store
5374 .transactional_orphan_sweep(backend.sql().as_ref(), false)
5375 .await
5376 .expect("completed V21 destructive sweep must be supported");
5377 assert_eq!(result.deleted, 1);
5378 assert!(store.exists(&bundle).await.unwrap());
5379 assert!(store.exists(&network).await.unwrap());
5380 assert!(!store.exists(&orphan).await.unwrap());
5381 }
5382
5383 #[tokio::test]
5392 async fn transactional_orphan_sweep_refuses_without_the_blob_gc_claims_migration() {
5393 let dir = tempfile::tempdir().unwrap();
5394 let db_path = dir.path().join("khive.db");
5395 let backend =
5396 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
5397 backend.entities().unwrap();
5398 {
5399 let reader = backend.pool().reader().unwrap();
5400 let present: bool = reader
5401 .conn()
5402 .query_row(
5403 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = 'table' \
5404 AND name = 'blob_gc_claims'",
5405 [],
5406 |row| row.get(0),
5407 )
5408 .unwrap();
5409 assert!(
5410 !present,
5411 "this test's premise requires blob_gc_claims to be absent"
5412 );
5413 }
5414
5415 let root = dir.path().join("blobs");
5416 let store = std::sync::Arc::new(
5417 FsBlobStore::new(root.clone(), 0)
5418 .unwrap()
5419 .with_orphan_sweep_grace(Duration::ZERO),
5420 );
5421 let orphan = store.put(b"direct-backend orphan".to_vec()).await.unwrap();
5422
5423 let sql = backend.sql();
5424 let error = store
5425 .transactional_orphan_sweep(sql.as_ref(), false)
5426 .await
5427 .expect_err("sweep must refuse a backend without the blob_gc_claims fencing set");
5428 assert!(
5429 matches!(error, StorageError::Unsupported { .. }),
5430 "expected StorageError::Unsupported, got {error:?}"
5431 );
5432 assert!(
5433 store.exists(&orphan).await.unwrap(),
5434 "a refused sweep must not have deleted anything"
5435 );
5436 }
5437
5438 #[tokio::test]
5439 async fn transactional_orphan_sweep_refuses_an_incomplete_cutover_marker() {
5440 let dir = tempfile::tempdir().unwrap();
5441 let db_path = dir.path().join("khive.db");
5442 let backend =
5443 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
5444 {
5445 let mut writer = backend.pool().writer().unwrap();
5446 prepare_completed_v21_gc_fixture(writer.conn_mut());
5447 writer
5448 .conn_mut()
5449 .execute_batch(
5450 "UPDATE attachment_cutover_state \
5451 SET state = 'incomplete', completed_at = NULL WHERE singleton = 1; \
5452 DELETE FROM _schema_migrations WHERE version = 21;",
5453 )
5454 .unwrap();
5455 }
5456
5457 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
5458 .unwrap()
5459 .with_orphan_sweep_grace(Duration::ZERO);
5460 let orphan = store
5461 .put(b"incomplete-cutover orphan".to_vec())
5462 .await
5463 .unwrap();
5464 let error = store
5465 .transactional_orphan_sweep(backend.sql().as_ref(), false)
5466 .await
5467 .expect_err("sweep must refuse every durable incomplete marker");
5468 assert!(matches!(error, StorageError::Unsupported { .. }));
5469 assert!(
5470 store.exists(&orphan).await.unwrap(),
5471 "refused incomplete-state sweep must preserve every blob"
5472 );
5473 }
5474
5475 #[tokio::test]
5480 async fn transactional_orphan_sweep_refuses_with_incomplete_fencing_triggers() {
5481 let dir = tempfile::tempdir().unwrap();
5482 let db_path = dir.path().join("khive.db");
5483 let backend =
5484 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
5485 {
5486 let mut writer = backend.pool().writer().unwrap();
5487 prepare_completed_v21_gc_fixture(writer.conn_mut());
5488 writer
5489 .conn_mut()
5490 .execute_batch("DROP TRIGGER attachments_reject_claimed_blob_update")
5491 .unwrap();
5492 writer
5493 .conn_mut()
5494 .execute(
5495 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5496 VALUES ('abandoned-partial-fence', \
5497 'dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd', \
5498 1)",
5499 [],
5500 )
5501 .unwrap();
5502 }
5503
5504 let root = dir.path().join("blobs");
5505 let store = std::sync::Arc::new(
5506 FsBlobStore::new(root.clone(), 0)
5507 .unwrap()
5508 .with_orphan_sweep_grace(Duration::ZERO),
5509 );
5510 let orphan = store.put(b"partial-fence orphan".to_vec()).await.unwrap();
5511
5512 let sql = backend.sql();
5513 let _root_guard = store.write_lock.clone().lock_owned().await;
5514 let error = tokio::time::timeout(
5515 Duration::from_secs(1),
5516 store.transactional_orphan_sweep(sql.as_ref(), false),
5517 )
5518 .await
5519 .expect("an incomplete V21 fence must refuse before the root wait")
5520 .expect_err("sweep must refuse when any V21 fencing trigger is missing");
5521 assert!(
5522 matches!(error, StorageError::Unsupported { .. }),
5523 "expected StorageError::Unsupported, got {error:?}"
5524 );
5525 assert!(
5526 store.exists(&orphan).await.unwrap(),
5527 "a refused sweep must not have deleted anything"
5528 );
5529 let remaining: i64 = backend
5530 .pool()
5531 .reader()
5532 .unwrap()
5533 .conn()
5534 .query_row(
5535 "SELECT COUNT(*) FROM blob_gc_claims \
5536 WHERE root_key = 'abandoned-partial-fence'",
5537 [],
5538 |row| row.get(0),
5539 )
5540 .unwrap();
5541 assert_eq!(remaining, 1, "a refused sweep must not recover claims");
5542 }
5543
5544 #[tokio::test]
5550 async fn transactional_orphan_sweep_refuses_same_named_noop_fencing_triggers() {
5551 let dir = tempfile::tempdir().unwrap();
5552 let db_path = dir.path().join("khive.db");
5553 let backend =
5554 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
5555 {
5556 let mut writer = backend.pool().writer().unwrap();
5557 prepare_completed_v21_gc_fixture(writer.conn_mut());
5558 writer
5559 .conn_mut()
5560 .execute_batch(
5561 "DROP TRIGGER attachments_reject_claimed_blob_insert; \
5562 DROP TRIGGER attachments_reject_claimed_blob_update; \
5563 CREATE TRIGGER attachments_reject_claimed_blob_insert \
5564 BEFORE INSERT ON attachments BEGIN SELECT 0; END; \
5565 CREATE TRIGGER attachments_reject_claimed_blob_update \
5566 BEFORE UPDATE OF content_ref ON attachments \
5567 BEGIN SELECT 0; END;",
5568 )
5569 .unwrap();
5570 }
5571
5572 let root = dir.path().join("blobs");
5573 let store = std::sync::Arc::new(
5574 FsBlobStore::new(root.clone(), 0)
5575 .unwrap()
5576 .with_orphan_sweep_grace(Duration::ZERO),
5577 );
5578 let orphan = store.put(b"noop-trigger orphan".to_vec()).await.unwrap();
5579
5580 let sql = backend.sql();
5581 let error = store
5582 .transactional_orphan_sweep(sql.as_ref(), false)
5583 .await
5584 .expect_err("sweep must refuse when the fencing triggers are same-named no-ops");
5585 assert!(
5586 matches!(error, StorageError::Unsupported { .. }),
5587 "expected StorageError::Unsupported, got {error:?}"
5588 );
5589 assert!(
5590 store.exists(&orphan).await.unwrap(),
5591 "a refused sweep must not have deleted anything"
5592 );
5593
5594 let reader = backend.pool().reader().unwrap();
5596 let leftovers: i64 = reader
5597 .conn()
5598 .query_row(
5599 "SELECT (SELECT COUNT(*) FROM blob_gc_claims \
5600 WHERE root_key GLOB '__fence_probe-*') \
5601 + (SELECT COUNT(*) FROM attachments \
5602 WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
5603 [],
5604 |row| row.get(0),
5605 )
5606 .unwrap();
5607 assert_eq!(leftovers, 0, "fence probe rows must not survive the probe");
5608 }
5609
5610 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5611 async fn fence_probe_refuses_id_collision_and_preserves_the_colliding_attachment() {
5612 let dir = tempfile::tempdir().unwrap();
5613 let db_path = dir.path().join("khive.db");
5614 let backend =
5615 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
5616 {
5617 let mut writer = backend.pool().writer().unwrap();
5618 prepare_completed_v21_gc_fixture(writer.conn_mut());
5619 writer
5620 .conn_mut()
5621 .execute(
5622 "INSERT INTO attachments \
5623 (record_uuid, substrate, role, content_ref, media_type, created_at) \
5624 VALUES ('victim-id', 'entity', 'content', \
5625 '2222222222222222222222222222222222222222222222222222222222222222', \
5626 'application/test', 7)",
5627 [],
5628 )
5629 .unwrap();
5630 }
5631
5632 let sql = backend.sql();
5633 let error = super::blob_gc_fence_probe_with_ids(
5634 sql.as_ref(),
5635 "victim-id".to_string(),
5636 "victim-update-id".to_string(),
5637 "victim-insert2-id".to_string(),
5638 "victim-update2-id".to_string(),
5639 "victim-claim-key".to_string(),
5640 )
5641 .await
5642 .expect_err("the probe must refuse when an id it would delete already names a row");
5643 assert!(
5644 matches!(
5645 &error,
5646 StorageError::WriterTaskRequestFailed {
5647 request_state:
5648 khive_storage::WriterTaskRequestState::TransactionRolledBack,
5649 source,
5650 } if matches!(source.as_ref(), StorageError::Unsupported { .. })
5651 ),
5652 "expected a proven-rollback wrapper retaining StorageError::Unsupported, got {error:?}"
5653 );
5654
5655 let reader = backend.pool().reader().unwrap();
5656 let (media_type, created_at): (String, i64) = reader
5657 .conn()
5658 .query_row(
5659 "SELECT media_type, created_at FROM attachments \
5660 WHERE record_uuid = 'victim-id' AND role = 'content'",
5661 [],
5662 |row| Ok((row.get(0)?, row.get(1)?)),
5663 )
5664 .expect("the colliding attachment must survive the refused probe untouched");
5665 assert_eq!(media_type, "application/test");
5666 assert_eq!(created_at, 7);
5667 }
5668
5669 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5670 async fn fence_probe_does_not_touch_an_unrelated_retained_entity_sequence() {
5671 let dir = tempfile::tempdir().unwrap();
5672 let db_path = dir.path().join("khive.db");
5673 let backend =
5674 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
5675 {
5676 let mut writer = backend.pool().writer().unwrap();
5677 prepare_completed_v21_gc_fixture(writer.conn_mut());
5678 writer
5682 .conn_mut()
5683 .execute(
5684 "INSERT INTO entities \
5685 (id, namespace, kind, name, tags, created_at, updated_at) \
5686 VALUES ('retained-id', 'local', 'document', 'gone entity', '[]', 7, 7)",
5687 [],
5688 )
5689 .unwrap();
5690 writer
5691 .conn_mut()
5692 .execute("DELETE FROM entities WHERE id = 'retained-id'", [])
5693 .unwrap();
5694 let retained: i64 = writer
5695 .conn_mut()
5696 .query_row(
5697 "SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
5698 [],
5699 |row| row.get(0),
5700 )
5701 .unwrap();
5702 assert_eq!(retained, 1, "fixture requires a retained-only ledger row");
5703 }
5704
5705 let sql = backend.sql();
5706 super::blob_gc_fence_probe_with_ids(
5707 sql.as_ref(),
5708 "retained-id".to_string(),
5709 "retained-update-id".to_string(),
5710 "retained-insert2-id".to_string(),
5711 "retained-update2-id".to_string(),
5712 "retained-claim-key".to_string(),
5713 )
5714 .await
5715 .expect("attachment probe has no reason to mutate an entity sequence row");
5716
5717 let reader = backend.pool().reader().unwrap();
5718 let survivors: i64 = reader
5719 .conn()
5720 .query_row(
5721 "SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
5722 [],
5723 |row| row.get(0),
5724 )
5725 .unwrap();
5726 assert_eq!(
5727 survivors, 1,
5728 "the retained entity ledger row must survive the attachment probe"
5729 );
5730 }
5731
5732 fn nul_embedded_canonical_ref() -> String {
5733 let mut polluted = "a".repeat(64);
5734 polluted.push('\0');
5735 polluted.push_str("zz");
5736 polluted
5737 }
5738
5739 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5744 async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref() {
5745 let dir = tempfile::tempdir().unwrap();
5746 let db_path = dir.path().join("khive.db");
5747 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5748 {
5749 let mut writer = backend.pool().writer().unwrap();
5750 prepare_completed_v21_gc_fixture(writer.conn_mut());
5751 writer
5752 .conn_mut()
5753 .execute(
5754 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5755 VALUES ('nul-claim-key', ?1, 0)",
5756 rusqlite::params![nul_embedded_canonical_ref()],
5757 )
5758 .unwrap();
5759 }
5760
5761 let sql = backend.sql();
5762 let error = super::validate_blob_gc_evidence(sql.as_ref())
5763 .await
5764 .expect_err("a NUL-embedded claim ref must refuse the sweep");
5765 assert!(
5766 error.to_string().contains("blob_gc_claims"),
5767 "expected the claims-table refusal, got {error:?}"
5768 );
5769 }
5770
5771 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5776 async fn blob_gc_evidence_rejects_a_nul_embedded_attachment_ref() {
5777 let dir = tempfile::tempdir().unwrap();
5778 let db_path = dir.path().join("khive.db");
5779 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5780 {
5781 let mut writer = backend.pool().writer().unwrap();
5782 prepare_completed_v21_gc_fixture(writer.conn_mut());
5783 writer
5788 .conn_mut()
5789 .execute_batch("PRAGMA ignore_check_constraints = ON")
5790 .unwrap();
5791 writer
5792 .conn_mut()
5793 .execute(
5794 "INSERT INTO attachments \
5795 (record_uuid, substrate, role, content_ref, created_at) \
5796 VALUES ('nul-attachment-id', 'entity', 'content', ?1, 0)",
5797 rusqlite::params![nul_embedded_canonical_ref()],
5798 )
5799 .expect("ignore_check_constraints must allow the corrupt row to insert");
5800 writer
5801 .conn_mut()
5802 .execute_batch("PRAGMA ignore_check_constraints = OFF")
5803 .unwrap();
5804 }
5805
5806 let sql = backend.sql();
5807 let error = super::validate_blob_gc_evidence(sql.as_ref())
5808 .await
5809 .expect_err("a NUL-embedded attachment ref must refuse the sweep");
5810 assert!(
5811 error.to_string().contains("attachments"),
5812 "expected the attachments-table refusal, got {error:?}"
5813 );
5814 }
5815
5816 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5820 async fn fence_probe_refuses_a_digest_restricted_trigger_rewrite() {
5821 let dir = tempfile::tempdir().unwrap();
5822 let db_path = dir.path().join("khive.db");
5823 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5824 {
5825 let mut writer = backend.pool().writer().unwrap();
5826 prepare_completed_v21_gc_fixture(writer.conn_mut());
5827 }
5828
5829 let sql = backend.sql();
5830 super::blob_gc_fence_probe(sql.as_ref())
5831 .await
5832 .expect("the healthy fence must pass all four probe arms");
5833
5834 {
5835 let mut writer = backend.pool().writer().unwrap();
5836 writer
5837 .conn_mut()
5838 .execute_batch(
5839 "DROP TRIGGER attachments_reject_claimed_blob_insert; \
5840 DROP TRIGGER attachments_reject_claimed_blob_update; \
5841 CREATE TRIGGER attachments_reject_claimed_blob_insert \
5842 BEFORE INSERT ON attachments \
5843 WHEN NEW.content_ref = \
5844 '0000000000000000000000000000000000000000000000000000000000000000' \
5845 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5846 WHERE content_ref = NEW.content_ref) \
5847 BEGIN \
5848 SELECT RAISE(ABORT, \
5849 'content_ref is reserved by an active blob sweep'); \
5850 END; \
5851 CREATE TRIGGER attachments_reject_claimed_blob_update \
5852 BEFORE UPDATE OF content_ref ON attachments \
5853 WHEN NEW.content_ref = \
5854 '0000000000000000000000000000000000000000000000000000000000000000' \
5855 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5856 WHERE content_ref = NEW.content_ref) \
5857 BEGIN \
5858 SELECT RAISE(ABORT, \
5859 'content_ref is reserved by an active blob sweep'); \
5860 END;",
5861 )
5862 .unwrap();
5863 }
5864
5865 let error = super::blob_gc_fence_probe(sql.as_ref())
5866 .await
5867 .expect_err("a sentinel-only fence must fail the second-digest arms");
5868 assert!(
5869 matches!(error, StorageError::Unsupported { .. }),
5870 "expected StorageError::Unsupported, got {error:?}"
5871 );
5872 assert!(
5873 error.to_string().contains("second-digest"),
5874 "the refusal must name a second-digest arm, got {error}"
5875 );
5876 }
5877
5878 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5881 async fn fence_probe_refuses_a_shape_restricted_trigger_rewrite() {
5882 let dir = tempfile::tempdir().unwrap();
5883 let db_path = dir.path().join("khive.db");
5884 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5885 {
5886 let mut writer = backend.pool().writer().unwrap();
5887 prepare_completed_v21_gc_fixture(writer.conn_mut());
5888 writer
5889 .conn_mut()
5890 .execute_batch(
5891 "DROP TRIGGER attachments_reject_claimed_blob_insert; \
5892 DROP TRIGGER attachments_reject_claimed_blob_update; \
5893 CREATE TRIGGER attachments_reject_claimed_blob_insert \
5894 BEFORE INSERT ON attachments \
5895 WHEN NEW.substrate = 'entity' \
5896 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5897 WHERE content_ref = NEW.content_ref) \
5898 BEGIN \
5899 SELECT RAISE(ABORT, \
5900 'content_ref is reserved by an active blob sweep'); \
5901 END; \
5902 CREATE TRIGGER attachments_reject_claimed_blob_update \
5903 BEFORE UPDATE OF content_ref ON attachments \
5904 WHEN NEW.substrate = 'entity' \
5905 AND EXISTS (SELECT 1 FROM blob_gc_claims \
5906 WHERE content_ref = NEW.content_ref) \
5907 BEGIN \
5908 SELECT RAISE(ABORT, \
5909 'content_ref is reserved by an active blob sweep'); \
5910 END;",
5911 )
5912 .unwrap();
5913 }
5914
5915 let sql = backend.sql();
5916 let error = super::blob_gc_fence_probe(sql.as_ref())
5917 .await
5918 .expect_err("an entity-shape-only fence must fail the note-shaped arms");
5919 assert!(
5920 matches!(error, StorageError::Unsupported { .. }),
5921 "expected StorageError::Unsupported, got {error:?}"
5922 );
5923 assert!(
5924 error.to_string().contains("second-digest"),
5925 "the refusal must name a second-digest arm, got {error}"
5926 );
5927 }
5928
5929 fn initialize_utf16le_database(db_path: &std::path::Path) {
5933 let conn = rusqlite::Connection::open(db_path).unwrap();
5934 conn.execute_batch(
5935 "PRAGMA encoding = 'UTF-16le'; \
5936 CREATE TABLE __encoding_pin (x INTEGER); \
5937 DROP TABLE __encoding_pin;",
5938 )
5939 .unwrap();
5940 }
5941
5942 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5946 async fn blob_gc_evidence_accepts_valid_refs_in_a_utf16le_database() {
5947 let dir = tempfile::tempdir().unwrap();
5948 let db_path = dir.path().join("khive.db");
5949 initialize_utf16le_database(&db_path);
5950 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5951 {
5952 let mut writer = backend.pool().writer().unwrap();
5953 let encoding: String = writer
5954 .conn_mut()
5955 .query_row("PRAGMA encoding", [], |row| row.get(0))
5956 .unwrap();
5957 assert_eq!(
5958 encoding, "UTF-16le",
5959 "the fixture database must actually be UTF-16le"
5960 );
5961 prepare_completed_v21_gc_fixture(writer.conn_mut());
5962 writer
5967 .conn_mut()
5968 .execute(
5969 "INSERT INTO attachments \
5970 (record_uuid, substrate, role, content_ref, created_at) \
5971 VALUES ('utf16-valid-attachment', 'entity', 'content', ?1, 0)",
5972 rusqlite::params!["a".repeat(64)],
5973 )
5974 .unwrap();
5975 writer
5976 .conn_mut()
5977 .execute(
5978 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5979 VALUES ('utf16-valid-claim-key', ?1, 0)",
5980 rusqlite::params!["b".repeat(64)],
5981 )
5982 .unwrap();
5983 }
5984
5985 let sql = backend.sql();
5986 super::validate_blob_gc_evidence(sql.as_ref())
5987 .await
5988 .expect("valid canonical refs must pass in a UTF-16LE database");
5989 }
5990
5991 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5994 async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref_in_a_utf16le_database() {
5995 let dir = tempfile::tempdir().unwrap();
5996 let db_path = dir.path().join("khive.db");
5997 initialize_utf16le_database(&db_path);
5998 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
5999 {
6000 let mut writer = backend.pool().writer().unwrap();
6001 let encoding: String = writer
6002 .conn_mut()
6003 .query_row("PRAGMA encoding", [], |row| row.get(0))
6004 .unwrap();
6005 assert_eq!(
6006 encoding, "UTF-16le",
6007 "the fixture database must actually be UTF-16le"
6008 );
6009 prepare_completed_v21_gc_fixture(writer.conn_mut());
6010 writer
6011 .conn_mut()
6012 .execute(
6013 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
6014 VALUES ('nul-claim-key-utf16', ?1, 0)",
6015 rusqlite::params![nul_embedded_canonical_ref()],
6016 )
6017 .unwrap();
6018 }
6019
6020 let sql = backend.sql();
6021 let error = super::validate_blob_gc_evidence(sql.as_ref())
6022 .await
6023 .expect_err("a NUL-embedded claim ref must refuse the sweep in UTF-16LE too");
6024 assert!(
6025 error.to_string().contains("blob_gc_claims"),
6026 "expected the claims-table refusal, got {error:?}"
6027 );
6028 }
6029
6030 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6031 async fn transactional_orphan_sweep_preserves_put_started_after_liveness_mark() {
6032 let dir = tempfile::tempdir().unwrap();
6033 let db_path = dir.path().join("khive.db");
6034 let backend =
6035 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
6036 {
6037 let mut writer = backend.pool().writer().unwrap();
6038 prepare_completed_v21_gc_fixture(writer.conn_mut());
6039 }
6040 let root = dir.path().join("blobs");
6041 let store = std::sync::Arc::new(
6042 FsBlobStore::new(root.clone(), 0)
6043 .unwrap()
6044 .with_orphan_sweep_grace(Duration::ZERO),
6045 );
6046 let orphan = store.put(b"old orphan".to_vec()).await.unwrap();
6047 let canonical_root = root.canonicalize().unwrap();
6048 let (marked, release, _done) = sync_hook::install(&canonical_root);
6049
6050 let sweep = {
6051 let store = store.clone();
6052 let sql = backend.sql();
6053 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
6054 };
6055 assert!(
6056 recv_blocking(marked).await,
6057 "sweep must finish its liveness mark"
6058 );
6059
6060 assert!(
6061 store.write_lock.try_lock().is_err(),
6062 "the sweep must hold the same root lock used by blob writers"
6063 );
6064 let (started_tx, started_rx) = std::sync::mpsc::channel();
6065 let new_ref = {
6066 let root = root.clone();
6067 tokio::task::spawn_blocking(move || {
6068 let _ = started_tx.send(());
6069 put_blocking(&root, 0, b"new concurrent blob".to_vec())
6070 })
6071 };
6072 assert!(recv_blocking(started_rx).await, "blob put must start");
6073
6074 release.send(()).unwrap();
6075 let sweep_result = sweep.await.unwrap().unwrap();
6076 let new_ref = new_ref.await.unwrap().unwrap();
6077
6078 assert_eq!(sweep_result.deleted, 1);
6079 assert!(!store.exists(&orphan).await.unwrap());
6080 assert!(
6081 store.exists(&new_ref).await.unwrap(),
6082 "a blob put started between the liveness mark and physical sweep must survive"
6083 );
6084 }
6085
6086 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6087 async fn transactional_orphan_sweep_releases_sqlite_writer_before_physical_delete() {
6088 let dir = tempfile::tempdir().unwrap();
6089 let db_path = dir.path().join("khive.db");
6090 let backend =
6091 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
6092 {
6093 let mut writer = backend.pool().writer().unwrap();
6094 prepare_completed_v21_gc_fixture(writer.conn_mut());
6095 }
6096 let root = dir.path().join("blobs");
6097 let store = std::sync::Arc::new(
6098 FsBlobStore::new(root.clone(), 0)
6099 .unwrap()
6100 .with_orphan_sweep_grace(Duration::ZERO),
6101 );
6102 let orphan = store
6103 .put(b"claim then delete outside sqlite".to_vec())
6104 .await
6105 .unwrap();
6106 let canonical_root = root.canonicalize().unwrap();
6107 let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
6108
6109 let sweep = {
6110 let store = store.clone();
6111 let sql = backend.sql();
6112 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
6113 };
6114 assert!(
6115 recv_blocking(claimed).await,
6116 "sweep must durably claim the orphan before physical deletion"
6117 );
6118 assert!(
6119 store.exists(&orphan).await.unwrap(),
6120 "the test seam must pause before the physical delete"
6121 );
6122 let external_database_lock = fs::OpenOptions::new()
6123 .read(true)
6124 .write(true)
6125 .open(database_gc_lock_path(&db_path))
6126 .unwrap();
6127 assert!(
6128 matches!(
6129 fs4::FileExt::try_lock(&external_database_lock),
6130 Err(fs4::TryLockError::WouldBlock)
6131 ),
6132 "the sweep must retain cross-process database ownership while SQLite's writer is free"
6133 );
6134
6135 let unrelated = rusqlite::Connection::open(&db_path).unwrap();
6140 unrelated.busy_timeout(Duration::from_millis(100)).unwrap();
6141 unrelated
6142 .execute(
6143 "INSERT INTO entities \
6144 (id, namespace, kind, name, tags, created_at, updated_at) \
6145 VALUES ('unrelated-writer', 'local', 'concept', 'unrelated', '[]', 1, 1)",
6146 [],
6147 )
6148 .expect("external filesystem work must not retain SQLite's writer lock");
6149
6150 let claimed_err = unrelated
6154 .execute(
6155 "INSERT INTO attachments \
6156 (record_uuid, substrate, role, content_ref, created_at) \
6157 VALUES ('racing-reference', 'entity', 'content', ?1, 1)",
6158 [orphan.as_str()],
6159 )
6160 .expect_err("a claimed content_ref must fail closed before deletion");
6161 assert!(
6162 claimed_err.to_string().contains("active blob sweep"),
6163 "unexpected claim error: {claimed_err}"
6164 );
6165
6166 release_delete.send(()).unwrap();
6167 let result = sweep.await.unwrap().unwrap();
6168 assert_eq!(result.deleted, 1);
6169 assert!(!store.exists(&orphan).await.unwrap());
6170
6171 let remaining_claims: i64 = unrelated
6172 .query_row(
6173 "SELECT COUNT(*) FROM blob_gc_claims WHERE content_ref = ?1",
6174 [orphan.as_str()],
6175 |row| row.get(0),
6176 )
6177 .unwrap();
6178 assert_eq!(
6179 remaining_claims, 0,
6180 "successful deletion releases the claim"
6181 );
6182 }
6183
6184 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6185 async fn cancelling_sweep_during_delete_keeps_owner_locks_until_blocking_work_finishes() {
6186 let dir = tempfile::tempdir().unwrap();
6187 let db_path = dir.path().join("khive.db");
6188 let backend =
6189 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
6190 {
6191 let mut writer = backend.pool().writer().unwrap();
6192 prepare_completed_v21_gc_fixture(writer.conn_mut());
6193 }
6194 let root = dir.path().join("blobs");
6195 let store = std::sync::Arc::new(
6196 FsBlobStore::new(root.clone(), 0)
6197 .unwrap()
6198 .with_orphan_sweep_grace(Duration::ZERO),
6199 );
6200 let orphan = store.put(b"cancelled sweep orphan".to_vec()).await.unwrap();
6201 let canonical_root = root.canonicalize().unwrap();
6202 let (claimed, release_delete, done) = sync_hook::install(&canonical_root);
6203
6204 let sweep = {
6205 let store = store.clone();
6206 let sql = backend.sql();
6207 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
6208 };
6209 assert!(recv_blocking(claimed).await);
6210 sweep.abort();
6211 assert!(sweep.await.unwrap_err().is_cancelled());
6212
6213 let external_root_lock = fs::OpenOptions::new()
6214 .read(true)
6215 .write(true)
6216 .open(root.join(ROOT_WRITE_LOCK_FILE))
6217 .unwrap();
6218 let external_database_lock = fs::OpenOptions::new()
6219 .read(true)
6220 .write(true)
6221 .open(database_gc_lock_path(&db_path))
6222 .unwrap();
6223 assert!(matches!(
6224 fs4::FileExt::try_lock(&external_root_lock),
6225 Err(fs4::TryLockError::WouldBlock)
6226 ));
6227 assert!(matches!(
6228 fs4::FileExt::try_lock(&external_database_lock),
6229 Err(fs4::TryLockError::WouldBlock)
6230 ));
6231
6232 release_delete.send(()).unwrap();
6233 let done_disconnected = tokio::task::spawn_blocking(move || done.recv().is_err())
6234 .await
6235 .unwrap();
6236 assert!(
6237 done_disconnected,
6238 "the cancelled outer task cannot send done"
6239 );
6240 assert!(fs4::FileExt::try_lock(&external_root_lock).is_ok());
6241 assert!(fs4::FileExt::try_lock(&external_database_lock).is_ok());
6242 drop(external_root_lock);
6243 drop(external_database_lock);
6244 assert!(!store.exists(&orphan).await.unwrap());
6245 let stranded_claims: i64 = rusqlite::Connection::open(&db_path)
6246 .unwrap()
6247 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
6248 .unwrap();
6249 assert_eq!(
6250 stranded_claims, 1,
6251 "cancellation leaves a fail-closed claim"
6252 );
6253
6254 let recovered = store
6255 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6256 .await
6257 .unwrap();
6258 assert_eq!(recovered.deleted, 0);
6259 let remaining: i64 = rusqlite::Connection::open(&db_path)
6260 .unwrap()
6261 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
6262 .unwrap();
6263 assert_eq!(remaining, 0, "the next exclusive owner recovers the claim");
6264 }
6265
6266 #[tokio::test]
6267 async fn transactional_orphan_sweep_recovers_stale_claims_fail_closed() {
6268 let dir = tempfile::tempdir().unwrap();
6269 let db_path = dir.path().join("khive.db");
6270 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6271 {
6272 let mut writer = backend.pool().writer().unwrap();
6273 prepare_completed_v21_gc_fixture(writer.conn_mut());
6274 }
6275 let root = dir.path().join("blobs");
6276 let store = FsBlobStore::new(root.clone(), 0)
6277 .unwrap()
6278 .with_orphan_sweep_grace(Duration::from_secs(60));
6279 let bytes = b"republished after a crashed claim".to_vec();
6280 let content_ref = store.put(bytes.clone()).await.unwrap();
6281 let canonical_root = root.canonicalize().unwrap();
6282 let root_key = blob_root_key(&canonical_root);
6283 let absent_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
6284 let former_probe_seed = "1111111111111111111111111111111111111111111111111111111111111111";
6285 {
6286 let writer = backend.pool().writer().unwrap();
6287 writer
6288 .conn()
6289 .execute(
6290 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
6291 VALUES (?1, ?2, 1), (?1, ?3, 1), (?1, ?4, 1)",
6292 rusqlite::params![
6293 root_key,
6294 content_ref.as_str(),
6295 absent_ref,
6296 former_probe_seed
6297 ],
6298 )
6299 .unwrap();
6300 }
6301
6302 assert_eq!(store.put(bytes).await.unwrap(), content_ref);
6308 let result = store
6309 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6310 .await
6311 .unwrap();
6312 assert_eq!(result.deleted, 0);
6313 assert_eq!(result.grace_period_skipped, 1);
6314 assert!(store.exists(&content_ref).await.unwrap());
6315
6316 let remaining: i64 = backend
6317 .pool()
6318 .writer()
6319 .unwrap()
6320 .conn()
6321 .query_row(
6322 "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?1",
6323 [blob_root_key(&canonical_root)],
6324 |row| row.get(0),
6325 )
6326 .unwrap();
6327 assert_eq!(remaining, 0, "the next sweep recovers stale claims");
6328 }
6329
6330 #[tokio::test]
6331 async fn transactional_orphan_sweep_recovers_claims_after_root_relocation() {
6332 let dir = tempfile::tempdir().unwrap();
6333 let db_path = dir.path().join("khive.db");
6334 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6335 {
6336 let mut writer = backend.pool().writer().unwrap();
6337 prepare_completed_v21_gc_fixture(writer.conn_mut());
6338 }
6339
6340 let old_root = dir.path().join("old-blobs");
6341 let bytes = b"claim must follow a relocated blob root".to_vec();
6342 let content_ref = {
6343 let old_store = FsBlobStore::new(old_root.clone(), 0)
6344 .unwrap()
6345 .with_orphan_sweep_grace(Duration::from_secs(60));
6346 old_store.put(bytes).await.unwrap()
6347 };
6348 let old_root_key = blob_root_key(&old_root.canonicalize().unwrap());
6349 backend
6350 .pool()
6351 .writer()
6352 .unwrap()
6353 .conn()
6354 .execute(
6355 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
6356 VALUES (?1, ?2, 1)",
6357 rusqlite::params![old_root_key, content_ref.as_str()],
6358 )
6359 .unwrap();
6360
6361 let new_root = dir.path().join("relocated-blobs");
6362 std::fs::rename(&old_root, &new_root).unwrap();
6363 let relocated_store = FsBlobStore::new(new_root, 0)
6364 .unwrap()
6365 .with_orphan_sweep_grace(Duration::from_secs(60));
6366 let result = relocated_store
6367 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6368 .await
6369 .unwrap();
6370
6371 assert_eq!(
6372 result.deleted, 0,
6373 "a fresh relocated blob remains protected"
6374 );
6375 assert_eq!(result.grace_period_skipped, 1);
6376 assert!(relocated_store.exists(&content_ref).await.unwrap());
6377 let remaining: i64 = backend
6378 .pool()
6379 .writer()
6380 .unwrap()
6381 .conn()
6382 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
6383 .unwrap();
6384 assert_eq!(
6385 remaining, 0,
6386 "exclusive database sweep ownership makes every pre-existing claim abandoned, \
6387 even when its old path-derived root key no longer matches"
6388 );
6389 }
6390
6391 #[tokio::test]
6392 async fn transactional_orphan_sweep_recovers_claims_copied_by_database_restore() {
6393 let dir = tempfile::tempdir().unwrap();
6394 let source_path = dir.path().join("source.db");
6395 let restored_path = dir.path().join("restored.db");
6396 let bytes = b"claim copied in an online database backup".to_vec();
6397 let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
6398 {
6399 let source = crate::StorageBackend::sqlite_for_test(&source_path).unwrap();
6400 let mut writer = source.pool().writer().unwrap();
6401 prepare_completed_v21_gc_fixture(writer.conn_mut());
6402 writer
6403 .conn()
6404 .execute(
6405 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
6406 VALUES ('source-root-before-backup', ?1, 1)",
6407 [content_ref.as_str()],
6408 )
6409 .unwrap();
6410 writer
6411 .conn()
6412 .execute_batch("PRAGMA wal_checkpoint(TRUNCATE)")
6413 .unwrap();
6414 }
6415 std::fs::copy(&source_path, &restored_path).unwrap();
6416
6417 let restored = crate::StorageBackend::sqlite_for_test(&restored_path).unwrap();
6418 let restored_root = dir.path().join("restored-blobs");
6419 let store = FsBlobStore::new(restored_root, 0)
6420 .unwrap()
6421 .with_orphan_sweep_grace(Duration::from_secs(60));
6422 assert_eq!(store.put(bytes).await.unwrap(), content_ref);
6423 let result = store
6424 .transactional_orphan_sweep(restored.sql().as_ref(), false)
6425 .await
6426 .unwrap();
6427
6428 assert_eq!(result.deleted, 0);
6429 assert_eq!(result.grace_period_skipped, 1);
6430 let remaining: i64 = restored
6431 .pool()
6432 .writer()
6433 .unwrap()
6434 .conn()
6435 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
6436 .unwrap();
6437 assert_eq!(remaining, 0, "restored claims are abandoned ownership");
6438 }
6439
6440 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6441 async fn transactional_orphan_sweep_bounds_each_durable_claim_batch() {
6442 let dir = tempfile::tempdir().unwrap();
6443 let db_path = dir.path().join("khive.db");
6444 let backend =
6445 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
6446 {
6447 let mut writer = backend.pool().writer().unwrap();
6448 prepare_completed_v21_gc_fixture(writer.conn_mut());
6449 }
6450 let root = dir.path().join("blobs");
6451 let store = std::sync::Arc::new(
6452 FsBlobStore::new(root.clone(), 0)
6453 .unwrap()
6454 .with_orphan_sweep_grace(Duration::ZERO),
6455 );
6456 let candidate_count = BLOB_GC_CLAIM_BATCH_SIZE * 2 + 1;
6457 for index in 0..candidate_count {
6458 store
6459 .put(format!("bounded claim candidate {index}").into_bytes())
6460 .await
6461 .unwrap();
6462 }
6463 let canonical_root = root.canonicalize().unwrap();
6464 let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
6465
6466 let sweep = {
6467 let store = store.clone();
6468 let sql = backend.sql();
6469 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
6470 };
6471 assert!(
6472 recv_blocking(claimed).await,
6473 "the first bounded claim batch must commit before deletion"
6474 );
6475 let active_claims: i64 = rusqlite::Connection::open(&db_path)
6476 .unwrap()
6477 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
6478 .unwrap();
6479 assert!(active_claims > 0);
6480 assert!(
6481 active_claims <= BLOB_GC_CLAIM_BATCH_SIZE as i64,
6482 "one transaction may expose at most {BLOB_GC_CLAIM_BATCH_SIZE} claim rows; \
6483 observed {active_claims}"
6484 );
6485
6486 release_delete.send(()).unwrap();
6487 let result = sweep.await.unwrap().unwrap();
6488 assert_eq!(result.deleted, candidate_count as u64);
6489 let remaining: i64 = rusqlite::Connection::open(&db_path)
6490 .unwrap()
6491 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
6492 .unwrap();
6493 assert_eq!(remaining, 0);
6494 }
6495
6496 #[tokio::test]
6497 async fn abandoned_claim_recovery_deletes_at_most_one_batch_per_writer_hold() {
6498 let dir = tempfile::tempdir().unwrap();
6499 let db_path = dir.path().join("khive.db");
6500 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6501 {
6502 let mut writer = backend.pool().writer().unwrap();
6503 crate::run_migrations(writer.conn_mut()).unwrap();
6504 let tx = writer.conn_mut().transaction().unwrap();
6505 for index in 0..(BLOB_GC_CLAIM_BATCH_SIZE + 1) {
6506 tx.execute(
6507 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
6508 VALUES ('abandoned-root', ?1, 1)",
6509 [format!("{index:064x}")],
6510 )
6511 .unwrap();
6512 }
6513 tx.commit().unwrap();
6514 }
6515
6516 let released = release_abandoned_blob_gc_claim_batch(backend.sql().as_ref())
6517 .await
6518 .unwrap();
6519 assert_eq!(released, BLOB_GC_CLAIM_BATCH_SIZE as u64);
6520 let remaining: i64 = backend
6521 .pool()
6522 .writer()
6523 .unwrap()
6524 .conn()
6525 .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
6526 .unwrap();
6527 assert_eq!(remaining, 1);
6528 }
6529
6530 #[tokio::test]
6531 async fn transactional_orphan_sweep_refuses_corrupt_liveness_and_claim_evidence() {
6532 let dir = tempfile::tempdir().unwrap();
6533 let db_path = dir.path().join("khive.db");
6534 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6535 {
6536 let mut writer = backend.pool().writer().unwrap();
6537 prepare_completed_v21_gc_fixture(writer.conn_mut());
6538 }
6539 let root = dir.path().join("blobs");
6540 let store = FsBlobStore::new(root.clone(), 0)
6541 .unwrap()
6542 .with_orphan_sweep_grace(Duration::ZERO);
6543 let orphan = store
6544 .put(b"must survive corrupt evidence".to_vec())
6545 .await
6546 .unwrap();
6547
6548 let conn = rusqlite::Connection::open(&db_path).unwrap();
6549 conn.execute_batch("PRAGMA ignore_check_constraints = ON")
6550 .unwrap();
6551 conn.execute(
6552 "INSERT INTO attachments \
6553 (record_uuid, substrate, role, content_ref, created_at) \
6554 VALUES ('corrupt-live', 'entity', 'content', 'not-a-content-ref', 1)",
6555 [],
6556 )
6557 .unwrap();
6558 conn.execute_batch("PRAGMA ignore_check_constraints = OFF")
6559 .unwrap();
6560 let live_error = store
6561 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6562 .await
6563 .expect_err("corrupt live evidence must fail closed");
6564 assert!(matches!(live_error, StorageError::InvalidInput { .. }));
6565 assert!(
6566 store.exists(&orphan).await.unwrap(),
6567 "no file may be removed after corrupt live evidence"
6568 );
6569
6570 conn.execute(
6571 "DELETE FROM attachments WHERE record_uuid = 'corrupt-live'",
6572 [],
6573 )
6574 .unwrap();
6575 let root_key = blob_root_key(&root.canonicalize().unwrap());
6576 conn.execute(
6577 "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
6578 VALUES (?1, 'also-not-a-content-ref', 1)",
6579 [root_key.as_str()],
6580 )
6581 .unwrap();
6582 let claim_error = store
6583 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6584 .await
6585 .expect_err("corrupt durable claim evidence must fail closed");
6586 assert!(matches!(claim_error, StorageError::InvalidInput { .. }));
6587 assert!(
6588 store.exists(&orphan).await.unwrap(),
6589 "no file may be removed after corrupt claim evidence"
6590 );
6591 let remaining: i64 = conn
6592 .query_row(
6593 "SELECT COUNT(*) FROM blob_gc_claims \
6594 WHERE root_key = ?1 AND content_ref = 'also-not-a-content-ref'",
6595 [root_key.as_str()],
6596 |row| row.get(0),
6597 )
6598 .unwrap();
6599 assert_eq!(
6600 remaining, 1,
6601 "corrupt claim evidence is not silently erased"
6602 );
6603 let probe_residue: i64 = conn
6604 .query_row(
6605 "SELECT (SELECT COUNT(*) FROM blob_gc_claims \
6606 WHERE root_key GLOB '__fence_probe-*') \
6607 + (SELECT COUNT(*) FROM attachments \
6608 WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
6609 [],
6610 |row| row.get(0),
6611 )
6612 .unwrap();
6613 assert_eq!(
6614 probe_residue, 0,
6615 "invalid evidence must abort before the functional fence probe"
6616 );
6617 }
6618
6619 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6620 async fn transactional_orphan_sweep_republishes_deduplicated_external_put() {
6621 let dir = tempfile::tempdir().unwrap();
6622 let db_path = dir.path().join("khive.db");
6623 let backend =
6624 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
6625 {
6626 let mut writer = backend.pool().writer().unwrap();
6627 prepare_completed_v21_gc_fixture(writer.conn_mut());
6628 }
6629 let root = dir.path().join("blobs");
6630 let store = std::sync::Arc::new(
6631 FsBlobStore::new(root.clone(), 0)
6632 .unwrap()
6633 .with_orphan_sweep_grace(Duration::ZERO),
6634 );
6635 let payload = b"existing orphan republished during sweep".to_vec();
6636 let orphan = store.put(payload.clone()).await.unwrap();
6637 let canonical_root = root.canonicalize().unwrap();
6638 let (marked, release, _done) = sync_hook::install(&canonical_root);
6639
6640 let sweep = {
6641 let store = store.clone();
6642 let sql = backend.sql();
6643 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
6644 };
6645 assert!(
6646 recv_blocking(marked).await,
6647 "sweep must finish its liveness mark"
6648 );
6649
6650 let external_lock = fs::OpenOptions::new()
6651 .read(true)
6652 .write(true)
6653 .open(root.join(ROOT_WRITE_LOCK_FILE))
6654 .unwrap();
6655 assert!(
6656 matches!(
6657 fs4::FileExt::try_lock(&external_lock),
6658 Err(fs4::TryLockError::WouldBlock)
6659 ),
6660 "the sweep must exclude a publisher using an independently opened root lock"
6661 );
6662
6663 let (started_tx, started_rx) = std::sync::mpsc::channel();
6664 let republished = {
6665 let root = root.clone();
6666 tokio::task::spawn_blocking(move || {
6667 let _ = started_tx.send(());
6668 put_blocking(&root, 0, payload)
6669 })
6670 };
6671 assert!(recv_blocking(started_rx).await, "blob put must start");
6672
6673 release.send(()).unwrap();
6674 let sweep_result = sweep.await.unwrap().unwrap();
6675 let republished = republished.await.unwrap().unwrap();
6676
6677 assert_eq!(sweep_result.deleted, 1);
6678 assert_eq!(republished, orphan);
6679 assert!(
6680 store.exists(&republished).await.unwrap(),
6681 "a deduplicated put concurrent with the sweep must not return a deleted reference"
6682 );
6683 }
6684
6685 #[tokio::test]
6686 async fn transactional_orphan_sweep_uses_all_attachment_refs_as_live() {
6687 let dir = tempfile::tempdir().unwrap();
6688 let db_path = dir.path().join("khive.db");
6689 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6690 {
6691 let mut writer = backend.pool().writer().unwrap();
6692 prepare_completed_v21_gc_fixture(writer.conn_mut());
6693 }
6694 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6695 .unwrap()
6696 .with_orphan_sweep_grace(Duration::ZERO);
6697 let live = store.put(b"live".to_vec()).await.unwrap();
6698 let soft_deleted = store.put(b"soft deleted".to_vec()).await.unwrap();
6699 let orphan = store.put(b"orphan".to_vec()).await.unwrap();
6700 {
6701 let writer = backend.pool().writer().unwrap();
6702 writer
6703 .conn()
6704 .execute_batch(
6705 "INSERT INTO entities \
6706 (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6707 VALUES ('live', 'local', 'document', 'live', '[]', 1, 1, NULL), \
6708 ('deleted', 'local', 'document', 'deleted', '[]', 1, 1, 2);",
6709 )
6710 .unwrap();
6711 writer
6712 .conn()
6713 .execute(
6714 "INSERT INTO attachments \
6715 (record_uuid, substrate, role, content_ref, created_at) \
6716 VALUES ('live', 'entity', 'content', ?1, 1), \
6717 ('deleted', 'entity', 'content', ?2, 1)",
6718 rusqlite::params![live.as_str(), soft_deleted.as_str()],
6719 )
6720 .unwrap();
6721 }
6722
6723 let dry_run = store
6724 .transactional_orphan_sweep(backend.sql().as_ref(), true)
6725 .await
6726 .unwrap();
6727 assert_eq!(dry_run.would_delete, 1);
6728 assert_eq!(dry_run.deleted, 0);
6729 assert!(store.exists(&soft_deleted).await.unwrap());
6730 assert!(store.exists(&orphan).await.unwrap());
6731
6732 let result = store
6733 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6734 .await
6735 .unwrap();
6736
6737 assert_eq!(result.scanned, 3);
6738 assert_eq!(result.deleted, 1);
6739 assert!(store.exists(&live).await.unwrap());
6740 assert!(
6741 store.exists(&soft_deleted).await.unwrap(),
6742 "soft delete retains attachment rows and their blobs"
6743 );
6744 assert!(!store.exists(&orphan).await.unwrap());
6745 }
6746
6747 #[tokio::test]
6748 async fn transactional_orphan_sweep_protects_a_freshly_published_blob_before_its_reference_commits(
6749 ) {
6750 let dir = tempfile::tempdir().unwrap();
6763 let db_path = dir.path().join("khive.db");
6764 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6765 {
6766 let mut writer = backend.pool().writer().unwrap();
6767 prepare_completed_v21_gc_fixture(writer.conn_mut());
6768 }
6769 let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
6772
6773 let blob = store
6776 .put(b"published, reference not yet committed".to_vec())
6777 .await
6778 .unwrap();
6779
6780 let result = store
6782 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6783 .await
6784 .unwrap();
6785
6786 assert_eq!(result.deleted, 0, "the blob must survive: {result:?}");
6787 assert_eq!(
6788 result.would_delete, 0,
6789 "not treated as a deletable orphan: {result:?}"
6790 );
6791 assert_eq!(
6792 result.grace_period_skipped, 1,
6793 "must be reported as grace-protected rather than silently ignored: {result:?}"
6794 );
6795 assert!(
6796 store.exists(&blob).await.unwrap(),
6797 "a blob still inside its publish grace period must survive the sweep"
6798 );
6799
6800 {
6803 let writer = backend.pool().writer().unwrap();
6804 writer
6805 .conn()
6806 .execute(
6807 "INSERT INTO entities \
6808 (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6809 VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
6810 [],
6811 )
6812 .unwrap();
6813 writer
6814 .conn()
6815 .execute(
6816 "INSERT INTO attachments \
6817 (record_uuid, substrate, role, content_ref, created_at) \
6818 VALUES ('e1', 'entity', 'content', ?1, 1)",
6819 [blob.as_str()],
6820 )
6821 .unwrap();
6822 }
6823
6824 let result = store
6827 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6828 .await
6829 .unwrap();
6830 assert_eq!(result.deleted, 0);
6831 assert!(store.exists(&blob).await.unwrap());
6832 }
6833
6834 #[tokio::test]
6835 async fn put_republishing_an_aged_orphan_restarts_its_grace_clock_before_the_reference_commits()
6836 {
6837 let dir = tempfile::tempdir().unwrap();
6845 let db_path = dir.path().join("khive.db");
6846 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6847 {
6848 let mut writer = backend.pool().writer().unwrap();
6849 prepare_completed_v21_gc_fixture(writer.conn_mut());
6850 }
6851 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6852 .unwrap()
6853 .with_orphan_sweep_grace(Duration::from_secs(60));
6854
6855 let bytes = b"old orphan re-published".to_vec();
6856 let first = store.put(bytes.clone()).await.unwrap();
6857
6858 let path = shard_path(store.root(), &first);
6861 let old_mtime = SystemTime::now() - Duration::from_secs(3600);
6862 fs::OpenOptions::new()
6863 .write(true)
6864 .open(&path)
6865 .unwrap()
6866 .set_modified(old_mtime)
6867 .unwrap();
6868
6869 let second = store.put(bytes).await.unwrap();
6872 assert_eq!(first, second);
6873
6874 let result = store
6877 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6878 .await
6879 .unwrap();
6880 assert_eq!(
6881 result.deleted, 0,
6882 "a dedup-republished blob must survive a sweep landing before its reference \
6883 commits: {result:?}"
6884 );
6885 assert_eq!(
6886 result.grace_period_skipped, 1,
6887 "must be reported as grace-protected, not silently ignored: {result:?}"
6888 );
6889 assert!(store.exists(&first).await.unwrap());
6890
6891 {
6893 let writer = backend.pool().writer().unwrap();
6894 writer
6895 .conn()
6896 .execute(
6897 "INSERT INTO entities \
6898 (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6899 VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
6900 [],
6901 )
6902 .unwrap();
6903 writer
6904 .conn()
6905 .execute(
6906 "INSERT INTO attachments \
6907 (record_uuid, substrate, role, content_ref, created_at) \
6908 VALUES ('e1', 'entity', 'content', ?1, 1)",
6909 [first.as_str()],
6910 )
6911 .unwrap();
6912 }
6913
6914 let result = store
6915 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6916 .await
6917 .unwrap();
6918 assert_eq!(result.deleted, 0);
6919 assert!(
6920 store.exists(&first).await.unwrap(),
6921 "the blob must stay live once its reference has committed"
6922 );
6923 }
6924
6925 #[tokio::test]
6926 async fn put_dedup_mtime_refresh_has_no_observable_effect_under_zero_grace_period() {
6927 let dir = tempfile::tempdir().unwrap();
6934 let db_path = dir.path().join("khive.db");
6935 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6936 {
6937 let mut writer = backend.pool().writer().unwrap();
6938 prepare_completed_v21_gc_fixture(writer.conn_mut());
6939 }
6940 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6941 .unwrap()
6942 .with_orphan_sweep_grace(Duration::ZERO);
6943 let bytes = b"zero grace dedup refresh".to_vec();
6944 let first = store.put(bytes.clone()).await.unwrap();
6945 let second = store.put(bytes.clone()).await.unwrap();
6946 assert_eq!(first, second);
6947
6948 let result = store
6949 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6950 .await
6951 .unwrap();
6952 assert_eq!(
6953 result.deleted, 1,
6954 "a zero grace period must still delete an unreferenced blob even after a dedup \
6955 put refreshed its mtime: {result:?}"
6956 );
6957 assert!(!store.exists(&first).await.unwrap());
6958 }
6959
6960 #[tokio::test]
6961 async fn transactional_orphan_sweep_still_removes_orphans_older_than_the_grace_period() {
6962 let dir = tempfile::tempdir().unwrap();
6967 let db_path = dir.path().join("khive.db");
6968 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
6969 {
6970 let mut writer = backend.pool().writer().unwrap();
6971 prepare_completed_v21_gc_fixture(writer.conn_mut());
6972 }
6973 let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6974 .unwrap()
6975 .with_orphan_sweep_grace(Duration::from_secs(60));
6976
6977 let orphan = store
6978 .put(b"actually orphaned, published long ago".to_vec())
6979 .await
6980 .unwrap();
6981 let path = shard_path(store.root(), &orphan);
6984 let old_mtime = SystemTime::now() - Duration::from_secs(3600);
6985 fs::OpenOptions::new()
6986 .write(true)
6987 .open(&path)
6988 .unwrap()
6989 .set_modified(old_mtime)
6990 .unwrap();
6991
6992 let result = store
6993 .transactional_orphan_sweep(backend.sql().as_ref(), false)
6994 .await
6995 .unwrap();
6996
6997 assert_eq!(
6998 result.deleted, 1,
6999 "an orphan older than the grace period must still be swept: {result:?}"
7000 );
7001 assert_eq!(result.grace_period_skipped, 0);
7002 assert!(!store.exists(&orphan).await.unwrap());
7003 }
7004
7005 #[test]
7006 fn resolve_blob_root_prefers_env_var() {
7007 let _guard = ENV_LOCK.lock().unwrap();
7008 std::env::set_var("KHIVE_BLOB_ROOT", "/tmp/env-override-root");
7009 let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
7010 std::env::remove_var("KHIVE_BLOB_ROOT");
7011 assert_eq!(resolved.unwrap(), PathBuf::from("/tmp/env-override-root"));
7012 }
7013
7014 #[test]
7015 fn resolve_blob_root_prefers_config_over_default() {
7016 let _guard = ENV_LOCK.lock().unwrap();
7017 std::env::remove_var("KHIVE_BLOB_ROOT");
7018 let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
7019 assert_eq!(resolved.unwrap(), PathBuf::from("/cfg/root"));
7020 }
7021
7022 #[test]
7023 fn resolve_blob_root_defaults_beside_db_dir() {
7024 let _guard = ENV_LOCK.lock().unwrap();
7025 std::env::remove_var("KHIVE_BLOB_ROOT");
7026 let resolved = resolve_blob_root(Some(Path::new("/db/dir")), None);
7027 assert_eq!(resolved.unwrap(), PathBuf::from("/db/dir/blobs"));
7028 }
7029
7030 #[test]
7031 fn resolve_blob_root_errors_with_no_env_config_or_db_dir() {
7032 let _guard = ENV_LOCK.lock().unwrap();
7033 std::env::remove_var("KHIVE_BLOB_ROOT");
7034 let resolved = resolve_blob_root(None, None);
7035 assert!(resolved.is_err());
7036 }
7037
7038 static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
7042
7043 #[cfg(unix)]
7044 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7045 async fn transactional_orphan_sweep_walk_ignores_a_leaf_swapped_for_an_outside_symlink_mid_scan(
7046 ) {
7047 let dir = tempfile::tempdir().unwrap();
7063 let db_path = dir.path().join("khive.db");
7064 let backend =
7065 std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
7066 {
7067 let mut writer = backend.pool().writer().unwrap();
7068 prepare_completed_v21_gc_fixture(writer.conn_mut());
7069 }
7070
7071 let root = dir.path().join("blobs");
7072 let store = std::sync::Arc::new(
7073 FsBlobStore::new(root.clone(), 0)
7074 .unwrap()
7075 .with_orphan_sweep_grace(Duration::from_secs(3600)),
7076 );
7077
7078 let real = store
7082 .put(b"real freshly-published blob".to_vec())
7083 .await
7084 .unwrap();
7085 let real_path = shard_path(&root, &real);
7086
7087 let outside_dir = dir.path().join("outside");
7092 fs::create_dir_all(&outside_dir).unwrap();
7093 let decoy_path = outside_dir.join(real.as_str());
7094 fs::write(&decoy_path, b"outside decoy, must never be observed").unwrap();
7095 let ancient = SystemTime::now() - Duration::from_secs(7200);
7096 fs::OpenOptions::new()
7097 .write(true)
7098 .open(&decoy_path)
7099 .unwrap()
7100 .set_modified(ancient)
7101 .unwrap();
7102
7103 let (reached, release) = walk_leaf_sync_hook::install(&root);
7104 let sweep = {
7105 let store = store.clone();
7106 let sql = backend.sql();
7107 tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), true).await })
7108 };
7109 assert!(
7110 recv_blocking(reached).await,
7111 "sweep walk must reach the leaf classification pause"
7112 );
7113
7114 fs::remove_file(&real_path).unwrap();
7117 std::os::unix::fs::symlink(&decoy_path, &real_path).unwrap();
7118
7119 release.send(()).unwrap();
7120 let result = sweep.await.unwrap().unwrap();
7121
7122 assert_eq!(
7123 result.would_delete, 0,
7124 "an outside decoy's stale mtime must never make an in-root candidate \
7125 eligible for deletion: {result:?}"
7126 );
7127 assert_eq!(
7128 result.grace_period_skipped, 0,
7129 "the swapped leaf is a symlink; `openat(..., O_NOFOLLOW)` refuses it, so it \
7130 must be dropped from candidates entirely rather than counted (real or \
7131 outside) at all: {result:?}"
7132 );
7133 assert_eq!(
7134 result.scanned, 0,
7135 "the symlinked leaf must never be scanned as a candidate: {result:?}"
7136 );
7137
7138 fs::remove_file(&real_path).unwrap();
7142 fs::write(&real_path, b"real freshly-published blob").unwrap();
7143 let control = store
7144 .transactional_orphan_sweep(backend.sql().as_ref(), true)
7145 .await
7146 .unwrap();
7147 assert_eq!(
7148 control.grace_period_skipped, 1,
7149 "control: the un-replaced root must still classify the real candidate as \
7150 grace-protected: {control:?}"
7151 );
7152 assert_eq!(control.would_delete, 0);
7153 }
7154
7155 #[cfg(unix)]
7156 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7157 async fn transactional_orphan_sweep_walk_still_finds_a_real_orphan_past_its_grace_period() {
7158 let dir = tempfile::tempdir().unwrap();
7164 let db_path = dir.path().join("khive.db");
7165 let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
7166 {
7167 let mut writer = backend.pool().writer().unwrap();
7168 prepare_completed_v21_gc_fixture(writer.conn_mut());
7169 }
7170 let root = dir.path().join("blobs");
7171 let store = FsBlobStore::new(root.clone(), 0)
7172 .unwrap()
7173 .with_orphan_sweep_grace(Duration::from_secs(60));
7174
7175 let orphan = store.put(b"aged real orphan".to_vec()).await.unwrap();
7176 let path = shard_path(&root, &orphan);
7177 let ancient = SystemTime::now() - Duration::from_secs(3600);
7178 fs::OpenOptions::new()
7179 .write(true)
7180 .open(&path)
7181 .unwrap()
7182 .set_modified(ancient)
7183 .unwrap();
7184
7185 let result = store
7186 .transactional_orphan_sweep(backend.sql().as_ref(), false)
7187 .await
7188 .unwrap();
7189 assert_eq!(
7190 result.deleted, 1,
7191 "a real orphan older than the grace period must still be swept: {result:?}"
7192 );
7193 assert!(!store.exists(&orphan).await.unwrap());
7194 }
7195}