use std::collections::HashMap;
use std::fs;
use std::io::{Read, Write};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex as StdMutex, OnceLock};
use std::time::{Duration, SystemTime};
use async_trait::async_trait;
use khive_storage::blob::{
BlobOrphanSweepConfig, BlobOrphanSweepResult, BlobStore, ContentRef, UploadId,
MAX_BLOB_WHOLE_BYTES,
};
use khive_storage::error::StorageError;
use khive_storage::types::{SqlRow, SqlStatement, SqlValue, StorageResult};
use khive_storage::{AtomicUnitOp, SqlAccess, StorageCapability};
use crate::error::SqliteError;
use uuid::Uuid;
#[path = "blob_uploads.rs"]
mod uploads;
const ROOT_WRITE_LOCK_FILE: &str = ".khive-blob-write.lock";
const DATABASE_GC_LOCK_SUFFIX: &str = ".khive-blob-gc.lock";
const BLOB_GC_CLAIM_BATCH_SIZE: usize = 128;
fn map_io_err(e: std::io::Error, op: &'static str) -> StorageError {
StorageError::driver(StorageCapability::Blob, op, e)
}
#[cfg(any(test, not(unix)))]
fn shard_path(root: &Path, content_ref: &ContentRef) -> PathBuf {
let hex = content_ref.as_str();
root.join(&hex[0..2]).join(&hex[2..4]).join(hex)
}
#[cfg(unix)]
fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
open_dir_no_follow(root)
}
#[cfg(windows)]
fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
use std::fs::OpenOptions;
use std::os::windows::fs::OpenOptionsExt;
const FILE_SHARE_READ: u32 = 0x1;
const FILE_SHARE_WRITE: u32 = 0x2;
const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
let handle = OpenOptions::new()
.read(true)
.share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
.custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
.open(root)?;
let file_type = handle.metadata()?.file_type();
if file_type.is_symlink() || !file_type.is_dir() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"blob store root is not a directory or is a reparse point: {}",
root.display()
),
));
}
Ok(handle)
}
#[cfg(not(any(unix, windows)))]
fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
let handle = std::fs::File::open(root)?;
if !handle.metadata()?.is_dir() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!("blob store root is not a directory: {}", root.display()),
));
}
Ok(handle)
}
#[cfg(unix)]
fn verify_blob_root_identity(root: &Path, root_handle: &std::fs::File) -> std::io::Result<()> {
use std::os::unix::fs::MetadataExt;
let current = open_dir_no_follow(root).map_err(|error| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"blob store root is no longer reachable as its initialization-time directory ({}): {error}",
root.display()
),
)
})?;
let expected = root_handle.metadata()?;
let current = current.metadata()?;
if expected.dev() != current.dev() || expected.ino() != current.ino() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"blob store root no longer names its initialization-time directory: {}",
root.display()
),
));
}
Ok(())
}
#[cfg(not(unix))]
fn verify_blob_root_identity(root: &Path, _root_handle: &std::fs::File) -> std::io::Result<()> {
let current = root.canonicalize().map_err(|error| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"blob store root is no longer reachable as its initialization-time directory ({}): {error}",
root.display()
),
)
})?;
if current != root {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"blob store root no longer names its initialization-time directory: {}",
root.display()
),
));
}
Ok(())
}
#[cfg(unix)]
fn unlink_blob_shard_file_no_follow(
root: &Path,
root_handle: &std::fs::File,
content_ref: &ContentRef,
) -> std::io::Result<()> {
use std::os::unix::io::AsRawFd;
verify_blob_root_identity(root, root_handle)?;
let hex = content_ref.as_str();
let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
unlink_entry_at(shard2_dir.as_raw_fd(), hex)
}
#[cfg(windows)]
fn unlink_blob_shard_file_no_follow(
root: &Path,
root_handle: &std::fs::File,
content_ref: &ContentRef,
) -> std::io::Result<()> {
use std::fs::OpenOptions;
use std::os::windows::ffi::OsStringExt;
use std::os::windows::fs::OpenOptionsExt;
use std::os::windows::io::AsRawHandle;
use windows_sys::Win32::Storage::FileSystem::{
FileDispositionInfo, GetFinalPathNameByHandleW, SetFileInformationByHandle,
FILE_DISPOSITION_INFO,
};
const FILE_SHARE_READ: u32 = 0x1;
const FILE_SHARE_WRITE: u32 = 0x2;
const DELETE: u32 = 0x0001_0000;
const FILE_READ_ATTRIBUTES: u32 = 0x80;
const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
const FINAL_PATH_FLAGS: u32 = 0x0;
fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
let dir = OpenOptions::new()
.read(true)
.share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
.custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
.open(path)?;
let file_type = dir.metadata()?.file_type();
if file_type.is_symlink() || !file_type.is_dir() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"refusing to unlink blob shard file through non-directory or \
reparse-point path component: {}",
path.display()
),
));
}
Ok(dir)
}
fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<std::path::PathBuf> {
let handle = file.as_raw_handle();
let mut buf: Vec<u16> = vec![0; 512];
loop {
let len = unsafe {
GetFinalPathNameByHandleW(
handle as _,
buf.as_mut_ptr(),
buf.len() as u32,
FINAL_PATH_FLAGS,
)
};
if len == 0 {
return Err(std::io::Error::last_os_error());
}
let len = len as usize;
if len <= buf.len() {
buf.truncate(len);
return Ok(std::path::PathBuf::from(std::ffi::OsString::from_wide(
&buf,
)));
}
buf.resize(len, 0);
}
}
verify_blob_root_identity(root, root_handle)?;
let hex = content_ref.as_str();
let shard1 = root.join(&hex[0..2]);
let shard2 = shard1.join(&hex[2..4]);
let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
let expected = final_path_by_handle(root_handle)?
.join(&hex[0..2])
.join(&hex[2..4])
.join(hex);
let target = OpenOptions::new()
.access_mode(DELETE | FILE_READ_ATTRIBUTES)
.share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
.custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
.open(shard2.join(hex))?;
let resolved = final_path_by_handle(&target)?;
if resolved != expected {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"refusing blob delete: handle resolved outside the verified blob root \
(expected {}, resolved {})",
expected.display(),
resolved.display()
),
));
}
let disposition = FILE_DISPOSITION_INFO { DeleteFile: true };
let ok = unsafe {
SetFileInformationByHandle(
target.as_raw_handle() as _,
FileDispositionInfo,
std::ptr::from_ref(&disposition).cast(),
std::mem::size_of::<FILE_DISPOSITION_INFO>() as u32,
)
};
if ok == 0 {
return Err(std::io::Error::last_os_error());
}
Ok(())
}
#[cfg(not(any(unix, windows)))]
fn unlink_blob_shard_file_no_follow(
root: &Path,
root_handle: &std::fs::File,
content_ref: &ContentRef,
) -> std::io::Result<()> {
verify_blob_root_identity(root, root_handle)?;
let hex = content_ref.as_str();
let shard1 = root.join(&hex[0..2]);
let shard2 = shard1.join(&hex[2..4]);
for component in [root, shard1.as_path(), shard2.as_path()] {
let metadata = fs::symlink_metadata(component)?;
if metadata.file_type().is_symlink() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"refusing to unlink blob shard file through symlinked path component: {}",
component.display()
),
));
}
}
fs::remove_file(shard2.join(hex))
}
#[cfg(unix)]
fn open_blob_shard_file_no_follow(
root: &Path,
root_handle: &std::fs::File,
content_ref: &ContentRef,
) -> std::io::Result<std::fs::File> {
verify_blob_root_identity(root, root_handle)?;
open_blob_shard_file_at_no_follow(root_handle, content_ref, libc::O_RDONLY)
}
#[cfg(windows)]
fn open_blob_shard_file_no_follow(
root: &Path,
root_handle: &std::fs::File,
content_ref: &ContentRef,
) -> std::io::Result<std::fs::File> {
use std::fs::OpenOptions;
use std::os::windows::ffi::OsStringExt;
use std::os::windows::fs::OpenOptionsExt;
use std::os::windows::io::AsRawHandle;
use windows_sys::Win32::Storage::FileSystem::GetFinalPathNameByHandleW;
const FILE_SHARE_READ: u32 = 0x1;
const FILE_SHARE_WRITE: u32 = 0x2;
const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
const FINAL_PATH_FLAGS: u32 = 0x0;
fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
let dir = OpenOptions::new()
.read(true)
.share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
.custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
.open(path)?;
let file_type = dir.metadata()?.file_type();
if file_type.is_symlink() || !file_type.is_dir() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"refusing to read a blob through non-directory or reparse-point component: {}",
path.display()
),
));
}
Ok(dir)
}
fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<PathBuf> {
let mut buf: Vec<u16> = vec![0; 512];
loop {
let len = unsafe {
GetFinalPathNameByHandleW(
file.as_raw_handle() as _,
buf.as_mut_ptr(),
buf.len() as u32,
FINAL_PATH_FLAGS,
)
};
if len == 0 {
return Err(std::io::Error::last_os_error());
}
let len = len as usize;
if len <= buf.len() {
buf.truncate(len);
return Ok(PathBuf::from(std::ffi::OsString::from_wide(&buf)));
}
buf.resize(len, 0);
}
}
verify_blob_root_identity(root, root_handle)?;
let hex = content_ref.as_str();
let shard1 = root.join(&hex[0..2]);
let shard2 = shard1.join(&hex[2..4]);
let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
let expected = final_path_by_handle(root_handle)?
.join(&hex[0..2])
.join(&hex[2..4])
.join(hex);
let target = OpenOptions::new()
.read(true)
.share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
.custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
.open(shard2.join(hex))?;
let file_type = target.metadata()?.file_type();
if file_type.is_symlink() || !file_type.is_file() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"refusing to read a blob leaf that is not a regular file",
));
}
let resolved = final_path_by_handle(&target)?;
if resolved != expected {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!(
"refusing blob read: handle resolved outside the verified blob root (expected {}, resolved {})",
expected.display(),
resolved.display()
),
));
}
Ok(target)
}
#[cfg(not(any(unix, windows)))]
fn open_blob_shard_file_no_follow(
root: &Path,
root_handle: &std::fs::File,
_content_ref: &ContentRef,
) -> std::io::Result<std::fs::File> {
verify_blob_root_identity(root, root_handle)?;
Err(std::io::Error::new(
std::io::ErrorKind::Unsupported,
"bounded verified blob reads require handle-relative no-follow file APIs",
))
}
#[cfg(unix)]
fn open_dir_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
use std::os::unix::ffi::OsStrExt;
use std::os::unix::io::FromRawFd;
let c_path = std::ffi::CString::new(path.as_os_str().as_bytes())
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let fd = unsafe {
libc::open(
c_path.as_ptr(),
libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
)
};
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(unsafe { std::fs::File::from_raw_fd(fd) })
}
#[cfg(unix)]
fn openat_dir_no_follow(
parent_fd: std::os::unix::io::RawFd,
name: &str,
) -> std::io::Result<std::fs::File> {
use std::os::unix::io::FromRawFd;
let c_name = std::ffi::CString::new(name)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let fd = unsafe {
libc::openat(
parent_fd,
c_name.as_ptr(),
libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
)
};
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(unsafe { std::fs::File::from_raw_fd(fd) })
}
#[cfg(unix)]
fn openat_regular_file_no_follow(
parent_fd: std::os::unix::io::RawFd,
name: &str,
access_flags: libc::c_int,
) -> std::io::Result<std::fs::File> {
use std::os::unix::io::FromRawFd;
let c_name = std::ffi::CString::new(name)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let fd = unsafe {
libc::openat(
parent_fd,
c_name.as_ptr(),
access_flags | libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK,
)
};
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
let file = unsafe { std::fs::File::from_raw_fd(fd) };
if !file.metadata()?.file_type().is_file() {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
format!("refusing non-regular blob store entry: {name}"),
));
}
Ok(file)
}
#[cfg(unix)]
fn open_blob_shard_file_at_no_follow(
root_handle: &std::fs::File,
content_ref: &ContentRef,
access_flags: libc::c_int,
) -> std::io::Result<std::fs::File> {
use std::os::unix::io::AsRawFd;
let hex = content_ref.as_str();
let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
openat_regular_file_no_follow(shard2_dir.as_raw_fd(), hex, access_flags)
}
#[cfg(unix)]
fn open_or_create_dir_at_no_follow(
parent_fd: std::os::unix::io::RawFd,
name: &str,
) -> std::io::Result<std::fs::File> {
let c_name = std::ffi::CString::new(name)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
match openat_dir_no_follow(parent_fd, name) {
Ok(dir) => Ok(dir),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
let rc = unsafe { libc::mkdirat(parent_fd, c_name.as_ptr(), 0o777) };
if rc != 0 {
let error = std::io::Error::last_os_error();
if error.kind() != std::io::ErrorKind::AlreadyExists {
return Err(error);
}
}
openat_dir_no_follow(parent_fd, name)
}
Err(error) => Err(error),
}
}
#[cfg(unix)]
fn create_regular_file_at_no_follow(
parent_fd: std::os::unix::io::RawFd,
name: &str,
mode: libc::mode_t,
) -> std::io::Result<std::fs::File> {
use std::os::unix::io::FromRawFd;
let c_name = std::ffi::CString::new(name)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let fd = unsafe {
libc::openat(
parent_fd,
c_name.as_ptr(),
libc::O_RDWR | libc::O_CREAT | libc::O_EXCL | libc::O_NOFOLLOW | libc::O_CLOEXEC,
mode as libc::c_uint,
)
};
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(unsafe { std::fs::File::from_raw_fd(fd) })
}
#[cfg(unix)]
fn unlink_entry_at(parent_fd: std::os::unix::io::RawFd, name: &str) -> std::io::Result<()> {
let c_name = std::ffi::CString::new(name)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let rc = unsafe { libc::unlinkat(parent_fd, c_name.as_ptr(), 0) };
if rc != 0 {
return Err(std::io::Error::last_os_error());
}
Ok(())
}
#[cfg(unix)]
fn rename_entry_at(
source_fd: std::os::unix::io::RawFd,
from: &str,
destination_fd: std::os::unix::io::RawFd,
to: &str,
) -> std::io::Result<()> {
let c_from = std::ffi::CString::new(from)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let c_to = std::ffi::CString::new(to)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
let rc = unsafe { libc::renameat(source_fd, c_from.as_ptr(), destination_fd, c_to.as_ptr()) };
if rc != 0 {
return Err(std::io::Error::last_os_error());
}
Ok(())
}
#[cfg(unix)]
fn available_space_at(root_handle: &std::fs::File) -> std::io::Result<u64> {
use std::os::unix::io::AsRawFd;
let mut stat: libc::statvfs = unsafe { std::mem::zeroed() };
let rc = unsafe { libc::fstatvfs(root_handle.as_raw_fd(), &mut stat) };
if rc != 0 {
return Err(std::io::Error::last_os_error());
}
#[allow(clippy::useless_conversion)]
Ok(stat.f_frsize.saturating_mul(u64::from(stat.f_bavail)))
}
#[cfg(unix)]
fn acquire_root_write_lock_at(root_handle: &std::fs::File) -> StorageResult<std::fs::File> {
use std::os::unix::io::{AsRawFd, FromRawFd};
let c_name = std::ffi::CString::new(ROOT_WRITE_LOCK_FILE)
.expect("the static blob root lock name contains no NUL");
let fd = unsafe {
libc::openat(
root_handle.as_raw_fd(),
c_name.as_ptr(),
libc::O_RDWR | libc::O_CREAT | libc::O_NOFOLLOW | libc::O_CLOEXEC,
0o666,
)
};
if fd < 0 {
return Err(map_io_err(
std::io::Error::last_os_error(),
"root_write_lock_open",
));
}
let lock_file = unsafe { std::fs::File::from_raw_fd(fd) };
if !lock_file
.metadata()
.map_err(|e| map_io_err(e, "root_write_lock_metadata"))?
.file_type()
.is_file()
{
return Err(map_io_err(
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"blob root write lock is not a regular file",
),
"root_write_lock_open",
));
}
fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
Ok(lock_file)
}
fn blob_root_key(root: &Path) -> String {
#[cfg(unix)]
let bytes = {
use std::os::unix::ffi::OsStrExt;
root.as_os_str().as_bytes().to_vec()
};
#[cfg(windows)]
let bytes = {
use std::os::windows::ffi::OsStrExt;
root.as_os_str()
.encode_wide()
.flat_map(u16::to_le_bytes)
.collect::<Vec<_>>()
};
#[cfg(not(any(unix, windows)))]
let bytes = root.to_string_lossy().as_bytes().to_vec();
blake3::hash(&bytes).to_hex().to_string()
}
pub fn resolve_blob_root(
db_dir: Option<&Path>,
config_root: Option<&Path>,
) -> Result<PathBuf, SqliteError> {
if let Ok(env_root) = std::env::var("KHIVE_BLOB_ROOT") {
if !env_root.trim().is_empty() {
return Ok(PathBuf::from(env_root));
}
}
if let Some(root) = config_root {
return Ok(root.to_path_buf());
}
if let Some(dir) = db_dir {
return Ok(dir.join("blobs"));
}
Err(SqliteError::InvalidData(
"cannot resolve a blob store root: no KHIVE_BLOB_ROOT env var, no configured \
root, and the database has no on-disk directory to default beside (in-memory \
backend)"
.to_string(),
))
}
fn crosses_floor(available: u64, required_write_bytes: u64, floor_bytes: u64) -> bool {
available.saturating_sub(required_write_bytes) < floor_bytes
}
#[cfg(any(test, not(unix)))]
fn put_blocking_with_space_probe<F>(
root: &Path,
floor_bytes: u64,
bytes: Vec<u8>,
available_space: F,
) -> StorageResult<ContentRef>
where
F: FnOnce(&Path) -> std::io::Result<u64>,
{
let digest = blake3::hash(&bytes);
let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
let target = shard_path(root, &content_ref);
if target.exists() {
let file = fs::OpenOptions::new()
.write(true)
.open(&target)
.map_err(|e| map_io_err(e, "put_touch_open"))?;
file.set_modified(SystemTime::now())
.map_err(|e| map_io_err(e, "put_touch_mtime"))?;
return Ok(content_ref);
}
let required_write_bytes = bytes.len() as u64;
let available = available_space(root).map_err(|e| map_io_err(e, "put_check_space"))?;
if crosses_floor(available, required_write_bytes, floor_bytes) {
return Err(StorageError::CapacityFloor {
capability: StorageCapability::Blob,
volume: root.display().to_string(),
available_bytes: available,
floor_bytes,
});
}
let shard_dir = target
.parent()
.expect("shard_path always nests under two directory levels");
fs::create_dir_all(shard_dir).map_err(|e| map_io_err(e, "put_mkdir"))?;
let mut tmp = tempfile::Builder::new()
.prefix(".tmp-")
.tempfile_in(shard_dir)
.map_err(|e| map_io_err(e, "put_tempfile"))?;
tmp.write_all(&bytes)
.map_err(|e| map_io_err(e, "put_write"))?;
tmp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
tmp.as_file()
.sync_all()
.map_err(|e| map_io_err(e, "put_fsync"))?;
let written_len = tmp
.as_file()
.metadata()
.map_err(|e| map_io_err(e, "put_verify"))?
.len();
if written_len != bytes.len() as u64 {
return Err(map_io_err(
std::io::Error::other(format!(
"temp file length {written_len} does not match {} written bytes",
bytes.len()
)),
"put_verify",
));
}
let temporary = tmp.into_temp_path();
publish_blob_path(&temporary, &target)?;
Ok(content_ref)
}
#[cfg(any(test, not(unix)))]
fn publish_blob_path(source: &Path, target: &Path) -> StorageResult<()> {
fs::rename(source, target).map_err(|error| map_io_err(error, "put_persist"))
}
#[cfg(any(test, not(unix)))]
fn put_blocking(root: &Path, floor_bytes: u64, bytes: Vec<u8>) -> StorageResult<ContentRef> {
let _root_write_guard = acquire_root_write_lock(root)?;
put_blocking_with_space_probe(root, floor_bytes, bytes, |path| fs4::available_space(path))
}
fn acquire_root_write_lock_anchored(
root: &Path,
root_handle: &std::fs::File,
) -> StorageResult<std::fs::File> {
verify_blob_root_identity(root, root_handle)
.map_err(|e| map_io_err(e, "root_write_lock_identity"))?;
#[cfg(unix)]
{
acquire_root_write_lock_at(root_handle)
}
#[cfg(not(unix))]
{
acquire_root_write_lock(root)
}
}
#[cfg(unix)]
struct BlobPublication {
#[cfg(test)]
hook: Option<sync_hook::Publication>,
}
#[cfg(unix)]
impl BlobPublication {
fn step<T>(
&self,
operation: &'static str,
action: impl FnOnce() -> std::io::Result<T>,
) -> StorageResult<T> {
self.io_step(operation, action)
.map_err(|error| map_io_err(error, operation))
}
fn io_step<T>(
&self,
_operation: &'static str,
action: impl FnOnce() -> std::io::Result<T>,
) -> std::io::Result<T> {
#[cfg(test)]
if let Some(hook) = &self.hook {
hook.before(_operation)?;
}
let result = action()?;
#[cfg(test)]
if let Some(hook) = &self.hook {
hook.completed(_operation);
}
Ok(result)
}
fn sync_directories(
&self,
root: &fs::File,
shard1: &fs::File,
shard2: &fs::File,
) -> StorageResult<()> {
for (operation, directory) in [
("put_sync_shard", shard2),
("put_sync_parent", shard1),
("put_sync_root", root),
] {
self.sync_directory(operation, directory)
.map_err(|error| map_io_err(error, operation))?;
}
Ok(())
}
fn sync_directory(&self, operation: &'static str, directory: &fs::File) -> std::io::Result<()> {
self.io_step(operation, || {
sync_directory(directory)?;
#[cfg(test)]
if let Some(hook) = &self.hook {
hook.directory_synced(operation, directory)?;
}
Ok(())
})
}
}
#[cfg(unix)]
fn sync_directory(directory: &fs::File) -> std::io::Result<()> {
use std::os::fd::AsRawFd;
loop {
if unsafe { libc::fsync(directory.as_raw_fd()) } == 0 {
return Ok(());
}
let error = std::io::Error::last_os_error();
if error.kind() != std::io::ErrorKind::Interrupted {
return Err(error);
}
}
}
#[cfg(unix)]
fn publish_blob_at(
root: &fs::File,
shard1: &fs::File,
shard2: &fs::File,
source: &fs::File,
temp_name: &str,
content_ref: &ContentRef,
publication: &BlobPublication,
) -> StorageResult<()> {
use std::os::fd::AsRawFd;
if let Err(error) = publication.step("put_persist", || {
rename_entry_at(
source.as_raw_fd(),
temp_name,
shard2.as_raw_fd(),
content_ref.as_str(),
)
}) {
let _ = unlink_entry_at(source.as_raw_fd(), temp_name);
return Err(error);
}
publication.sync_directories(root, shard1, shard2)
}
#[cfg(unix)]
fn put_blocking_from_root_handle(
root: &Path,
root_handle: &std::fs::File,
floor_bytes: u64,
bytes: Vec<u8>,
publication: &BlobPublication,
) -> StorageResult<ContentRef> {
use std::os::unix::io::AsRawFd;
let _root_write_guard = acquire_root_write_lock_anchored(root, root_handle)?;
let digest = blake3::hash(&bytes);
let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
let hex = content_ref.as_str();
let existing = (|| -> std::io::Result<_> {
let shard1 = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
let shard2 = openat_dir_no_follow(shard1.as_raw_fd(), &hex[2..4])?;
let file = openat_regular_file_no_follow(shard2.as_raw_fd(), hex, libc::O_WRONLY)?;
Ok((file, shard1, shard2))
})();
match existing {
Ok((file, shard1, shard2)) => {
file.set_modified(SystemTime::now())
.map_err(|e| map_io_err(e, "put_touch_mtime"))?;
publication.step("put_fsync", || file.sync_all())?;
publication.sync_directories(root_handle, &shard1, &shard2)?;
return Ok(content_ref);
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(map_io_err(error, "put_touch_open")),
}
let required_write_bytes = bytes.len() as u64;
let available =
available_space_at(root_handle).map_err(|e| map_io_err(e, "put_check_space"))?;
if crosses_floor(available, required_write_bytes, floor_bytes) {
return Err(StorageError::CapacityFloor {
capability: StorageCapability::Blob,
volume: root.display().to_string(),
available_bytes: available,
floor_bytes,
});
}
let shard1_dir = open_or_create_dir_at_no_follow(root_handle.as_raw_fd(), &hex[0..2])
.map_err(|e| map_io_err(e, "put_mkdir"))?;
let shard2_dir = open_or_create_dir_at_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])
.map_err(|e| map_io_err(e, "put_mkdir"))?;
let temp_name = format!(".tmp-{}", Uuid::new_v4());
let mut temp = create_regular_file_at_no_follow(shard2_dir.as_raw_fd(), &temp_name, 0o600)
.map_err(|e| map_io_err(e, "put_tempfile"))?;
let write_result = (|| -> StorageResult<()> {
temp.write_all(&bytes)
.map_err(|e| map_io_err(e, "put_write"))?;
temp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
publication.step("put_fsync", || temp.sync_all())?;
let written_len = temp
.metadata()
.map_err(|e| map_io_err(e, "put_verify"))?
.len();
if written_len != bytes.len() as u64 {
return Err(map_io_err(
std::io::Error::other(format!(
"temp file length {written_len} does not match {} written bytes",
bytes.len()
)),
"put_verify",
));
}
Ok(())
})();
drop(temp);
if let Err(error) = write_result {
let _ = unlink_entry_at(shard2_dir.as_raw_fd(), &temp_name);
return Err(error);
}
publish_blob_at(
root_handle,
&shard1_dir,
&shard2_dir,
&shard2_dir,
&temp_name,
&content_ref,
publication,
)?;
Ok(content_ref)
}
#[cfg(not(unix))]
fn put_blocking_from_root_handle(
root: &Path,
root_handle: &std::fs::File,
floor_bytes: u64,
bytes: Vec<u8>,
) -> StorageResult<ContentRef> {
verify_blob_root_identity(root, root_handle).map_err(|e| map_io_err(e, "put_root_identity"))?;
put_blocking(root, floor_bytes, bytes)
}
#[cfg(any(test, not(unix)))]
fn acquire_root_write_lock(root: &Path) -> StorageResult<fs::File> {
let lock_file = fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(root.join(ROOT_WRITE_LOCK_FILE))
.map_err(|e| map_io_err(e, "root_write_lock_open"))?;
fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
Ok(lock_file)
}
fn database_gc_lock_path(database_path: &Path) -> PathBuf {
let mut lock_path = database_path.as_os_str().to_os_string();
lock_path.push(DATABASE_GC_LOCK_SUFFIX);
PathBuf::from(lock_path)
}
fn acquire_database_gc_lock(database_path: Option<&Path>) -> StorageResult<Option<fs::File>> {
let Some(database_path) = database_path else {
return Ok(None);
};
let lock_path = database_gc_lock_path(database_path);
let lock_file = fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&lock_path)
.map_err(|e| map_io_err(e, "database_gc_lock_open"))?;
fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "database_gc_lock_acquire"))?;
Ok(Some(lock_file))
}
#[cfg(all(unix, target_os = "macos"))]
fn errno_location() -> *mut libc::c_int {
unsafe { libc::__error() }
}
#[cfg(all(unix, not(target_os = "macos")))]
fn errno_location() -> *mut libc::c_int {
unsafe { libc::__errno_location() }
}
#[cfg(unix)]
fn clear_errno() {
unsafe { *errno_location() = 0 };
}
#[cfg(unix)]
fn current_errno() -> libc::c_int {
unsafe { *errno_location() }
}
#[cfg(unix)]
fn read_dir_names_no_follow(dir_fd: std::os::unix::io::RawFd) -> std::io::Result<Vec<String>> {
use std::os::unix::io::IntoRawFd;
let reopened = openat_dir_no_follow(dir_fd, ".")?;
let owned_fd = reopened.into_raw_fd();
let dirp = unsafe { libc::fdopendir(owned_fd) };
if dirp.is_null() {
let err = std::io::Error::last_os_error();
unsafe { libc::close(owned_fd) };
return Err(err);
}
let mut names = Vec::new();
loop {
clear_errno();
let entry = unsafe { libc::readdir(dirp) };
if entry.is_null() {
if current_errno() != 0 {
let err = std::io::Error::last_os_error();
unsafe { libc::closedir(dirp) };
return Err(err);
}
break;
}
let first = unsafe { *(*entry).d_name.as_ptr() };
if first == b'.' as libc::c_char {
continue;
}
let name = unsafe { std::ffi::CStr::from_ptr((*entry).d_name.as_ptr()) }
.to_string_lossy()
.into_owned();
names.push(name);
}
unsafe { libc::closedir(dirp) };
Ok(names)
}
#[cfg(unix)]
fn walk_blob_files_from_root_handle(
root_handle: &std::fs::File,
_root: &Path,
) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
use std::os::unix::io::AsRawFd;
let mut out = Vec::new();
for l1_name in read_dir_names_no_follow(root_handle.as_raw_fd())? {
let l1_dir = match openat_dir_no_follow(root_handle.as_raw_fd(), &l1_name) {
Ok(dir) => dir,
Err(_) => continue,
};
for l2_name in read_dir_names_no_follow(l1_dir.as_raw_fd())? {
let l2_dir = match openat_dir_no_follow(l1_dir.as_raw_fd(), &l2_name) {
Ok(dir) => dir,
Err(_) => continue,
};
for leaf_name in read_dir_names_no_follow(l2_dir.as_raw_fd())? {
let Ok(content_ref) = ContentRef::from_hex(leaf_name.clone()) else {
continue;
};
#[cfg(all(test, unix))]
if let Some(hook) = walk_leaf_sync_hook::take(_root) {
let _ = hook.reached.send(());
let _ = hook.release.recv();
}
let file = match openat_regular_file_no_follow(
l2_dir.as_raw_fd(),
&leaf_name,
libc::O_RDONLY,
) {
Ok(file) => file,
Err(_) => continue,
};
let mtime = file.metadata().ok().and_then(|meta| meta.modified().ok());
out.push((content_ref, mtime));
}
}
}
Ok(out)
}
#[cfg(not(unix))]
fn walk_blob_files_from_root_handle(
_root_handle: &std::fs::File,
_root: &Path,
) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
Err(std::io::Error::new(
std::io::ErrorKind::Unsupported,
"orphan sweep candidate enumeration requires descriptor-relative directory \
reads, available only on unix in this release; refusing to classify via \
path-based reads",
))
}
fn within_publish_grace(
mtime: Option<SystemTime>,
now: SystemTime,
grace_period: Duration,
) -> bool {
let age = mtime.and_then(|mtime| now.duration_since(mtime).ok());
match age {
Some(age) => age < grace_period,
None => true,
}
}
#[derive(Debug)]
struct PreparedTransactionalSweep {
result: BlobOrphanSweepResult,
candidates: Vec<(ContentRef, bool)>,
}
fn prepare_transactional_sweep(
files: Vec<(ContentRef, Option<SystemTime>)>,
grace_period: Duration,
) -> PreparedTransactionalSweep {
let now = SystemTime::now();
let mut result = BlobOrphanSweepResult::default();
let mut candidates = Vec::with_capacity(files.len());
for (content_ref, mtime) in files {
result.scanned += 1;
let within_grace = within_publish_grace(mtime, now, grace_period);
candidates.push((content_ref, within_grace));
}
PreparedTransactionalSweep { result, candidates }
}
#[derive(Debug)]
struct BlobGcBatchRows {
grace_period_skipped: u64,
would_delete: u64,
claimed_rows: Vec<SqlRow>,
}
fn required_nonnegative_count(
value: Option<SqlValue>,
operation: &'static str,
) -> StorageResult<u64> {
match value {
Some(SqlValue::Integer(value)) if value >= 0 => Ok(value as u64),
other => Err(StorageError::Internal(format!(
"{operation} returned an invalid count: {other:?}"
))),
}
}
fn invalid_content_ref(message: String) -> StorageError {
StorageError::InvalidInput {
capability: StorageCapability::Blob,
operation: "transactional_orphan_sweep".into(),
message,
}
}
async fn blob_gc_fencing_complete(sql: &dyn SqlAccess) -> StorageResult<bool> {
let mut reader = sql.reader().await?;
let present = required_nonnegative_count(
reader
.query_scalar(SqlStatement {
sql: "SELECT COUNT(*) FROM sqlite_master \
WHERE (type = 'table' AND name IN ( \
'blob_gc_claims', 'attachments', \
'attachment_cutover_state')) \
OR (type = 'index' AND name IN ( \
'idx_blob_gc_claims_content_ref', \
'idx_attachments_content_ref')) \
OR (type = 'trigger' AND name IN ( \
'attachments_reject_claimed_blob_insert', \
'attachments_reject_claimed_blob_update'))"
.to_string(),
params: vec![],
label: Some("blob_gc_fencing_complete".to_string()),
})
.await?,
"blob_gc_fencing_complete",
)?;
if present != 7 {
return Ok(false);
}
let legacy_objects = required_nonnegative_count(
reader
.query_scalar(SqlStatement {
sql: "SELECT \
(SELECT COUNT(*) FROM pragma_table_info('entities') \
WHERE name = 'content_ref') \
+ (SELECT COUNT(*) FROM sqlite_master \
WHERE (type = 'index' AND name = 'idx_entities_content_ref') \
OR (type = 'trigger' AND name IN ( \
'entities_reject_claimed_blob_insert', \
'entities_reject_claimed_blob_update')))"
.to_string(),
params: vec![],
label: Some("blob_gc_legacy_fencing_absent".to_string()),
})
.await?,
"blob_gc_legacy_fencing_absent",
)?;
if legacy_objects != 0 {
return Ok(false);
}
let complete = required_nonnegative_count(
reader
.query_scalar(SqlStatement {
sql: "SELECT COUNT(*) FROM attachment_cutover_state AS cutover \
WHERE cutover.singleton = 1 \
AND cutover.state = 'complete' \
AND cutover.completed_at IS NOT NULL \
AND (SELECT COUNT(*) FROM _schema_migrations \
WHERE version = ?1 \
AND name = 'attachments_first_class') = 1 \
AND (SELECT COUNT(*) FROM _schema_migrations) = ?1 \
AND (SELECT MIN(version) FROM _schema_migrations) = 1 \
AND (SELECT MAX(version) FROM _schema_migrations) = ?1"
.to_string(),
params: vec![SqlValue::Integer(i64::from(
crate::migrations::ATTACHMENT_CUTOVER_VERSION,
))],
label: Some("blob_gc_cutover_complete".to_string()),
})
.await?,
"blob_gc_cutover_complete",
)?;
Ok(complete == 1)
}
fn unsupported_blob_gc_epoch() -> StorageError {
StorageError::Unsupported {
capability: StorageCapability::Blob,
operation: "transactional_orphan_sweep".into(),
message: "transactional blob GC requires a complete V21 attachment cutover with \
the attachment claim-fencing set; refusing both report-only and \
destructive sweep in this database epoch"
.into(),
}
}
const BLOB_GC_FENCE_PROBE_REF: &str =
"0000000000000000000000000000000000000000000000000000000000000000";
const BLOB_GC_FENCE_TRIGGER_MESSAGE: &str = "content_ref is reserved by an active blob sweep";
async fn blob_gc_fence_probe(sql: &dyn SqlAccess) -> StorageResult<()> {
let run = Uuid::new_v4().simple().to_string();
blob_gc_fence_probe_with_ids(
sql,
format!("__blob-gc-fence-probe-insert-{run}__"),
format!("__blob-gc-fence-probe-update-{run}__"),
format!("__blob-gc-fence-probe-insert2-{run}__"),
format!("__blob-gc-fence-probe-update2-{run}__"),
format!("__fence_probe-{run}__"),
)
.await
}
async fn blob_gc_fence_probe_with_ids(
sql: &dyn SqlAccess,
insert_id: String,
update_id: String,
insert2_id: String,
update2_id: String,
claim_key: String,
) -> StorageResult<()> {
fn fence_rejection(result: Result<u64, StorageError>) -> Result<bool, String> {
match result {
Ok(_) => Ok(false),
Err(error) => {
let text = error.to_string();
if text.contains(BLOB_GC_FENCE_TRIGGER_MESSAGE) {
Ok(true)
} else {
Err(text)
}
}
}
}
fn required_seed(value: Option<SqlValue>) -> StorageResult<String> {
match value {
Some(SqlValue::Text(seed)) => Ok(seed),
_ => Err(StorageError::Unsupported {
capability: StorageCapability::Blob,
operation: "transactional_orphan_sweep".into(),
message: "the blob GC fence probe could not select an unclaimed \
canonical seed; refusing deletion so a later sweep can retry"
.into(),
}),
}
}
let op: AtomicUnitOp = Box::new(move |writer| {
Box::pin(async move {
let preexisting = writer
.query_row(SqlStatement {
sql: "SELECT (SELECT COUNT(*) FROM attachments \
WHERE record_uuid IN (?1, ?2, ?3, ?4)) \
+ (SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?5)"
.to_string(),
params: vec![
SqlValue::Text(insert_id.clone()),
SqlValue::Text(update_id.clone()),
SqlValue::Text(insert2_id.clone()),
SqlValue::Text(update2_id.clone()),
SqlValue::Text(claim_key.clone()),
],
label: Some("blob_gc_fence_probe_ownership_guard".to_string()),
})
.await?
.and_then(|row| row.columns.first().map(|c| c.value.clone()));
match preexisting {
Some(SqlValue::Integer(0)) => {}
Some(SqlValue::Integer(_)) => {
return Err(StorageError::Unsupported {
capability: StorageCapability::Blob,
operation: "transactional_orphan_sweep".into(),
message: "the blob GC fence probe's row ids collide with existing \
rows; refusing to probe rather than delete data the \
probe does not own"
.into(),
});
}
_ => {
return Err(StorageError::Internal(
"blob GC fence probe ownership guard returned no count".into(),
));
}
}
let seed_ref = writer
.query_row(SqlStatement {
sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
SELECT 1, lower(hex(randomblob(32))) \
UNION ALL \
SELECT attempt + 1, lower(hex(randomblob(32))) \
FROM candidates WHERE attempt < 8 \
) \
SELECT candidate.content_ref FROM candidates AS candidate \
WHERE candidate.content_ref <> ?1 \
AND NOT EXISTS ( \
SELECT 1 FROM blob_gc_claims \
WHERE content_ref = candidate.content_ref \
) \
LIMIT 1"
.to_string(),
params: vec![SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string())],
label: Some("blob_gc_fence_probe_select_seed".to_string()),
})
.await?
.and_then(|row| row.columns.first().map(|column| column.value.clone()));
let seed_ref = required_seed(seed_ref)?;
let seed2_ref = writer
.query_row(SqlStatement {
sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
SELECT 1, lower(hex(randomblob(32))) \
UNION ALL \
SELECT attempt + 1, lower(hex(randomblob(32))) \
FROM candidates WHERE attempt < 8 \
) \
SELECT candidate.content_ref FROM candidates AS candidate \
WHERE candidate.content_ref NOT IN (?1, ?2) \
AND NOT EXISTS ( \
SELECT 1 FROM blob_gc_claims \
WHERE content_ref = candidate.content_ref \
) \
LIMIT 1"
.to_string(),
params: vec![
SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
SqlValue::Text(seed_ref.clone()),
],
label: Some("blob_gc_fence_probe_select_seed2".to_string()),
})
.await?
.and_then(|row| row.columns.first().map(|column| column.value.clone()));
let seed2_ref = required_seed(seed2_ref)?;
let probe2_ref = writer
.query_row(SqlStatement {
sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
SELECT 1, lower(hex(randomblob(32))) \
UNION ALL \
SELECT attempt + 1, lower(hex(randomblob(32))) \
FROM candidates WHERE attempt < 8 \
) \
SELECT candidate.content_ref FROM candidates AS candidate \
WHERE candidate.content_ref NOT IN (?1, ?2, ?3) \
AND NOT EXISTS ( \
SELECT 1 FROM blob_gc_claims \
WHERE content_ref = candidate.content_ref \
) \
LIMIT 1"
.to_string(),
params: vec![
SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
SqlValue::Text(seed_ref.clone()),
SqlValue::Text(seed2_ref.clone()),
],
label: Some("blob_gc_fence_probe_select_probe2".to_string()),
})
.await?
.and_then(|row| row.columns.first().map(|column| column.value.clone()));
let probe2_ref = required_seed(probe2_ref)?;
writer
.execute(SqlStatement {
sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES (?1, ?2, 0), (?1, ?3, 0)"
.to_string(),
params: vec![
SqlValue::Text(claim_key.clone()),
SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
SqlValue::Text(probe2_ref.clone()),
],
label: Some("blob_gc_fence_probe_claim".to_string()),
})
.await?;
let insert_attempt = writer
.execute(SqlStatement {
sql: "INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES (?1, 'entity', 'content', ?2, 0)"
.to_string(),
params: vec![
SqlValue::Text(insert_id.clone()),
SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
],
label: Some("blob_gc_fence_probe_insert_arm".to_string()),
})
.await;
let insert_fenced = fence_rejection(insert_attempt);
writer
.execute(SqlStatement {
sql: "INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES (?1, 'entity', 'content', ?2, 0)"
.to_string(),
params: vec![SqlValue::Text(update_id.clone()), SqlValue::Text(seed_ref)],
label: Some("blob_gc_fence_probe_update_arm_seed".to_string()),
})
.await?;
let update_attempt = writer
.execute(SqlStatement {
sql: "UPDATE attachments SET content_ref = ?1 \
WHERE record_uuid = ?2 AND role = 'content'"
.to_string(),
params: vec![
SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
SqlValue::Text(update_id.clone()),
],
label: Some("blob_gc_fence_probe_update_arm".to_string()),
})
.await;
let update_fenced = fence_rejection(update_attempt);
let insert2_attempt = writer
.execute(SqlStatement {
sql: "INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES (?1, 'note', 'evidence', ?2, 0)"
.to_string(),
params: vec![
SqlValue::Text(insert2_id.clone()),
SqlValue::Text(probe2_ref.clone()),
],
label: Some("blob_gc_fence_probe_insert2_arm".to_string()),
})
.await;
let insert2_fenced = fence_rejection(insert2_attempt);
writer
.execute(SqlStatement {
sql: "INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES (?1, 'note', 'evidence', ?2, 0)"
.to_string(),
params: vec![
SqlValue::Text(update2_id.clone()),
SqlValue::Text(seed2_ref),
],
label: Some("blob_gc_fence_probe_update2_arm_seed".to_string()),
})
.await?;
let update2_attempt = writer
.execute(SqlStatement {
sql: "UPDATE attachments SET content_ref = ?1 \
WHERE record_uuid = ?2 AND role = 'evidence'"
.to_string(),
params: vec![
SqlValue::Text(probe2_ref),
SqlValue::Text(update2_id.clone()),
],
label: Some("blob_gc_fence_probe_update2_arm".to_string()),
})
.await;
let update2_fenced = fence_rejection(update2_attempt);
writer
.execute(SqlStatement {
sql: "DELETE FROM attachments WHERE record_uuid IN (?1, ?2, ?3, ?4)"
.to_string(),
params: vec![
SqlValue::Text(insert_id.clone()),
SqlValue::Text(update_id.clone()),
SqlValue::Text(insert2_id.clone()),
SqlValue::Text(update2_id.clone()),
],
label: Some("blob_gc_fence_probe_cleanup_attachments".to_string()),
})
.await?;
writer
.execute(SqlStatement {
sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
params: vec![SqlValue::Text(claim_key)],
label: Some("blob_gc_fence_probe_cleanup_claim".to_string()),
})
.await?;
Ok(
Box::new((insert_fenced, update_fenced, insert2_fenced, update2_fenced))
as Box<dyn std::any::Any + Send>,
)
})
});
let outcome = sql.atomic_unit(op).await?;
let (insert_fenced, update_fenced, insert2_fenced, update2_fenced) = *outcome
.downcast::<(
Result<bool, String>,
Result<bool, String>,
Result<bool, String>,
Result<bool, String>,
)>()
.map_err(|_| {
StorageError::Internal("blob GC fence probe returned an unexpected outcome type".into())
})?;
let arm_verdict = |arm: &str, fenced: Result<bool, String>| -> StorageResult<()> {
match fenced {
Ok(true) => Ok(()),
Ok(false) => Err(StorageError::Unsupported {
capability: StorageCapability::Blob,
operation: "transactional_orphan_sweep".into(),
message: format!(
"the V21 fencing triggers exist by name but did not reject a claimed \
content_ref on the attachment {arm} path; refusing unfenced deletion"
),
}),
Err(other) => Err(StorageError::Unsupported {
capability: StorageCapability::Blob,
operation: "transactional_orphan_sweep".into(),
message: format!(
"the blob GC fence probe could not verify the attachment {arm} fence \
(unexpected rejection: {other}); refusing unfenced deletion"
),
}),
}
};
arm_verdict("INSERT", insert_fenced)?;
arm_verdict("UPDATE", update_fenced)?;
arm_verdict("second-digest INSERT", insert2_fenced)?;
arm_verdict("second-digest UPDATE", update2_fenced)
}
async fn validate_blob_gc_evidence(sql: &dyn SqlAccess) -> StorageResult<()> {
let mut reader = sql.reader().await?;
let canonical_bytes = match reader
.query_row(SqlStatement {
sql: "SELECT length(CAST('x' AS BLOB))".to_string(),
params: vec![],
label: Some("blob_gc_validate_encoding_width".to_string()),
})
.await?
.and_then(|row| row.columns.first().map(|column| column.value.clone()))
{
Some(SqlValue::Integer(width)) if (1..=4).contains(&width) => width * 64,
other => {
return Err(invalid_content_ref(format!(
"the text-encoding width probe returned {other:?}; refusing GC validation"
)));
}
};
let invalid_claim = reader
.query_row(SqlStatement {
sql: "SELECT content_ref FROM blob_gc_claims \
WHERE typeof(content_ref) <> 'text' \
OR length(content_ref) <> 64 \
OR length(CAST(content_ref AS BLOB)) <> ?1 \
OR content_ref GLOB '*[^0-9a-f]*' \
LIMIT 1"
.to_string(),
params: vec![SqlValue::Integer(canonical_bytes)],
label: Some("blob_gc_validate_existing_claims".to_string()),
})
.await?;
if invalid_claim.is_some() {
return Err(invalid_content_ref(
"blob_gc_claims.content_ref contained a non-canonical value".into(),
));
}
let invalid_live = reader
.query_row(SqlStatement {
sql: "SELECT content_ref FROM attachments \
WHERE typeof(content_ref) <> 'text' \
OR length(content_ref) <> 64 \
OR length(CAST(content_ref AS BLOB)) <> ?1 \
OR content_ref GLOB '*[^0-9a-f]*' \
LIMIT 1"
.to_string(),
params: vec![SqlValue::Integer(canonical_bytes)],
label: Some("blob_gc_validate_live_refs".to_string()),
})
.await?;
if invalid_live.is_some() {
return Err(invalid_content_ref(
"attachments.content_ref contained a non-canonical value".into(),
));
}
Ok(())
}
async fn release_abandoned_blob_gc_claim_batch(sql: &dyn SqlAccess) -> StorageResult<u64> {
let op: AtomicUnitOp = Box::new(move |writer| {
Box::pin(async move {
let released = writer
.execute(SqlStatement {
sql: "DELETE FROM blob_gc_claims \
WHERE rowid IN ( \
SELECT rowid FROM blob_gc_claims \
ORDER BY rowid LIMIT ?1 \
)"
.to_string(),
params: vec![SqlValue::Integer(BLOB_GC_CLAIM_BATCH_SIZE as i64)],
label: Some("blob_gc_release_abandoned_claim_batch".to_string()),
})
.await?;
Ok(Box::new(released) as Box<dyn std::any::Any + Send>)
})
});
let released = sql.atomic_unit(op).await?;
released.downcast::<u64>().map(|count| *count).map_err(|_| {
StorageError::Internal(
"transactional orphan sweep returned an unexpected recovery count type".into(),
)
})
}
async fn claim_blob_gc_batch(
sql: &dyn SqlAccess,
root_key: String,
candidates: &[(ContentRef, bool)],
dry_run: bool,
) -> StorageResult<BlobGcBatchRows> {
debug_assert!(candidates.len() <= BLOB_GC_CLAIM_BATCH_SIZE);
let eligible_refs = candidates
.iter()
.filter(|(_, within_grace)| !within_grace)
.map(|(content_ref, _)| content_ref.to_string())
.collect::<Vec<_>>();
let grace_refs = candidates
.iter()
.filter(|(_, within_grace)| *within_grace)
.map(|(content_ref, _)| content_ref.to_string())
.collect::<Vec<_>>();
let eligible_json = serde_json::to_string(&eligible_refs).map_err(|error| {
StorageError::Internal(format!(
"failed to prepare blob GC eligible candidate batch: {error}"
))
})?;
let grace_json = serde_json::to_string(&grace_refs).map_err(|error| {
StorageError::Internal(format!(
"failed to prepare blob GC grace candidate batch: {error}"
))
})?;
let claimed_at = chrono::Utc::now().timestamp_micros();
let op: AtomicUnitOp = Box::new(move |writer| {
Box::pin(async move {
let grace_period_skipped = required_nonnegative_count(
writer
.query_scalar(SqlStatement {
sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
WHERE NOT EXISTS ( \
SELECT 1 FROM attachments \
WHERE content_ref = candidate.value \
)"
.to_string(),
params: vec![SqlValue::Text(grace_json)],
label: Some("blob_gc_count_grace_candidates_batch".to_string()),
})
.await?,
"blob_gc_count_grace_candidates_batch",
)?;
if dry_run {
let would_delete = required_nonnegative_count(
writer
.query_scalar(SqlStatement {
sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
WHERE NOT EXISTS ( \
SELECT 1 FROM attachments \
WHERE content_ref = candidate.value \
)"
.to_string(),
params: vec![SqlValue::Text(eligible_json)],
label: Some("blob_gc_count_dry_run_candidates_batch".to_string()),
})
.await?,
"blob_gc_count_dry_run_candidates_batch",
)?;
return Ok(Box::new(BlobGcBatchRows {
grace_period_skipped,
would_delete,
claimed_rows: Vec::new(),
}) as Box<dyn std::any::Any + Send>);
}
writer
.execute(SqlStatement {
sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
SELECT ?1, candidate.value, ?3 \
FROM json_each(?2) AS candidate \
WHERE NOT EXISTS ( \
SELECT 1 FROM attachments \
WHERE content_ref = candidate.value \
)"
.to_string(),
params: vec![
SqlValue::Text(root_key.clone()),
SqlValue::Text(eligible_json),
SqlValue::Integer(claimed_at),
],
label: Some("blob_gc_claim_candidate_batch".to_string()),
})
.await?;
let claimed_rows = writer
.query_all(SqlStatement {
sql: "SELECT content_ref FROM blob_gc_claims \
WHERE root_key = ?1 ORDER BY content_ref"
.to_string(),
params: vec![SqlValue::Text(root_key)],
label: Some("blob_gc_claimed_candidate_batch".to_string()),
})
.await?;
Ok(Box::new(BlobGcBatchRows {
grace_period_skipped,
would_delete: claimed_rows.len() as u64,
claimed_rows,
}) as Box<dyn std::any::Any + Send>)
})
});
let rows = sql.atomic_unit(op).await?;
rows.downcast::<BlobGcBatchRows>()
.map(|rows| *rows)
.map_err(|_| {
StorageError::Internal(
"transactional orphan sweep returned an unexpected batch-row type".into(),
)
})
}
fn parse_blob_gc_claim_rows(rows: Vec<SqlRow>) -> StorageResult<Vec<ContentRef>> {
let mut claimed = Vec::with_capacity(rows.len());
for row in rows {
let raw = match row.get("content_ref") {
Some(SqlValue::Text(raw)) => raw.clone(),
_ => {
return Err(invalid_content_ref(
"blob_gc_claims.content_ref contained a non-text value".into(),
));
}
};
claimed.push(ContentRef::from_hex(raw).map_err(invalid_content_ref)?);
}
Ok(claimed)
}
async fn release_blob_gc_batch(sql: &dyn SqlAccess, root_key: String) -> StorageResult<()> {
let cleanup: AtomicUnitOp = Box::new(move |writer| {
Box::pin(async move {
writer
.execute(SqlStatement {
sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
params: vec![SqlValue::Text(root_key)],
label: Some("blob_gc_release_claim_batch".to_string()),
})
.await?;
Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
})
});
sql.atomic_unit(cleanup).await?;
Ok(())
}
type SweepLockMap = HashMap<Option<PathBuf>, Arc<DatabaseGcProcessLock>>;
#[derive(Debug, Default)]
struct DatabaseGcProcessLock {
held: StdMutex<bool>,
released: std::sync::Condvar,
#[cfg(test)]
waiters: std::sync::atomic::AtomicUsize,
}
impl DatabaseGcProcessLock {
fn acquire(self: &Arc<Self>) -> DatabaseGcProcessGuard {
let mut held = self
.held
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
while *held {
#[cfg(test)]
self.waiters
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
held = self
.released
.wait(held)
.unwrap_or_else(std::sync::PoisonError::into_inner);
#[cfg(test)]
self.waiters
.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
}
*held = true;
DatabaseGcProcessGuard {
lock: Arc::clone(self),
}
}
fn try_acquire(self: &Arc<Self>) -> Option<DatabaseGcProcessGuard> {
let mut held = self
.held
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if *held {
return None;
}
*held = true;
Some(DatabaseGcProcessGuard {
lock: Arc::clone(self),
})
}
}
#[derive(Debug)]
struct DatabaseGcProcessGuard {
lock: Arc<DatabaseGcProcessLock>,
}
impl Drop for DatabaseGcProcessGuard {
fn drop(&mut self) {
let mut held = self
.lock
.held
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
debug_assert!(*held, "database GC process owner released twice");
*held = false;
self.lock.released.notify_one();
}
}
fn database_sweep_locks() -> &'static StdMutex<SweepLockMap> {
static REGISTRY: OnceLock<StdMutex<SweepLockMap>> = OnceLock::new();
REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
}
fn sweep_lock_for_database(database_path: Option<&Path>) -> Arc<DatabaseGcProcessLock> {
let key = database_path.map(Path::to_path_buf);
let mut locks = database_sweep_locks()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
locks
.entry(key)
.or_insert_with(|| Arc::new(DatabaseGcProcessLock::default()))
.clone()
}
#[cfg(test)]
pub(crate) fn database_gc_waiter_count(database_path: Option<&Path>) -> usize {
sweep_lock_for_database(database_path)
.waiters
.load(std::sync::atomic::Ordering::SeqCst)
}
pub struct DatabaseGcOwnerGuard {
_process_guard: DatabaseGcProcessGuard,
_advisory_guard: Option<fs::File>,
database_path: Option<PathBuf>,
}
pub(crate) fn acquire_database_gc_owner_for_path_blocking(
database_path: Option<PathBuf>,
) -> StorageResult<DatabaseGcOwnerGuard> {
let process_guard = sweep_lock_for_database(database_path.as_deref()).acquire();
let advisory_guard = acquire_database_gc_lock(database_path.as_deref())?;
Ok(DatabaseGcOwnerGuard {
_process_guard: process_guard,
_advisory_guard: advisory_guard,
database_path,
})
}
pub(crate) fn try_acquire_database_gc_owner_for_path(
database_path: PathBuf,
) -> StorageResult<DatabaseGcOwnerGuard> {
let process_guard = sweep_lock_for_database(Some(&database_path))
.try_acquire()
.ok_or_else(|| {
StorageError::Internal(format!(
"database GC owner for {} is already held; retry schema migration through the \
coordinated backend boot path",
database_path.display()
))
})?;
let lock_path = database_gc_lock_path(&database_path);
let advisory_guard = fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&lock_path)
.map_err(|error| map_io_err(error, "database_gc_lock_open"))?;
fs4::FileExt::try_lock(&advisory_guard)
.map_err(|error| map_io_err(error.into(), "database_gc_lock_try_acquire"))?;
Ok(DatabaseGcOwnerGuard {
_process_guard: process_guard,
_advisory_guard: Some(advisory_guard),
database_path: Some(database_path),
})
}
impl std::fmt::Debug for DatabaseGcOwnerGuard {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("DatabaseGcOwnerGuard")
.field("database_path", &self.database_path)
.finish_non_exhaustive()
}
}
impl DatabaseGcOwnerGuard {
pub fn database_path(&self) -> Option<&Path> {
self.database_path.as_deref()
}
}
pub async fn acquire_database_gc_owner(sql: &dyn SqlAccess) -> StorageResult<DatabaseGcOwnerGuard> {
let database_path = sql.database_path();
tokio::task::spawn_blocking(move || acquire_database_gc_owner_for_path_blocking(database_path))
.await
.map_err(|error| {
StorageError::driver(StorageCapability::Blob, "acquire_database_gc_owner", error)
})?
}
fn root_write_locks() -> &'static StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>> {
static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>>> =
OnceLock::new();
REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
}
fn write_lock_for_root(root: &Path) -> std::io::Result<Arc<tokio::sync::Mutex<()>>> {
let canonical = root.canonicalize()?;
let mut locks = root_write_locks()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Ok(locks
.entry(canonical)
.or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
.clone())
}
#[derive(Debug)]
pub struct FsBlobStore {
root: PathBuf,
root_handle: Arc<fs::File>,
floor_bytes: u64,
write_lock: Arc<tokio::sync::Mutex<()>>,
orphan_sweep_grace: Duration,
}
impl FsBlobStore {
pub const DEFAULT_FLOOR_BYTES: u64 = 100_000_000_000;
pub const DEFAULT_ORPHAN_SWEEP_GRACE: Duration = Duration::from_secs(3600);
pub fn new(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
match fs::create_dir(&root) {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
let parent = root
.parent()
.filter(|path| !path.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."));
return Err(std::io::Error::new(
std::io::ErrorKind::NotFound,
format!("blob root parent missing: {}", parent.display()),
)
.into());
}
Err(error) => return Err(error.into()),
}
let store = Self::open_existing(root, floor_bytes)?;
#[cfg(unix)]
{
use std::os::fd::AsRawFd;
let publication = BlobPublication {
#[cfg(test)]
hook: sync_hook::take(&store.root).and_then(|hook| hook.publication),
};
let parent = openat_dir_no_follow(store.root_handle.as_raw_fd(), "..")?;
publication.sync_directory("init_sync_root", &store.root_handle)?;
publication.sync_directory("init_sync_parent", &parent)?;
}
Ok(store)
}
pub fn open_existing(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
let root = root.canonicalize()?;
let metadata = fs::metadata(&root)?;
if !metadata.is_dir() {
return Err(SqliteError::InvalidData(format!(
"blob store root is not a directory: {}",
root.display()
)));
}
let root_handle = Arc::new(open_blob_root_handle(&root)?);
let write_lock = write_lock_for_root(&root)?;
verify_blob_root_identity(&root, &root_handle)?;
Ok(Self {
root,
root_handle,
floor_bytes,
write_lock,
orphan_sweep_grace: Self::DEFAULT_ORPHAN_SWEEP_GRACE,
})
}
pub fn with_orphan_sweep_grace(mut self, grace_period: Duration) -> Self {
self.orphan_sweep_grace = grace_period;
self
}
pub fn root(&self) -> &Path {
&self.root
}
}
#[async_trait]
impl BlobStore for FsBlobStore {
async fn begin_upload(&self, declared_size: u64) -> StorageResult<UploadId> {
uploads::begin(self, declared_size).await
}
async fn append_part(&self, id: &UploadId, bytes: Vec<u8>) -> StorageResult<u64> {
uploads::append(self, id.clone(), bytes).await
}
async fn commit_upload(&self, id: &UploadId, content_ref: &ContentRef) -> StorageResult<()> {
uploads::commit(self, id.clone(), content_ref.clone()).await
}
async fn abort_upload(&self, id: &UploadId) -> StorageResult<()> {
uploads::abort(self, id.clone()).await
}
async fn sweep_uploads(&self, idle_for: Duration) -> StorageResult<u64> {
uploads::sweep(self, idle_for).await
}
async fn put(&self, bytes: Vec<u8>) -> StorageResult<ContentRef> {
let owned_guard = self.write_lock.clone().lock_owned().await;
let root = self.root.clone();
let root_handle = Arc::clone(&self.root_handle);
let floor_bytes = self.floor_bytes;
#[cfg(test)]
let hook = sync_hook::take(&root);
tokio::task::spawn_blocking(move || {
#[cfg_attr(not(test), allow(clippy::let_and_return))]
let result = {
let _owned_guard = owned_guard;
#[cfg(test)]
if let Some(h) = &hook {
let _ = h.reached.send(());
let _ = h.release.recv();
}
put_blocking_from_root_handle(
&root,
&root_handle,
floor_bytes,
bytes,
#[cfg(unix)]
&BlobPublication {
#[cfg(test)]
hook: hook.as_ref().and_then(|hook| hook.publication.clone()),
},
)
};
#[cfg(test)]
if let Some(h) = &hook {
let _ = h.done.send(());
}
result
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Blob, "put", e))?
}
async fn get_bounded_verified(
&self,
content_ref: &ContentRef,
max_bytes: u64,
) -> StorageResult<Vec<u8>> {
if max_bytes > MAX_BLOB_WHOLE_BYTES {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Blob,
operation: "get_bounded_verified".into(),
message: format!(
"max_bytes {max_bytes} exceeds the {MAX_BLOB_WHOLE_BYTES}-byte portable whole-buffer envelope"
),
});
}
let root = self.root.clone();
let root_handle = Arc::clone(&self.root_handle);
let content_ref = content_ref.clone();
#[cfg(test)]
let read_hook = bounded_read_sync_hook::take(&root);
tokio::task::spawn_blocking(move || {
let mut file = open_blob_shard_file_no_follow(&root, &root_handle, &content_ref)
.map_err(|e| {
if e.kind() == std::io::ErrorKind::NotFound {
StorageError::NotFound {
capability: StorageCapability::Blob,
resource: "blob",
key: content_ref.to_string(),
}
} else if e.kind() == std::io::ErrorKind::Unsupported {
StorageError::Unsupported {
capability: StorageCapability::Blob,
operation: "get_bounded_verified".into(),
message: e.to_string(),
}
} else {
map_io_err(e, "get_bounded_verified.open")
}
})?;
let metadata_bytes = file
.metadata()
.map_err(|e| map_io_err(e, "get_bounded_verified.metadata"))?
.len();
#[cfg(test)]
if let Some(hook) = &read_hook {
let _ = hook.reached.send(());
let _ = hook.release.recv();
}
if metadata_bytes > max_bytes {
return Err(StorageError::BlobTooLarge {
content_ref,
max_bytes,
observed_at_least: metadata_bytes,
});
}
let mut bytes = Vec::with_capacity(metadata_bytes as usize);
(&mut file)
.take(max_bytes + 1)
.read_to_end(&mut bytes)
.map_err(|e| map_io_err(e, "get_bounded_verified.read"))?;
let actual_bytes = bytes.len() as u64;
if actual_bytes > max_bytes {
return Err(StorageError::BlobTooLarge {
content_ref,
max_bytes,
observed_at_least: actual_bytes,
});
}
if metadata_bytes != actual_bytes {
return Err(StorageError::BlobSizeMismatch {
content_ref,
metadata_bytes,
actual_bytes,
});
}
let actual = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
if actual != content_ref {
return Err(StorageError::BlobDigestMismatch {
expected: content_ref,
actual,
});
}
Ok(bytes)
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Blob, "get_bounded_verified", e))?
}
async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool> {
let root = self.root.clone();
let root_handle = Arc::clone(&self.root_handle);
let content_ref = content_ref.clone();
tokio::task::spawn_blocking(move || {
match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
Ok(_) => Ok(true),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(error) => Err(map_io_err(error, "exists")),
}
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Blob, "exists", e))?
}
async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>> {
let root = self.root.clone();
let root_handle = Arc::clone(&self.root_handle);
let content_ref = content_ref.clone();
tokio::task::spawn_blocking(move || {
match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
Ok(file) => file
.metadata()
.map(|metadata| Some(metadata.len()))
.map_err(|error| map_io_err(error, "size")),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(error) => Err(map_io_err(error, "size")),
}
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Blob, "size", e))?
}
async fn delete(&self, content_ref: &ContentRef) -> StorageResult<bool> {
let root = self.root.clone();
let root_handle = Arc::clone(&self.root_handle);
let content_ref = content_ref.clone();
tokio::task::spawn_blocking(move || {
match unlink_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
Ok(()) => Ok(true),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(e) => Err(map_io_err(e, "delete")),
}
})
.await
.map_err(|e| StorageError::driver(StorageCapability::Blob, "delete", e))?
}
async fn orphan_sweep(
&self,
config: &BlobOrphanSweepConfig,
) -> StorageResult<BlobOrphanSweepResult> {
let _ = config;
Err(StorageError::Unsupported {
capability: StorageCapability::Blob,
operation: "orphan_sweep".into(),
message: "caller-snapshot orphan_sweep is disabled in this compatibility release; \
it cannot prove a completed V21 attachment epoch, use \
transactional_orphan_sweep instead"
.into(),
})
}
async fn transactional_orphan_sweep(
&self,
sql: &dyn SqlAccess,
dry_run: bool,
) -> StorageResult<BlobOrphanSweepResult> {
if !blob_gc_fencing_complete(sql).await? {
return Err(unsupported_blob_gc_epoch());
}
let database_path = sql.database_path();
let lock_database_path = database_path.clone();
#[cfg(test)]
let hook_database_path = database_path.clone();
let (database_guard, database_file_guard) = tokio::task::spawn_blocking(move || {
let process_guard = sweep_lock_for_database(lock_database_path.as_deref()).acquire();
let file_guard = acquire_database_gc_lock(lock_database_path.as_deref())?;
#[cfg(test)]
if let Some(hook) = db_ownership_sync_hook::take(hook_database_path.as_deref()) {
let _ = hook.reached.send(());
let _ = hook.release.recv();
}
Ok::<_, StorageError>((process_guard, file_guard))
})
.await
.map_err(|e| {
StorageError::driver(
StorageCapability::Blob,
"transactional_orphan_sweep_lock",
e,
)
})??;
if !blob_gc_fencing_complete(sql).await? {
return Err(unsupported_blob_gc_epoch());
}
let root_guard = self.write_lock.clone().lock_owned().await;
let root = self.root.clone();
let root_handle = Arc::clone(&self.root_handle);
let scan_root = root.clone();
let scan_root_handle = Arc::clone(&root_handle);
let grace_period = self.orphan_sweep_grace;
let (write_guards, canonical_root, prepared) = tokio::task::spawn_blocking(move || {
verify_blob_root_identity(&scan_root, &scan_root_handle)
.map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
let canonical_root = scan_root;
let root_write_guard =
acquire_root_write_lock_anchored(&canonical_root, &scan_root_handle)?;
let candidates = walk_blob_files_from_root_handle(&scan_root_handle, &canonical_root)
.map_err(|e| map_io_err(e, "transactional_orphan_sweep_walk"))?;
verify_blob_root_identity(&canonical_root, &scan_root_handle)
.map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
let prepared = prepare_transactional_sweep(candidates, grace_period);
Ok::<_, StorageError>((
(
database_guard,
database_file_guard,
root_guard,
root_write_guard,
),
canonical_root,
prepared,
))
})
.await
.map_err(|e| {
StorageError::driver(
StorageCapability::Blob,
"transactional_orphan_sweep_walk",
e,
)
})??;
let root_key = blob_root_key(&canonical_root);
validate_blob_gc_evidence(sql).await?;
blob_gc_fence_probe(sql).await?;
if !dry_run {
loop {
let released = release_abandoned_blob_gc_claim_batch(sql).await?;
if released < BLOB_GC_CLAIM_BATCH_SIZE as u64 {
break;
}
}
}
let mut write_guards = write_guards;
let mut result = prepared.result;
let mut delete_error = None;
#[cfg(test)]
let mut hook: Option<sync_hook::Hook> = None;
#[cfg(not(test))]
let mut hook: Option<()> = None;
#[cfg(test)]
let mut hook_paused = false;
for candidates in prepared.candidates.chunks(BLOB_GC_CLAIM_BATCH_SIZE) {
let batch = claim_blob_gc_batch(sql, root_key.clone(), candidates, dry_run).await?;
result.grace_period_skipped += batch.grace_period_skipped;
result.would_delete += batch.would_delete;
if dry_run {
continue;
}
let claimed_refs = parse_blob_gc_claim_rows(batch.claimed_rows)?;
if claimed_refs.is_empty() {
continue;
}
#[cfg(test)]
if hook.is_none() {
hook = sync_hook::take(&root);
}
#[cfg(test)]
let pause_hook = hook.is_some() && !hook_paused;
#[cfg(test)]
if pause_hook {
hook_paused = true;
}
let delete_root = canonical_root.clone();
let delete_root_handle = Arc::clone(&root_handle);
let (returned_guards, deleted, batch_delete_error, returned_hook) =
tokio::task::spawn_blocking(move || {
#[cfg(test)]
if pause_hook {
if let Some(hook) = &hook {
let _ = hook.reached.send(());
let _ = hook.release.recv();
}
}
let mut deleted = 0_u64;
let mut first_error = None;
for content_ref in claimed_refs {
match unlink_blob_shard_file_no_follow(
&delete_root,
&delete_root_handle,
&content_ref,
) {
Ok(()) => deleted += 1,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => {
first_error =
Some(map_io_err(error, "transactional_orphan_sweep_delete"));
break;
}
}
}
(write_guards, deleted, first_error, hook)
})
.await
.map_err(|error| {
StorageError::driver(
StorageCapability::Blob,
"transactional_orphan_sweep_delete",
error,
)
})?;
write_guards = returned_guards;
hook = returned_hook;
result.deleted += deleted;
release_blob_gc_batch(sql, root_key.clone()).await?;
if batch_delete_error.is_some() {
delete_error = batch_delete_error;
break;
}
}
drop(write_guards);
#[cfg(test)]
if let Some(hook) = hook {
let _ = hook.done.send(());
}
#[cfg(not(test))]
let _ = hook;
if let Some(error) = delete_error {
return Err(error);
}
Ok(result)
}
}
#[cfg(test)]
mod sync_hook {
use std::collections::{HashMap, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::mpsc::{Receiver, Sender};
use std::sync::{Mutex as StdMutex, OnceLock};
pub(super) struct Hook {
pub(super) reached: Sender<()>,
pub(super) release: Receiver<()>,
pub(super) done: Sender<()>,
#[cfg(unix)]
pub(super) publication: Option<Publication>,
}
fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
}
pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>, Receiver<()>) {
let canonical = root
.canonicalize()
.expect("root must exist before installing a sync_hook");
let (reached_tx, reached_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let (done_tx, done_rx) = std::sync::mpsc::channel();
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.entry(canonical)
.or_default()
.push_back(Hook {
reached: reached_tx,
release: release_rx,
done: done_tx,
#[cfg(unix)]
publication: None,
});
(reached_rx, release_tx, done_rx)
}
pub(super) fn take(root: &Path) -> Option<Hook> {
let canonical = root.canonicalize().ok()?;
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get_mut(&canonical)
.and_then(VecDeque::pop_front)
}
#[cfg(unix)]
type StepAction = (&'static str, Box<dyn FnOnce() + Send>);
#[cfg(unix)]
type DirectorySync = (&'static str, u64, u64);
#[cfg(unix)]
#[derive(Clone)]
pub(super) struct Publication {
pub(super) completed: std::sync::Arc<StdMutex<Vec<&'static str>>>,
pub(super) directories: std::sync::Arc<StdMutex<Vec<DirectorySync>>>,
fail_at: Option<&'static str>,
action: std::sync::Arc<StdMutex<Option<StepAction>>>,
}
#[cfg(unix)]
impl Publication {
pub(super) fn before(&self, operation: &'static str) -> std::io::Result<()> {
if self.fail_at == Some(operation) {
return Err(std::io::Error::other("injected publication failure"));
}
let action = {
let mut slot = self.action.lock().unwrap();
if slot.as_ref().is_some_and(|(at, _)| *at == operation) {
slot.take()
} else {
None
}
};
if let Some((_, action)) = action {
action();
}
Ok(())
}
pub(super) fn on_step(
&self,
operation: &'static str,
action: impl FnOnce() + Send + 'static,
) {
*self.action.lock().unwrap() = Some((operation, Box::new(action)));
}
pub(super) fn completed(&self, operation: &'static str) {
self.completed.lock().unwrap().push(operation);
}
pub(super) fn directory_synced(
&self,
operation: &'static str,
directory: &std::fs::File,
) -> std::io::Result<()> {
use std::os::unix::fs::MetadataExt;
let metadata = directory.metadata()?;
self.directories
.lock()
.unwrap()
.push((operation, metadata.dev(), metadata.ino()));
Ok(())
}
}
#[cfg(unix)]
pub(super) fn install_publication(root: &Path, fail_at: Option<&'static str>) -> Publication {
let publication = Publication {
completed: std::sync::Arc::default(),
directories: std::sync::Arc::default(),
fail_at,
action: std::sync::Arc::default(),
};
let (reached, _) = std::sync::mpsc::channel();
let (_, release) = std::sync::mpsc::channel();
let (done, _) = std::sync::mpsc::channel();
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.entry(root.canonicalize().unwrap())
.or_default()
.push_back(Hook {
reached,
release,
done,
publication: Some(publication.clone()),
});
publication
}
}
#[cfg(all(test, unix))]
mod walk_leaf_sync_hook {
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::mpsc::{Receiver, Sender};
use std::sync::{Mutex as StdMutex, OnceLock};
pub(super) struct Hook {
pub(super) reached: Sender<()>,
pub(super) release: Receiver<()>,
}
fn registry() -> &'static StdMutex<HashMap<PathBuf, Hook>> {
static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Hook>>> = OnceLock::new();
REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
}
pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
let canonical = root
.canonicalize()
.expect("root must exist before installing a walk_leaf_sync_hook");
let (reached_tx, reached_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(
canonical,
Hook {
reached: reached_tx,
release: release_rx,
},
);
(reached_rx, release_tx)
}
pub(super) fn take(root: &Path) -> Option<Hook> {
let canonical = root.canonicalize().ok()?;
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&canonical)
}
}
#[cfg(test)]
mod bounded_read_sync_hook {
use std::collections::{HashMap, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::mpsc::{Receiver, Sender};
use std::sync::{Mutex as StdMutex, OnceLock};
pub(super) struct Hook {
pub(super) reached: Sender<()>,
pub(super) release: Receiver<()>,
}
fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
}
pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
let canonical = root
.canonicalize()
.expect("root must exist before installing a bounded-read hook");
let (reached_tx, reached_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.entry(canonical)
.or_default()
.push_back(Hook {
reached: reached_tx,
release: release_rx,
});
(reached_rx, release_tx)
}
pub(super) fn take(root: &Path) -> Option<Hook> {
let canonical = root.canonicalize().ok()?;
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get_mut(&canonical)
.and_then(VecDeque::pop_front)
}
}
#[cfg(test)]
mod db_ownership_sync_hook {
use std::collections::{HashMap, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::mpsc::{Receiver, Sender};
use std::sync::{Mutex as StdMutex, OnceLock};
pub(super) struct Hook {
pub(super) reached: Sender<()>,
pub(super) release: Receiver<()>,
}
fn registry() -> &'static StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>> {
static REGISTRY: OnceLock<StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>>> =
OnceLock::new();
REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
}
pub(super) fn install(database_path: Option<&Path>) -> (Receiver<()>, Sender<()>) {
let key = database_path.map(Path::to_path_buf);
let (reached_tx, reached_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.entry(key)
.or_default()
.push_back(Hook {
reached: reached_tx,
release: release_rx,
});
(reached_rx, release_tx)
}
pub(super) fn take(database_path: Option<&Path>) -> Option<Hook> {
let key = database_path.map(Path::to_path_buf);
registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get_mut(&key)
.and_then(VecDeque::pop_front)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn store(floor_bytes: u64) -> (tempfile::TempDir, FsBlobStore) {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
let store = FsBlobStore::new(root, floor_bytes)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO);
(dir, store)
}
#[cfg(unix)]
#[tokio::test]
async fn publication_barriers_cover_fresh_shards_and_dedup() {
use std::os::unix::fs::MetadataExt;
let (_dir, store) = store(0);
let bytes = b"publication barrier order".to_vec();
let hook = sync_hook::install_publication(store.root(), None);
let content_ref = store.put(bytes.clone()).await.unwrap();
assert_eq!(
*hook.completed.lock().unwrap(),
[
"put_fsync",
"put_persist",
"put_sync_shard",
"put_sync_parent",
"put_sync_root"
]
);
let path = shard_path(store.root(), &content_ref);
let directory_identity = |operation, path: &Path| {
let metadata = fs::metadata(path).unwrap();
(operation, metadata.dev(), metadata.ino())
};
let expected_directories = [
directory_identity("put_sync_shard", path.parent().unwrap()),
directory_identity("put_sync_parent", path.parent().unwrap().parent().unwrap()),
directory_identity("put_sync_root", store.root()),
];
assert_eq!(*hook.directories.lock().unwrap(), expected_directories);
let inode = fs::metadata(&path).unwrap().ino();
let reopened = FsBlobStore::open_existing(store.root().to_path_buf(), 0).unwrap();
let hook = sync_hook::install_publication(reopened.root(), None);
assert_eq!(reopened.put(bytes.clone()).await.unwrap(), content_ref);
assert_eq!(
*hook.completed.lock().unwrap(),
[
"put_fsync",
"put_sync_shard",
"put_sync_parent",
"put_sync_root"
]
);
assert_eq!(*hook.directories.lock().unwrap(), expected_directories);
assert_eq!(fs::metadata(&path).unwrap().ino(), inode);
assert_eq!(
reopened
.get_bounded_verified(&content_ref, bytes.len() as u64)
.await
.unwrap(),
bytes
);
}
#[cfg(unix)]
#[test]
fn publication_barriers_initialize_root_but_not_read_only_open() {
use std::os::unix::fs::MetadataExt;
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
fs::create_dir(&root).unwrap();
let hook = sync_hook::install_publication(&root, None);
fs::remove_dir(&root).unwrap();
let store = FsBlobStore::new(root.clone(), 0).unwrap();
assert_eq!(
*hook.completed.lock().unwrap(),
["init_sync_root", "init_sync_parent"]
);
let root_metadata = fs::metadata(&root).unwrap();
let parent_metadata = fs::metadata(dir.path()).unwrap();
assert_eq!(
*hook.directories.lock().unwrap(),
[
("init_sync_root", root_metadata.dev(), root_metadata.ino()),
(
"init_sync_parent",
parent_metadata.dev(),
parent_metadata.ino()
),
]
);
let hook = sync_hook::install_publication(store.root(), Some("init_sync_root"));
let reopened = FsBlobStore::open_existing(root.clone(), 0).unwrap();
assert!(hook.completed.lock().unwrap().is_empty());
assert!(FsBlobStore::new(root, 0).is_err());
assert!(hook.completed.lock().unwrap().is_empty());
assert_eq!(reopened.root(), store.root());
}
#[cfg(unix)]
#[test]
fn publication_barriers_retry_each_initialization_failure() {
for operation in ["init_sync_root", "init_sync_parent"] {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
fs::create_dir(&root).unwrap();
let hook = sync_hook::install_publication(&root, Some(operation));
fs::remove_dir(&root).unwrap();
assert!(FsBlobStore::new(root.clone(), 0).is_err(), "{operation}");
assert!(root.is_dir());
assert!(!hook.completed.lock().unwrap().contains(&operation));
let retry = sync_hook::install_publication(&root, None);
FsBlobStore::new(root, 0).unwrap();
assert_eq!(
*retry.completed.lock().unwrap(),
["init_sync_root", "init_sync_parent"]
);
}
}
#[test]
fn publication_barriers_refuse_missing_root_parent_without_creating_it() {
let dir = tempfile::tempdir().unwrap();
let parent = dir.path().join("missing").join("nested");
let root = parent.join("blobs");
let error = match FsBlobStore::new(root.clone(), 0) {
Ok(_) => panic!("missing parent must not be created"),
Err(error) => error,
};
assert!(
error.to_string().contains("blob root parent missing"),
"{error}"
);
assert!(
error.to_string().contains(&parent.display().to_string()),
"{error}"
);
assert!(!dir.path().join("missing").exists());
fs::create_dir_all(&parent).unwrap();
FsBlobStore::new(root, 0).unwrap();
}
#[cfg(unix)]
#[test]
fn publication_barriers_fsync_propagates_kernel_errors() {
let (socket, _peer) = std::os::unix::net::UnixStream::pair().unwrap();
let handle = fs::File::from(std::os::fd::OwnedFd::from(socket));
assert!(sync_directory(&handle).is_err());
}
#[cfg(unix)]
#[tokio::test]
async fn publication_barriers_fail_before_rename_without_publishing() {
for operation in ["put_fsync", "put_persist"] {
let (_dir, store) = store(0);
let bytes = b"unpublished failure".to_vec();
let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
let hook = sync_hook::install_publication(store.root(), Some(operation));
let error = store.put(bytes.clone()).await.unwrap_err();
assert!(error.to_string().contains(operation), "{error}");
let path = shard_path(store.root(), &content_ref);
assert!(!path.exists());
assert_eq!(fs::read_dir(path.parent().unwrap()).unwrap().count(), 0);
assert!(!hook.completed.lock().unwrap().contains(&operation));
assert_eq!(store.put(bytes).await.unwrap(), content_ref);
}
}
#[cfg(unix)]
#[tokio::test]
async fn publication_barriers_retry_after_each_directory_failure() {
use std::os::unix::fs::MetadataExt;
for operation in ["put_sync_shard", "put_sync_parent", "put_sync_root"] {
let (_dir, store) = store(0);
let bytes = b"retry a visible but unacknowledged blob".to_vec();
let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
sync_hook::install_publication(store.root(), Some(operation));
let error = store.put(bytes.clone()).await.unwrap_err();
assert!(error.to_string().contains(operation), "{error}");
let path = shard_path(store.root(), &content_ref);
let inode = fs::metadata(&path).unwrap().ino();
assert_eq!(fs::read(&path).unwrap(), bytes);
let reopened = FsBlobStore::open_existing(store.root().to_path_buf(), 0).unwrap();
let failed_retry = sync_hook::install_publication(reopened.root(), Some(operation));
let error = reopened.put(bytes.clone()).await.unwrap_err();
assert!(error.to_string().contains(operation), "{error}");
assert!(!failed_retry
.completed
.lock()
.unwrap()
.contains(&"put_persist"));
let repaired = sync_hook::install_publication(reopened.root(), None);
assert_eq!(reopened.put(bytes.clone()).await.unwrap(), content_ref);
assert_eq!(
*repaired.completed.lock().unwrap(),
[
"put_fsync",
"put_sync_shard",
"put_sync_parent",
"put_sync_root"
]
);
assert_eq!(fs::metadata(&path).unwrap().ino(), inode);
assert_eq!(fs::read(&path).unwrap(), bytes);
}
}
#[cfg(unix)]
#[tokio::test]
async fn publication_barriers_fault_is_scoped_to_one_root_and_put() {
let (_a, first) = store(0);
let (_b, second) = store(0);
sync_hook::install_publication(first.root(), Some("put_sync_shard"));
second.put(b"second".to_vec()).await.unwrap();
assert!(first.put(b"first".to_vec()).await.is_err());
first.put(b"first".to_vec()).await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn publication_barriers_keep_open_handles_when_root_path_changes() {
let (dir, store) = store(0);
let root = store.root().to_path_buf();
let moved = dir.path().join("moved-root");
let hook = sync_hook::install_publication(&root, None);
let moved_for_hook = moved.clone();
hook.on_step("put_sync_shard", move || {
fs::rename(&root, moved_for_hook).unwrap();
std::os::unix::fs::symlink(&root, &root).unwrap();
});
let bytes = b"pinned publication".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
assert_eq!(fs::read(shard_path(&moved, &content_ref)).unwrap(), bytes);
assert_eq!(
*hook.completed.lock().unwrap(),
[
"put_fsync",
"put_persist",
"put_sync_shard",
"put_sync_parent",
"put_sync_root"
]
);
assert!(store.put(bytes).await.is_err());
}
#[cfg(unix)]
#[tokio::test]
async fn publication_barriers_reopen_in_a_fresh_process() {
let (_dir, store) = store(0);
let bytes = b"fresh process publication".to_vec();
let content_ref = store.put(bytes).await.unwrap();
let result = std::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"stores::blob::tests::publication_barriers_process_reader",
"--ignored",
"--nocapture",
])
.env("KHIVE_TEST_BLOB_PUBLICATION_ROOT", store.root())
.env("KHIVE_TEST_BLOB_PUBLICATION_REF", content_ref.as_str())
.output()
.unwrap();
assert!(result.status.success(), "{result:?}");
assert!(String::from_utf8_lossy(&result.stdout).contains("verified published object"));
}
#[cfg(unix)]
#[test]
#[ignore = "subprocess helper for publication_barriers_reopen_in_a_fresh_process"]
fn publication_barriers_process_reader() {
let root = PathBuf::from(std::env::var_os("KHIVE_TEST_BLOB_PUBLICATION_ROOT").unwrap());
let content_ref =
ContentRef::from_hex(std::env::var("KHIVE_TEST_BLOB_PUBLICATION_REF").unwrap())
.unwrap();
let store = FsBlobStore::open_existing(root, 0).unwrap();
let bytes = tokio::runtime::Runtime::new()
.unwrap()
.block_on(store.get_bounded_verified(&content_ref, 1024))
.unwrap();
assert_eq!(bytes, b"fresh process publication");
println!("verified published object");
}
fn prepare_v20_gc_fixture(conn: &mut rusqlite::Connection) {
conn.execute_batch(include_str!("../../sql/schema-migrations-table.sql"))
.expect("create migration ledger");
for migration in crate::MIGRATIONS
.iter()
.filter(|migration| migration.version <= 20)
{
let tx = conn.transaction().expect("begin historical migration");
tx.execute_batch(migration.up)
.expect("apply historical migration body");
tx.execute(
"INSERT INTO _schema_migrations (version, name, applied_at) \
VALUES (?1, ?2, 0)",
rusqlite::params![migration.version, migration.name],
)
.expect("record historical migration");
tx.commit().expect("commit historical migration");
}
}
fn prepare_completed_v21_gc_fixture(conn: &mut rusqlite::Connection) {
prepare_v20_gc_fixture(conn);
crate::migrations::stage_attachment_cutover(conn).expect("stage canonical completed V21");
crate::migrations::finalize_attachment_cutover(conn)
.expect("finalize canonical completed V21");
let version = crate::migrations::read_schema_version(conn)
.expect("read canonical completed V21 ledger");
assert_eq!(
version,
crate::migrations::ATTACHMENT_CUTOVER_VERSION,
"GC-gate fixtures need a completed-V21 ledger; a later migration chain \
must provide a pinned through-V21 fixture builder for these tests"
);
}
#[tokio::test]
async fn completed_v21_gc_gate_requires_new_indexes_and_absent_legacy_column() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
assert!(blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap());
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute_batch("DROP INDEX idx_attachments_content_ref")
.unwrap();
}
assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap());
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE INDEX idx_attachments_content_ref \
ON attachments(content_ref); \
ALTER TABLE entities ADD COLUMN content_ref TEXT",
)
.unwrap();
}
assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap());
}
#[tokio::test]
async fn completed_v21_gate_acceptance_matrix_rejects_each_removed_fact_independently() {
let cases: &[(&str, &str)] = &[
(
"blob_gc_claims_content_ref_index_dropped",
"DROP INDEX idx_blob_gc_claims_content_ref",
),
(
"v21_ledger_row_deleted",
"DELETE FROM _schema_migrations WHERE version = 21",
),
(
"v21_ledger_row_renamed",
"UPDATE _schema_migrations SET name = 'not_attachments_first_class' \
WHERE version = 21",
),
(
"below_v21_ledger_rows_deleted",
"DELETE FROM _schema_migrations WHERE version < 21",
),
(
"marker_row_deleted",
"DELETE FROM attachment_cutover_state WHERE singleton = 1",
),
(
"marker_completed_at_null_while_state_complete",
"DROP TABLE attachment_cutover_state; \
CREATE TABLE attachment_cutover_state ( \
singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
state TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
started_at INTEGER NOT NULL, \
completed_at INTEGER \
) STRICT; \
INSERT INTO attachment_cutover_state \
(singleton, state, started_at, completed_at) \
VALUES (1, 'complete', 21, NULL)",
),
(
"insert_fence_dropped_alone",
"DROP TRIGGER attachments_reject_claimed_blob_insert",
),
];
for (case, mutation_sql) in cases {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute_batch(mutation_sql)
.unwrap_or_else(|e| panic!("case {case}: failed to apply mutation: {e}"));
}
assert!(
!blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap(),
"case {case}: gate must reject with this fact removed"
);
let store = Arc::new(
FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store
.put(format!("gate matrix orphan for {case}").into_bytes())
.await
.unwrap();
let abandoned_ref = "c".repeat(64);
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('gate-matrix-abandoned', ?1, 1)",
[abandoned_ref.as_str()],
)
.unwrap();
}
let _root_guard = store.write_lock.clone().lock_owned().await;
for dry_run in [true, false] {
let outcome = tokio::time::timeout(
Duration::from_secs(1),
store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
)
.await
.unwrap_or_else(|_| panic!("case {case}: refusal must precede the root wait"));
assert!(
matches!(outcome, Err(StorageError::Unsupported { .. })),
"case {case} dry_run={dry_run}: expected Unsupported, got {outcome:?}"
);
}
assert!(
store.exists(&orphan).await.unwrap(),
"case {case}: a refused sweep must not delete anything"
);
let reader = backend.pool().reader().unwrap();
let remaining: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = 'gate-matrix-abandoned'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(remaining, 1, "case {case}: refusal must not recover claims");
let probe_claims: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims WHERE root_key GLOB '__fence_probe-*'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
probe_claims, 0,
"case {case}: refusal must leave no fence-probe claim residue"
);
let probe_attachments: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM attachments \
WHERE record_uuid GLOB '__blob-gc-fence-probe-*'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
probe_attachments, 0,
"case {case}: refusal must leave no fence-probe attachment residue"
);
}
}
#[tokio::test]
async fn gate_rejects_ledger_ahead_of_binary_latest() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
assert!(
blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap(),
"control: completed fixture at the binary's latest version must pass"
);
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO _schema_migrations (version, name, applied_at) \
VALUES (?1, 'post_cutover_feature', unixepoch())",
[i64::from(crate::migrations::latest_schema_version()) + 1],
)
.unwrap();
}
assert!(
!blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap(),
"a ledger ahead of the binary's latest schema version must fail closed"
);
}
#[tokio::test]
async fn gate_rejects_incomplete_ledger_behind_v21() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
assert!(
blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap(),
"control: the untouched completed fixture must pass"
);
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute_batch("DELETE FROM _schema_migrations WHERE version < 21")
.unwrap();
let (v21_named, max_version): (i64, i64) = writer
.conn()
.query_row(
"SELECT (SELECT COUNT(*) FROM _schema_migrations \
WHERE version = 21 AND name = 'attachments_first_class'), \
(SELECT MAX(version) FROM _schema_migrations)",
[],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
assert_eq!(
(v21_named, max_version),
(1, 21),
"fixture must keep the named V21 row and MAX(version) = 21 so the \
rejection can only come from the contiguity clause"
);
}
assert!(
!blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap(),
"an incomplete ledger behind V21 must fail closed despite a valid terminal row"
);
}
#[tokio::test]
async fn completed_v21_gate_fails_closed_when_marker_read_errors() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute_batch(
"DROP TABLE attachment_cutover_state; \
CREATE TABLE attachment_cutover_state ( \
singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
state TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
started_at INTEGER NOT NULL \
) STRICT; \
INSERT INTO attachment_cutover_state (singleton, state, started_at) \
VALUES (1, 'complete', 21)",
)
.unwrap();
}
let gate_error = blob_gc_fencing_complete(backend.sql().as_ref())
.await
.expect_err("a marker read error must propagate, not silently resolve to false");
assert!(
!matches!(gate_error, StorageError::Unsupported { .. }),
"a read error is a distinct failure from the typed epoch refusal: {gate_error:?}"
);
let store = Arc::new(
FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store
.put(b"marker read error orphan".to_vec())
.await
.unwrap();
let _root_guard = store.write_lock.clone().lock_owned().await;
for dry_run in [true, false] {
let outcome = tokio::time::timeout(
Duration::from_secs(1),
store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
)
.await
.expect("a marker read error must fail before waiting on the root lock");
assert!(
outcome.is_err(),
"dry_run={dry_run}: expected the sweep to fail closed, got {outcome:?}"
);
}
assert!(store.exists(&orphan).await.unwrap());
}
#[test]
fn database_sweep_owner_is_keyed_by_database_not_blob_root() {
let dir = tempfile::tempdir().unwrap();
let database = dir.path().join("khive.db");
let same_database_a = sweep_lock_for_database(Some(&database));
let same_database_b = sweep_lock_for_database(Some(&database));
let other_database = sweep_lock_for_database(Some(&dir.path().join("other.db")));
let mut expected_lock_path = database.as_os_str().to_os_string();
expected_lock_path.push(DATABASE_GC_LOCK_SUFFIX);
assert!(Arc::ptr_eq(&same_database_a, &same_database_b));
assert!(!Arc::ptr_eq(&same_database_a, &other_database));
assert_eq!(
database_gc_lock_path(&database),
PathBuf::from(expected_lock_path)
);
}
#[tokio::test]
async fn database_gc_owner_holds_process_and_advisory_fences_until_drop() {
let dir = tempfile::tempdir().unwrap();
let database = dir.path().join("owner.db");
let backend = crate::StorageBackend::sqlite_for_test(&database).unwrap();
let owner = acquire_database_gc_owner(backend.sql().as_ref())
.await
.unwrap();
let canonical_database = owner
.database_path()
.expect("file-backed owner path")
.to_path_buf();
assert!(
sweep_lock_for_database(Some(&canonical_database))
.try_acquire()
.is_none(),
"boot and sweep must share one process-local database owner"
);
let external = fs::OpenOptions::new()
.read(true)
.write(true)
.open(database_gc_lock_path(&canonical_database))
.unwrap();
assert!(
matches!(
fs4::FileExt::try_lock(&external),
Err(fs4::TryLockError::WouldBlock)
),
"the reusable owner must also retain the cross-process advisory fence"
);
drop(owner);
fs4::FileExt::try_lock(&external).expect("owner drop releases advisory fence");
}
#[cfg(unix)]
#[test]
fn database_gc_lock_path_preserves_non_utf8_identity() {
use std::os::unix::ffi::{OsStrExt, OsStringExt};
let database = PathBuf::from(std::ffi::OsString::from_vec(
b"khive-non-utf8-\xff.db".to_vec(),
));
let lock_path = database_gc_lock_path(&database);
let mut expected = database.as_os_str().as_bytes().to_vec();
expected.extend_from_slice(DATABASE_GC_LOCK_SUFFIX.as_bytes());
assert_eq!(lock_path.as_os_str().as_bytes(), expected);
}
#[cfg(unix)]
#[tokio::test]
async fn unlink_blob_shard_refuses_symlinked_shard_dir_and_still_sweeps_real_shard() {
let (dir, store) = store(0);
let root = dir.path().join("blobs");
let real = store.put(b"real blob content".to_vec()).await.unwrap();
let outside = tempfile::tempdir().unwrap();
let victim = outside.path().join("victim.txt");
fs::write(&victim, b"do not delete me").unwrap();
let real_prefix = &real.as_str()[0..2];
let attack_prefix = if real_prefix == "aa" { "bb" } else { "aa" };
let fake_ref = ContentRef::from_hex(format!("{attack_prefix}{}", "0".repeat(62))).unwrap();
let fake_hex = fake_ref.as_str().to_string();
let shard1 = root.join(attack_prefix);
std::os::unix::fs::symlink(outside.path(), &shard1).unwrap();
fs::create_dir_all(outside.path().join(&fake_hex[2..4])).unwrap();
fs::write(
outside.path().join(&fake_hex[2..4]).join(&fake_hex),
b"decoy",
)
.unwrap();
let root_handle = open_blob_root_handle(&root).unwrap();
let error = unlink_blob_shard_file_no_follow(&root, &root_handle, &fake_ref).unwrap_err();
assert!(
matches!(
error.raw_os_error(),
Some(libc::ELOOP) | Some(libc::ENOTDIR)
),
"opening a symlinked shard directory must be refused, not followed; got: {error}"
);
assert!(
victim.exists(),
"the file outside the blob root must never be touched by a refused shard-dir open"
);
unlink_blob_shard_file_no_follow(&root, &root_handle, &real).unwrap();
assert!(!store.exists(&real).await.unwrap());
}
async fn recv_blocking(rx: std::sync::mpsc::Receiver<()>) -> bool {
tokio::task::spawn_blocking(move || rx.recv().is_ok())
.await
.expect("recv_blocking thread panicked")
}
#[tokio::test]
async fn put_bounded_get_roundtrip() {
let (_dir, store) = store(0);
let bytes = b"hello blob store".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
let fetched = store
.get_bounded_verified(&content_ref, bytes.len() as u64)
.await
.unwrap();
assert_eq!(fetched, bytes);
}
#[tokio::test]
async fn bounded_verified_get_accepts_exact_and_portable_maximum_limits() {
let (_dir, store) = store(0);
let bytes = b"bounded fs blob".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
assert_eq!(
store
.get_bounded_verified(&content_ref, bytes.len() as u64)
.await
.unwrap(),
bytes
);
assert_eq!(
store
.get_bounded_verified(&content_ref, MAX_BLOB_WHOLE_BYTES)
.await
.unwrap(),
bytes
);
}
#[cfg(unix)]
#[tokio::test]
async fn bounded_verified_get_resolves_a_configured_symlink_root_once() {
use std::os::unix::fs::symlink;
let dir = tempfile::tempdir().unwrap();
let target = dir.path().join("blob-target");
fs::create_dir(&target).unwrap();
let configured = dir.path().join("blob-configured");
symlink(&target, &configured).unwrap();
let store = FsBlobStore::new(configured, 0).unwrap();
assert_eq!(store.root(), target.canonicalize().unwrap());
let bytes = b"symlink-configured root".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
assert_eq!(
store
.get_bounded_verified(&content_ref, bytes.len() as u64)
.await
.unwrap(),
bytes
);
}
#[cfg(unix)]
#[tokio::test]
async fn fs_blob_store_refuses_root_replacement_before_put() {
use std::os::unix::fs::symlink;
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
let store = FsBlobStore::new(root.clone(), 0).unwrap();
let original_root = dir.path().join("blobs-original");
fs::rename(&root, &original_root).unwrap();
let redirected_root = tempfile::tempdir().unwrap();
symlink(redirected_root.path(), &root).unwrap();
let bytes = b"must not land in a replaced root".to_vec();
let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
store
.put(bytes)
.await
.expect_err("a root replaced after construction must be refused");
assert!(
!shard_path(redirected_root.path(), &content_ref).exists(),
"put must not publish into the tree selected by the replacement symlink"
);
assert!(
!shard_path(&original_root, &content_ref).exists(),
"a refused put must not mutate the initialization-time root either"
);
}
#[cfg(unix)]
#[tokio::test]
async fn fs_blob_store_refuses_ancestor_replacement_for_existing_blob_operations() {
use std::os::unix::fs::symlink;
let dir = tempfile::tempdir().unwrap();
let ancestor = dir.path().join("store-parent");
let root = ancestor.join("blobs");
fs::create_dir(&ancestor).unwrap();
let store = FsBlobStore::new(root.clone(), 0).unwrap();
let bytes = b"same bytes in both trees".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
let original_ancestor = dir.path().join("store-parent-original");
fs::rename(&ancestor, &original_ancestor).unwrap();
let original_blob = shard_path(&original_ancestor.join("blobs"), &content_ref);
let redirected_ancestor = tempfile::tempdir().unwrap();
let redirected_blob = shard_path(&redirected_ancestor.path().join("blobs"), &content_ref);
fs::create_dir_all(redirected_blob.parent().unwrap()).unwrap();
fs::write(&redirected_blob, &bytes).unwrap();
symlink(redirected_ancestor.path(), &ancestor).unwrap();
store
.get_bounded_verified(&content_ref, bytes.len() as u64)
.await
.expect_err("a read through a replaced root ancestor must be refused");
store
.exists(&content_ref)
.await
.expect_err("exists through a replaced root ancestor must be refused");
store
.size(&content_ref)
.await
.expect_err("size through a replaced root ancestor must be refused");
store
.delete(&content_ref)
.await
.expect_err("delete through a replaced root ancestor must be refused");
assert!(
original_blob.exists(),
"refusal must preserve the initialization-time blob"
);
assert!(
redirected_blob.exists(),
"refusal must not read as authority or delete the redirected blob"
);
}
#[tokio::test]
async fn bounded_verified_get_rejects_same_size_digest_corruption() {
let (_dir, store) = store(0);
let expected_bytes = b"expected".to_vec();
let actual_bytes = b"mutated!".to_vec();
let expected = store.put(expected_bytes).await.unwrap();
let actual = ContentRef::from_digest_bytes(blake3::hash(&actual_bytes).as_bytes());
fs::write(shard_path(store.root(), &expected), actual_bytes).unwrap();
let err = store.get_bounded_verified(&expected, 8).await.unwrap_err();
assert!(matches!(
err,
StorageError::BlobDigestMismatch {
expected: ref got_expected,
actual: ref got_actual,
} if got_expected == &expected && got_actual == &actual
));
}
#[tokio::test]
async fn bounded_verified_get_stops_at_max_plus_one_after_file_growth() {
let (_dir, store) = store(0);
let store = Arc::new(store);
let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
let path = shard_path(store.root(), &content_ref);
let (reached, release) = bounded_read_sync_hook::install(store.root());
let read_store = Arc::clone(&store);
let read_ref = content_ref.clone();
let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
assert!(
recv_blocking(reached).await,
"read must reach the metadata seam"
);
let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
writer.write_all(b"efgh-poison-tail").unwrap();
writer.flush().unwrap();
release.send(()).unwrap();
let err = read.await.unwrap().unwrap_err();
assert!(matches!(
err,
StorageError::BlobTooLarge {
content_ref: ref got,
max_bytes: 4,
observed_at_least: 5,
} if got == &content_ref
));
}
#[tokio::test]
async fn bounded_verified_get_reports_growth_within_limit_as_size_mismatch() {
let (_dir, store) = store(0);
let store = Arc::new(store);
let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
let path = shard_path(store.root(), &content_ref);
let (reached, release) = bounded_read_sync_hook::install(store.root());
let read_store = Arc::clone(&store);
let read_ref = content_ref.clone();
let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 8).await });
assert!(
recv_blocking(reached).await,
"read must reach the metadata seam"
);
let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
writer.write_all(b"ef").unwrap();
writer.flush().unwrap();
release.send(()).unwrap();
let err = read.await.unwrap().unwrap_err();
assert!(matches!(
err,
StorageError::BlobSizeMismatch {
content_ref: ref got,
metadata_bytes: 4,
actual_bytes: 6,
} if got == &content_ref
));
}
#[tokio::test]
async fn bounded_verified_get_reports_truncation_as_size_mismatch() {
let (_dir, store) = store(0);
let store = Arc::new(store);
let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
let path = shard_path(store.root(), &content_ref);
let (reached, release) = bounded_read_sync_hook::install(store.root());
let read_store = Arc::clone(&store);
let read_ref = content_ref.clone();
let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
assert!(
recv_blocking(reached).await,
"read must reach the metadata seam"
);
let mut writer = fs::OpenOptions::new()
.write(true)
.truncate(true)
.open(&path)
.unwrap();
writer.write_all(b"abc").unwrap();
writer.flush().unwrap();
release.send(()).unwrap();
let err = read.await.unwrap().unwrap_err();
assert!(matches!(
err,
StorageError::BlobSizeMismatch {
content_ref: ref got,
metadata_bytes: 4,
actual_bytes: 3,
} if got == &content_ref
));
}
#[cfg(unix)]
#[tokio::test]
async fn bounded_verified_get_keeps_the_opened_inode_when_the_path_is_replaced() {
let (_dir, store) = store(0);
let store = Arc::new(store);
let original = b"original".to_vec();
let replacement = b"replaced".to_vec();
let content_ref = store.put(original.clone()).await.unwrap();
let path = shard_path(store.root(), &content_ref);
let moved_path = path.with_extension("opened-inode");
let max_bytes = original.len() as u64;
let (reached, release) = bounded_read_sync_hook::install(store.root());
let read_store = Arc::clone(&store);
let read_ref = content_ref.clone();
let read =
tokio::spawn(
async move { read_store.get_bounded_verified(&read_ref, max_bytes).await },
);
assert!(
recv_blocking(reached).await,
"read must reach the metadata seam"
);
fs::rename(&path, &moved_path).unwrap();
fs::write(&path, replacement).unwrap();
release.send(()).unwrap();
assert_eq!(read.await.unwrap().unwrap(), original);
}
#[cfg(unix)]
#[tokio::test]
async fn bounded_verified_get_refuses_a_symlink_leaf() {
use std::os::unix::fs::symlink;
let dir = tempfile::tempdir().unwrap();
let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
let outside_bytes = b"outside but digest matching".to_vec();
let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
let outside = dir.path().join("outside");
fs::write(&outside, &outside_bytes).unwrap();
let leaf = shard_path(store.root(), &content_ref);
fs::create_dir_all(leaf.parent().unwrap()).unwrap();
symlink(&outside, &leaf).unwrap();
let err = store
.get_bounded_verified(&content_ref, outside_bytes.len() as u64)
.await
.unwrap_err();
assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
}
#[cfg(unix)]
#[tokio::test]
async fn bounded_verified_get_refuses_a_symlinked_shard_component() {
use std::os::unix::fs::symlink;
let dir = tempfile::tempdir().unwrap();
let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
let outside_bytes = b"outside through shard link".to_vec();
let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
let hex = content_ref.as_str();
let outside_shard1 = dir.path().join("outside-shard1");
let outside_shard2 = outside_shard1.join(&hex[2..4]);
fs::create_dir_all(&outside_shard2).unwrap();
fs::write(outside_shard2.join(hex), &outside_bytes).unwrap();
symlink(&outside_shard1, store.root().join(&hex[0..2])).unwrap();
let err = store
.get_bounded_verified(&content_ref, outside_bytes.len() as u64)
.await
.unwrap_err();
assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
}
#[tokio::test]
async fn put_content_ref_matches_blake3_digest() {
let (_dir, store) = store(0);
let bytes = b"digest check".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
let expected = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
assert_eq!(content_ref, expected);
}
#[tokio::test]
async fn put_dedups_identical_content() {
let (_dir, store) = store(0);
let bytes = b"same bytes twice".to_vec();
let first = store.put(bytes.clone()).await.unwrap();
let second = store.put(bytes.clone()).await.unwrap();
assert_eq!(first, second);
assert_eq!(
store
.get_bounded_verified(&first, bytes.len() as u64)
.await
.unwrap(),
bytes
);
}
#[tokio::test]
async fn exists_reflects_put_and_delete() {
let (_dir, store) = store(0);
let bytes = b"exists check".to_vec();
let content_ref = store.put(bytes).await.unwrap();
assert!(store.exists(&content_ref).await.unwrap());
assert!(store.delete(&content_ref).await.unwrap());
assert!(!store.exists(&content_ref).await.unwrap());
}
#[tokio::test]
async fn delete_missing_content_ref_returns_false() {
let (_dir, store) = store(0);
let missing = ContentRef::from_hex("f".repeat(64)).unwrap();
assert!(!store.delete(&missing).await.unwrap());
}
#[tokio::test]
async fn size_reports_byte_length_for_a_present_object() {
let (_dir, store) = store(0);
let bytes = b"size check".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
assert_eq!(
store.size(&content_ref).await.unwrap(),
Some(bytes.len() as u64)
);
}
#[tokio::test]
async fn size_returns_none_for_an_absent_object() {
let (_dir, store) = store(0);
let missing = ContentRef::from_hex("9".repeat(64)).unwrap();
assert_eq!(store.size(&missing).await.unwrap(), None);
}
#[tokio::test]
async fn bounded_get_missing_content_ref_returns_not_found() {
let (_dir, store) = store(0);
let missing = ContentRef::from_hex("e".repeat(64)).unwrap();
let err = store
.get_bounded_verified(&missing, MAX_BLOB_WHOLE_BYTES)
.await
.unwrap_err();
assert!(matches!(err, StorageError::NotFound { .. }));
}
#[tokio::test]
async fn put_refuses_below_free_space_floor() {
let (_dir, store) = store(u64::MAX);
let err = store.put(b"too big a floor".to_vec()).await.unwrap_err();
match err {
StorageError::CapacityFloor {
floor_bytes,
available_bytes,
..
} => {
assert_eq!(floor_bytes, u64::MAX);
assert!(available_bytes < u64::MAX);
}
other => panic!("expected CapacityFloor, got {other:?}"),
}
}
#[tokio::test]
async fn capacity_floor_error_names_the_floor_and_volume() {
let (_dir, store) = store(u64::MAX);
let err = store.put(b"x".to_vec()).await.unwrap_err();
let msg = err.to_string();
assert!(
msg.contains(&u64::MAX.to_string()),
"must name the floor: {msg}"
);
assert!(msg.contains("Blob"), "must name the capability: {msg}");
}
#[test]
fn crosses_floor_is_write_size_aware_at_the_exact_boundary() {
assert!(crosses_floor(101, 2, 100));
assert!(!crosses_floor(101, 1, 100));
}
#[test]
fn crosses_floor_accepts_a_write_that_lands_exactly_on_the_floor() {
assert!(!crosses_floor(100, 0, 100));
}
#[test]
fn crosses_floor_rejects_a_write_that_lands_one_byte_under_the_floor() {
assert!(crosses_floor(100, 1, 100));
}
#[test]
fn crosses_floor_saturates_instead_of_underflowing_when_write_exceeds_available() {
assert!(crosses_floor(10, 100, 50));
assert!(!crosses_floor(10, 100, 0));
}
#[test]
fn put_refuses_a_write_that_would_cross_the_floor_even_though_available_alone_clears_it() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
fs::create_dir_all(&root).unwrap();
let err = put_blocking_with_space_probe(&root, 100, vec![7u8; 2], |_| Ok(101)).unwrap_err();
assert!(
matches!(err, StorageError::CapacityFloor { .. }),
"a write-size-aware floor check must reject a write that pushes the volume \
below the floor even though available space alone still clears it: {err:?}"
);
}
#[test]
fn a_later_put_checks_a_fresh_capacity_snapshot() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
fs::create_dir_all(&root).unwrap();
let first = put_blocking_with_space_probe(&root, 100, vec![1u8; 2], |_| Ok(102));
let second = put_blocking_with_space_probe(&root, 100, vec![2u8; 2], |_| Ok(101));
assert!(
first.is_ok(),
"the first put may land on the floor: {first:?}"
);
assert!(
matches!(second, Err(StorageError::CapacityFloor { .. })),
"the later put must use its lower capacity snapshot: {second:?}"
);
}
#[tokio::test]
async fn concurrent_puts_from_two_independently_constructed_stores_share_the_root_lock() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
fs::create_dir_all(&root).unwrap();
let canonical_root = root.canonicalize().unwrap();
let store_a = std::sync::Arc::new(FsBlobStore::new(root.clone(), 0).unwrap());
let store_b = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
let (a_reached, a_release, _a_done) = sync_hook::install(&canonical_root);
let a = {
let store_a = store_a.clone();
tokio::spawn(async move { store_a.put(b"store_a payload".to_vec()).await })
};
assert!(
recv_blocking(a_reached).await,
"store_a's put must reach the sync_hook checkpoint"
);
assert!(
store_b.write_lock.try_lock().is_err(),
"store_b's write_lock was NOT held while store_a's put held its guard -- the two \
independently constructed stores do NOT share one lock"
);
a_release.send(()).unwrap();
let result_a = a.await.unwrap();
assert!(result_a.is_ok(), "store_a's put must succeed: {result_a:?}");
let result_b = store_b.put(b"store_b payload".to_vec()).await;
assert!(result_b.is_ok(), "store_b's put must succeed: {result_b:?}");
}
#[tokio::test]
async fn aborting_the_outer_put_future_does_not_release_the_guard_before_persist_completes() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("blobs");
fs::create_dir_all(&root).unwrap();
let canonical_root = root.canonicalize().unwrap();
let store = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
let (reached, release, done) = sync_hook::install(&canonical_root);
let handle = {
let store = store.clone();
tokio::spawn(async move { store.put(b"cancellation race payload".to_vec()).await })
};
assert!(
recv_blocking(reached).await,
"put must reach the sync_hook checkpoint -- owned guard already moved into the \
closure -- before this test can mean anything"
);
handle.abort();
let abort_result = handle.await;
match &abort_result {
Err(e) if e.is_cancelled() => {}
other => panic!(
"the outer task must actually have been cancelled for this test to be \
meaningful: {other:?}"
),
}
let shared_lock = write_lock_for_root(&canonical_root).unwrap();
assert!(
shared_lock.try_lock().is_err(),
"the guard must still be held by the detached blocking write immediately after \
the outer future was cancelled -- if this is free, the guard was released with \
the aborted frame instead of moving into the spawn_blocking closure"
);
release.send(()).unwrap();
assert!(
recv_blocking(done).await,
"the detached write must signal completion once it actually persists"
);
assert!(
shared_lock.try_lock().is_ok(),
"the guard must be free once the detached write's completion was observed"
);
}
#[tokio::test]
async fn orphan_sweep_is_disabled_in_both_modes_regardless_of_live_refs() {
let (_dir, store) = store(0);
let blob = store
.put(b"never swept by this API".to_vec())
.await
.unwrap();
let mut live_refs = std::collections::HashSet::new();
live_refs.insert(blob.clone());
for dry_run in [true, false] {
let error = store
.orphan_sweep(&BlobOrphanSweepConfig {
live_refs: live_refs.clone(),
dry_run,
})
.await
.expect_err("caller-snapshot orphan_sweep must be disabled");
assert!(
matches!(error, StorageError::Unsupported { .. }),
"expected typed Unsupported, got {error:?}"
);
}
assert!(store.exists(&blob).await.unwrap());
}
#[tokio::test]
async fn transactional_orphan_sweep_refuses_v20_before_root_or_claim_mutation() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_v20_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = Arc::new(
FsBlobStore::new(root, 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let bundle = store.put(b"legacy model bundle".to_vec()).await.unwrap();
let network = store.put(b"legacy FANN network".to_vec()).await.unwrap();
let orphan = store.put(b"ordinary old orphan".to_vec()).await.unwrap();
let abandoned_ref = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO entities \
(id, namespace, kind, entity_type, name, tags, created_at, updated_at, \
content_ref) \
VALUES ('legacy-model', 'local', 'artifact', 'moodboard_model', \
'legacy model', '[]', 1, 1, ?1)",
[bundle.as_str()],
)
.unwrap();
writer
.conn()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('abandoned-before-compat', ?1, 1)",
[abandoned_ref],
)
.unwrap();
}
let _root_guard = store.write_lock.clone().lock_owned().await;
for dry_run in [true, false] {
let outcome = tokio::time::timeout(
Duration::from_secs(1),
store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
)
.await
.expect("V20 refusal must happen before waiting for the held root lock");
let error = outcome.expect_err("V20 transactional sweep must be disabled");
match error {
StorageError::Unsupported {
capability: StorageCapability::Blob,
operation,
message,
} => {
assert_eq!(operation, "transactional_orphan_sweep");
assert!(
message.contains("complete V21 attachment cutover"),
"unexpected compatibility diagnostic: {message}"
);
}
other => panic!("expected typed Unsupported refusal, got {other:?}"),
}
}
let reader = backend.pool().reader().unwrap();
let abandoned: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims \
WHERE root_key = 'abandoned-before-compat' AND content_ref = ?1",
[abandoned_ref],
|row| row.get(0),
)
.unwrap();
assert_eq!(abandoned, 1, "V20 refusal must not clean abandoned claims");
drop(reader);
assert!(store.exists(&bundle).await.unwrap());
assert!(store.exists(&network).await.unwrap());
assert!(store.exists(&orphan).await.unwrap());
}
#[tokio::test]
async fn transactional_orphan_sweep_refuses_incomplete_v21_marker_without_mutation() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
let abandoned_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute(
"UPDATE attachment_cutover_state \
SET state = 'incomplete', completed_at = NULL \
WHERE singleton = 1",
[],
)
.unwrap();
writer
.conn_mut()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('abandoned-incomplete-v21', ?1, 1)",
[abandoned_ref],
)
.unwrap();
}
let store = Arc::new(
FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store.put(b"incomplete V21 orphan".to_vec()).await.unwrap();
let _root_guard = store.write_lock.clone().lock_owned().await;
for dry_run in [true, false] {
let outcome = tokio::time::timeout(
Duration::from_secs(1),
store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
)
.await
.expect("incomplete V21 must refuse before waiting for the root lock");
assert!(
matches!(outcome, Err(StorageError::Unsupported { .. })),
"incomplete V21 must return typed Unsupported: {outcome:?}"
);
}
let remaining: i64 = backend
.pool()
.reader()
.unwrap()
.conn()
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims \
WHERE root_key = 'abandoned-incomplete-v21' AND content_ref = ?1",
[abandoned_ref],
|row| row.get(0),
)
.unwrap();
assert_eq!(remaining, 1, "refusal must not recover abandoned claims");
assert!(store.exists(&orphan).await.unwrap());
}
#[tokio::test]
async fn both_sweep_apis_refuse_v20_and_incomplete_v21_epochs_in_both_modes() {
for incomplete_v21 in [false, true] {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
if incomplete_v21 {
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute(
"UPDATE attachment_cutover_state \
SET state = 'incomplete', completed_at = NULL \
WHERE singleton = 1",
[],
)
.unwrap();
} else {
prepare_v20_gc_fixture(writer.conn_mut());
}
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO);
let orphan = store
.put(format!("both-apis orphan (incomplete_v21={incomplete_v21})").into_bytes())
.await
.unwrap();
let known_claim_ref = "f".repeat(64);
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('both-api-known-claim', ?1, 1)",
[known_claim_ref.as_str()],
)
.unwrap();
}
let assert_known_claim_unchanged = |arm: &str| {
let remaining: i64 = backend
.pool()
.reader()
.unwrap()
.conn()
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims \
WHERE root_key = 'both-api-known-claim' AND content_ref = ?1",
[known_claim_ref.as_str()],
|row| row.get(0),
)
.unwrap();
assert_eq!(
remaining, 1,
"incomplete_v21={incomplete_v21} arm={arm}: refusal must not mutate \
an existing claim"
);
};
for dry_run in [true, false] {
let snapshot_error = store
.orphan_sweep(&BlobOrphanSweepConfig {
live_refs: std::collections::HashSet::new(),
dry_run,
})
.await
.expect_err("orphan_sweep must refuse regardless of epoch");
assert!(
matches!(snapshot_error, StorageError::Unsupported { .. }),
"incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
from orphan_sweep, got {snapshot_error:?}"
);
assert_known_claim_unchanged(&format!("orphan_sweep dry_run={dry_run}"));
let transactional_error = store
.transactional_orphan_sweep(backend.sql().as_ref(), dry_run)
.await
.expect_err("transactional_orphan_sweep must refuse this epoch");
assert!(
matches!(transactional_error, StorageError::Unsupported { .. }),
"incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
from transactional_orphan_sweep, got {transactional_error:?}"
);
assert_known_claim_unchanged(&format!(
"transactional_orphan_sweep dry_run={dry_run}"
));
}
assert!(
store.exists(&orphan).await.unwrap(),
"incomplete_v21={incomplete_v21}: a refused sweep must not delete anything"
);
}
}
#[tokio::test]
async fn transactional_orphan_sweep_recheck_refuses_before_root_lock_when_epoch_regresses_after_db_ownership(
) {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
assert!(blob_gc_fencing_complete(backend.sql().as_ref())
.await
.unwrap());
let blob_root = dir.path().join("blobs");
let store = Arc::new(
FsBlobStore::new(blob_root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store
.put(b"epoch regressed after database ownership".to_vec())
.await
.unwrap();
let canonical_root = blob_root.canonicalize().unwrap();
let _root_write_guard = acquire_root_write_lock(&canonical_root).unwrap();
let canonical_db_path = backend.sql().database_path();
let (reached, release) = db_ownership_sync_hook::install(canonical_db_path.as_deref());
let sweep_store = store.clone();
let sweep_backend = backend.clone();
let handle = tokio::spawn(async move {
sweep_store
.transactional_orphan_sweep(sweep_backend.sql().as_ref(), false)
.await
});
let reached_signal = tokio::time::timeout(Duration::from_secs(1), recv_blocking(reached))
.await
.expect("the sweep must reach database ownership before this test's timeout");
assert!(reached_signal, "hook sender was dropped before signaling");
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute("DELETE FROM _schema_migrations WHERE version = 21", [])
.unwrap();
}
release.send(()).unwrap();
let outcome = tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect(
"the recheck must refuse before ever waiting on the externally held root lock -- \
under the old (pre-fix) ordering this join times out instead, because the \
sweep blocks acquiring the OS-level root lock held above",
)
.unwrap();
assert!(
matches!(outcome, Err(StorageError::Unsupported { .. })),
"expected the regressed epoch to be caught immediately after database ownership: \
{outcome:?}"
);
assert!(store.exists(&orphan).await.unwrap());
}
#[tokio::test]
async fn transactional_orphan_sweep_accepts_completed_v21_attachment_liveness() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO);
let bundle = store.put(b"V21 model bundle".to_vec()).await.unwrap();
let network = store.put(b"V21 FANN network".to_vec()).await.unwrap();
let orphan = store.put(b"V21 true orphan".to_vec()).await.unwrap();
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO entities \
(id, namespace, kind, entity_type, name, tags, created_at, updated_at) \
VALUES ('model', 'local', 'artifact', 'moodboard_model', \
'model', '[]', 1, 1)",
[],
)
.unwrap();
writer
.conn()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('model', 'entity', 'content', ?1, 1), \
('model', 'entity', 'fann-network', ?2, 1)",
rusqlite::params![bundle.as_str(), network.as_str()],
)
.unwrap();
}
let dry_run = store
.transactional_orphan_sweep(backend.sql().as_ref(), true)
.await
.expect("completed V21 dry run must be supported");
assert_eq!(dry_run.would_delete, 1);
assert_eq!(dry_run.deleted, 0);
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.expect("completed V21 destructive sweep must be supported");
assert_eq!(result.deleted, 1);
assert!(store.exists(&bundle).await.unwrap());
assert!(store.exists(&network).await.unwrap());
assert!(!store.exists(&orphan).await.unwrap());
}
#[tokio::test]
async fn transactional_orphan_sweep_refuses_without_the_blob_gc_claims_migration() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
backend.entities().unwrap();
{
let reader = backend.pool().reader().unwrap();
let present: bool = reader
.conn()
.query_row(
"SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = 'table' \
AND name = 'blob_gc_claims'",
[],
|row| row.get(0),
)
.unwrap();
assert!(
!present,
"this test's premise requires blob_gc_claims to be absent"
);
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store.put(b"direct-backend orphan".to_vec()).await.unwrap();
let sql = backend.sql();
let error = store
.transactional_orphan_sweep(sql.as_ref(), false)
.await
.expect_err("sweep must refuse a backend without the blob_gc_claims fencing set");
assert!(
matches!(error, StorageError::Unsupported { .. }),
"expected StorageError::Unsupported, got {error:?}"
);
assert!(
store.exists(&orphan).await.unwrap(),
"a refused sweep must not have deleted anything"
);
}
#[tokio::test]
async fn transactional_orphan_sweep_refuses_an_incomplete_cutover_marker() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute_batch(
"UPDATE attachment_cutover_state \
SET state = 'incomplete', completed_at = NULL WHERE singleton = 1; \
DELETE FROM _schema_migrations WHERE version = 21;",
)
.unwrap();
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO);
let orphan = store
.put(b"incomplete-cutover orphan".to_vec())
.await
.unwrap();
let error = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.expect_err("sweep must refuse every durable incomplete marker");
assert!(matches!(error, StorageError::Unsupported { .. }));
assert!(
store.exists(&orphan).await.unwrap(),
"refused incomplete-state sweep must preserve every blob"
);
}
#[tokio::test]
async fn transactional_orphan_sweep_refuses_with_incomplete_fencing_triggers() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute_batch("DROP TRIGGER attachments_reject_claimed_blob_update")
.unwrap();
writer
.conn_mut()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('abandoned-partial-fence', \
'dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd', \
1)",
[],
)
.unwrap();
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store.put(b"partial-fence orphan".to_vec()).await.unwrap();
let sql = backend.sql();
let _root_guard = store.write_lock.clone().lock_owned().await;
let error = tokio::time::timeout(
Duration::from_secs(1),
store.transactional_orphan_sweep(sql.as_ref(), false),
)
.await
.expect("an incomplete V21 fence must refuse before the root wait")
.expect_err("sweep must refuse when any V21 fencing trigger is missing");
assert!(
matches!(error, StorageError::Unsupported { .. }),
"expected StorageError::Unsupported, got {error:?}"
);
assert!(
store.exists(&orphan).await.unwrap(),
"a refused sweep must not have deleted anything"
);
let remaining: i64 = backend
.pool()
.reader()
.unwrap()
.conn()
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims \
WHERE root_key = 'abandoned-partial-fence'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(remaining, 1, "a refused sweep must not recover claims");
}
#[tokio::test]
async fn transactional_orphan_sweep_refuses_same_named_noop_fencing_triggers() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute_batch(
"DROP TRIGGER attachments_reject_claimed_blob_insert; \
DROP TRIGGER attachments_reject_claimed_blob_update; \
CREATE TRIGGER attachments_reject_claimed_blob_insert \
BEFORE INSERT ON attachments BEGIN SELECT 0; END; \
CREATE TRIGGER attachments_reject_claimed_blob_update \
BEFORE UPDATE OF content_ref ON attachments \
BEGIN SELECT 0; END;",
)
.unwrap();
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store.put(b"noop-trigger orphan".to_vec()).await.unwrap();
let sql = backend.sql();
let error = store
.transactional_orphan_sweep(sql.as_ref(), false)
.await
.expect_err("sweep must refuse when the fencing triggers are same-named no-ops");
assert!(
matches!(error, StorageError::Unsupported { .. }),
"expected StorageError::Unsupported, got {error:?}"
);
assert!(
store.exists(&orphan).await.unwrap(),
"a refused sweep must not have deleted anything"
);
let reader = backend.pool().reader().unwrap();
let leftovers: i64 = reader
.conn()
.query_row(
"SELECT (SELECT COUNT(*) FROM blob_gc_claims \
WHERE root_key GLOB '__fence_probe-*') \
+ (SELECT COUNT(*) FROM attachments \
WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(leftovers, 0, "fence probe rows must not survive the probe");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fence_probe_refuses_id_collision_and_preserves_the_colliding_attachment() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, media_type, created_at) \
VALUES ('victim-id', 'entity', 'content', \
'2222222222222222222222222222222222222222222222222222222222222222', \
'application/test', 7)",
[],
)
.unwrap();
}
let sql = backend.sql();
let error = super::blob_gc_fence_probe_with_ids(
sql.as_ref(),
"victim-id".to_string(),
"victim-update-id".to_string(),
"victim-insert2-id".to_string(),
"victim-update2-id".to_string(),
"victim-claim-key".to_string(),
)
.await
.expect_err("the probe must refuse when an id it would delete already names a row");
assert!(
matches!(
&error,
StorageError::WriterTaskRequestFailed {
request_state:
khive_storage::WriterTaskRequestState::TransactionRolledBack,
source,
} if matches!(source.as_ref(), StorageError::Unsupported { .. })
),
"expected a proven-rollback wrapper retaining StorageError::Unsupported, got {error:?}"
);
let reader = backend.pool().reader().unwrap();
let (media_type, created_at): (String, i64) = reader
.conn()
.query_row(
"SELECT media_type, created_at FROM attachments \
WHERE record_uuid = 'victim-id' AND role = 'content'",
[],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.expect("the colliding attachment must survive the refused probe untouched");
assert_eq!(media_type, "application/test");
assert_eq!(created_at, 7);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fence_probe_does_not_touch_an_unrelated_retained_entity_sequence() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute(
"INSERT INTO entities \
(id, namespace, kind, name, tags, created_at, updated_at) \
VALUES ('retained-id', 'local', 'document', 'gone entity', '[]', 7, 7)",
[],
)
.unwrap();
writer
.conn_mut()
.execute("DELETE FROM entities WHERE id = 'retained-id'", [])
.unwrap();
let retained: i64 = writer
.conn_mut()
.query_row(
"SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(retained, 1, "fixture requires a retained-only ledger row");
}
let sql = backend.sql();
super::blob_gc_fence_probe_with_ids(
sql.as_ref(),
"retained-id".to_string(),
"retained-update-id".to_string(),
"retained-insert2-id".to_string(),
"retained-update2-id".to_string(),
"retained-claim-key".to_string(),
)
.await
.expect("attachment probe has no reason to mutate an entity sequence row");
let reader = backend.pool().reader().unwrap();
let survivors: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
survivors, 1,
"the retained entity ledger row must survive the attachment probe"
);
}
fn nul_embedded_canonical_ref() -> String {
let mut polluted = "a".repeat(64);
polluted.push('\0');
polluted.push_str("zz");
polluted
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('nul-claim-key', ?1, 0)",
rusqlite::params![nul_embedded_canonical_ref()],
)
.unwrap();
}
let sql = backend.sql();
let error = super::validate_blob_gc_evidence(sql.as_ref())
.await
.expect_err("a NUL-embedded claim ref must refuse the sweep");
assert!(
error.to_string().contains("blob_gc_claims"),
"expected the claims-table refusal, got {error:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn blob_gc_evidence_rejects_a_nul_embedded_attachment_ref() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute_batch("PRAGMA ignore_check_constraints = ON")
.unwrap();
writer
.conn_mut()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('nul-attachment-id', 'entity', 'content', ?1, 0)",
rusqlite::params![nul_embedded_canonical_ref()],
)
.expect("ignore_check_constraints must allow the corrupt row to insert");
writer
.conn_mut()
.execute_batch("PRAGMA ignore_check_constraints = OFF")
.unwrap();
}
let sql = backend.sql();
let error = super::validate_blob_gc_evidence(sql.as_ref())
.await
.expect_err("a NUL-embedded attachment ref must refuse the sweep");
assert!(
error.to_string().contains("attachments"),
"expected the attachments-table refusal, got {error:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fence_probe_refuses_a_digest_restricted_trigger_rewrite() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let sql = backend.sql();
super::blob_gc_fence_probe(sql.as_ref())
.await
.expect("the healthy fence must pass all four probe arms");
{
let mut writer = backend.pool().writer().unwrap();
writer
.conn_mut()
.execute_batch(
"DROP TRIGGER attachments_reject_claimed_blob_insert; \
DROP TRIGGER attachments_reject_claimed_blob_update; \
CREATE TRIGGER attachments_reject_claimed_blob_insert \
BEFORE INSERT ON attachments \
WHEN NEW.content_ref = \
'0000000000000000000000000000000000000000000000000000000000000000' \
AND EXISTS (SELECT 1 FROM blob_gc_claims \
WHERE content_ref = NEW.content_ref) \
BEGIN \
SELECT RAISE(ABORT, \
'content_ref is reserved by an active blob sweep'); \
END; \
CREATE TRIGGER attachments_reject_claimed_blob_update \
BEFORE UPDATE OF content_ref ON attachments \
WHEN NEW.content_ref = \
'0000000000000000000000000000000000000000000000000000000000000000' \
AND EXISTS (SELECT 1 FROM blob_gc_claims \
WHERE content_ref = NEW.content_ref) \
BEGIN \
SELECT RAISE(ABORT, \
'content_ref is reserved by an active blob sweep'); \
END;",
)
.unwrap();
}
let error = super::blob_gc_fence_probe(sql.as_ref())
.await
.expect_err("a sentinel-only fence must fail the second-digest arms");
assert!(
matches!(error, StorageError::Unsupported { .. }),
"expected StorageError::Unsupported, got {error:?}"
);
assert!(
error.to_string().contains("second-digest"),
"the refusal must name a second-digest arm, got {error}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn fence_probe_refuses_a_shape_restricted_trigger_rewrite() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute_batch(
"DROP TRIGGER attachments_reject_claimed_blob_insert; \
DROP TRIGGER attachments_reject_claimed_blob_update; \
CREATE TRIGGER attachments_reject_claimed_blob_insert \
BEFORE INSERT ON attachments \
WHEN NEW.substrate = 'entity' \
AND EXISTS (SELECT 1 FROM blob_gc_claims \
WHERE content_ref = NEW.content_ref) \
BEGIN \
SELECT RAISE(ABORT, \
'content_ref is reserved by an active blob sweep'); \
END; \
CREATE TRIGGER attachments_reject_claimed_blob_update \
BEFORE UPDATE OF content_ref ON attachments \
WHEN NEW.substrate = 'entity' \
AND EXISTS (SELECT 1 FROM blob_gc_claims \
WHERE content_ref = NEW.content_ref) \
BEGIN \
SELECT RAISE(ABORT, \
'content_ref is reserved by an active blob sweep'); \
END;",
)
.unwrap();
}
let sql = backend.sql();
let error = super::blob_gc_fence_probe(sql.as_ref())
.await
.expect_err("an entity-shape-only fence must fail the note-shaped arms");
assert!(
matches!(error, StorageError::Unsupported { .. }),
"expected StorageError::Unsupported, got {error:?}"
);
assert!(
error.to_string().contains("second-digest"),
"the refusal must name a second-digest arm, got {error}"
);
}
fn initialize_utf16le_database(db_path: &std::path::Path) {
let conn = rusqlite::Connection::open(db_path).unwrap();
conn.execute_batch(
"PRAGMA encoding = 'UTF-16le'; \
CREATE TABLE __encoding_pin (x INTEGER); \
DROP TABLE __encoding_pin;",
)
.unwrap();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn blob_gc_evidence_accepts_valid_refs_in_a_utf16le_database() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
initialize_utf16le_database(&db_path);
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
let encoding: String = writer
.conn_mut()
.query_row("PRAGMA encoding", [], |row| row.get(0))
.unwrap();
assert_eq!(
encoding, "UTF-16le",
"the fixture database must actually be UTF-16le"
);
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('utf16-valid-attachment', 'entity', 'content', ?1, 0)",
rusqlite::params!["a".repeat(64)],
)
.unwrap();
writer
.conn_mut()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('utf16-valid-claim-key', ?1, 0)",
rusqlite::params!["b".repeat(64)],
)
.unwrap();
}
let sql = backend.sql();
super::validate_blob_gc_evidence(sql.as_ref())
.await
.expect("valid canonical refs must pass in a UTF-16LE database");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref_in_a_utf16le_database() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
initialize_utf16le_database(&db_path);
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
let encoding: String = writer
.conn_mut()
.query_row("PRAGMA encoding", [], |row| row.get(0))
.unwrap();
assert_eq!(
encoding, "UTF-16le",
"the fixture database must actually be UTF-16le"
);
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn_mut()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('nul-claim-key-utf16', ?1, 0)",
rusqlite::params![nul_embedded_canonical_ref()],
)
.unwrap();
}
let sql = backend.sql();
let error = super::validate_blob_gc_evidence(sql.as_ref())
.await
.expect_err("a NUL-embedded claim ref must refuse the sweep in UTF-16LE too");
assert!(
error.to_string().contains("blob_gc_claims"),
"expected the claims-table refusal, got {error:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn transactional_orphan_sweep_preserves_put_started_after_liveness_mark() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store.put(b"old orphan".to_vec()).await.unwrap();
let canonical_root = root.canonicalize().unwrap();
let (marked, release, _done) = sync_hook::install(&canonical_root);
let sweep = {
let store = store.clone();
let sql = backend.sql();
tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
};
assert!(
recv_blocking(marked).await,
"sweep must finish its liveness mark"
);
assert!(
store.write_lock.try_lock().is_err(),
"the sweep must hold the same root lock used by blob writers"
);
let (started_tx, started_rx) = std::sync::mpsc::channel();
let new_ref = {
let root = root.clone();
tokio::task::spawn_blocking(move || {
let _ = started_tx.send(());
put_blocking(&root, 0, b"new concurrent blob".to_vec())
})
};
assert!(recv_blocking(started_rx).await, "blob put must start");
release.send(()).unwrap();
let sweep_result = sweep.await.unwrap().unwrap();
let new_ref = new_ref.await.unwrap().unwrap();
assert_eq!(sweep_result.deleted, 1);
assert!(!store.exists(&orphan).await.unwrap());
assert!(
store.exists(&new_ref).await.unwrap(),
"a blob put started between the liveness mark and physical sweep must survive"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn transactional_orphan_sweep_releases_sqlite_writer_before_physical_delete() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store
.put(b"claim then delete outside sqlite".to_vec())
.await
.unwrap();
let canonical_root = root.canonicalize().unwrap();
let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
let sweep = {
let store = store.clone();
let sql = backend.sql();
tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
};
assert!(
recv_blocking(claimed).await,
"sweep must durably claim the orphan before physical deletion"
);
assert!(
store.exists(&orphan).await.unwrap(),
"the test seam must pause before the physical delete"
);
let external_database_lock = fs::OpenOptions::new()
.read(true)
.write(true)
.open(database_gc_lock_path(&db_path))
.unwrap();
assert!(
matches!(
fs4::FileExt::try_lock(&external_database_lock),
Err(fs4::TryLockError::WouldBlock)
),
"the sweep must retain cross-process database ownership while SQLite's writer is free"
);
let unrelated = rusqlite::Connection::open(&db_path).unwrap();
unrelated.busy_timeout(Duration::from_millis(100)).unwrap();
unrelated
.execute(
"INSERT INTO entities \
(id, namespace, kind, name, tags, created_at, updated_at) \
VALUES ('unrelated-writer', 'local', 'concept', 'unrelated', '[]', 1, 1)",
[],
)
.expect("external filesystem work must not retain SQLite's writer lock");
let claimed_err = unrelated
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('racing-reference', 'entity', 'content', ?1, 1)",
[orphan.as_str()],
)
.expect_err("a claimed content_ref must fail closed before deletion");
assert!(
claimed_err.to_string().contains("active blob sweep"),
"unexpected claim error: {claimed_err}"
);
release_delete.send(()).unwrap();
let result = sweep.await.unwrap().unwrap();
assert_eq!(result.deleted, 1);
assert!(!store.exists(&orphan).await.unwrap());
let remaining_claims: i64 = unrelated
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims WHERE content_ref = ?1",
[orphan.as_str()],
|row| row.get(0),
)
.unwrap();
assert_eq!(
remaining_claims, 0,
"successful deletion releases the claim"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cancelling_sweep_during_delete_keeps_owner_locks_until_blocking_work_finishes() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let orphan = store.put(b"cancelled sweep orphan".to_vec()).await.unwrap();
let canonical_root = root.canonicalize().unwrap();
let (claimed, release_delete, done) = sync_hook::install(&canonical_root);
let sweep = {
let store = store.clone();
let sql = backend.sql();
tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
};
assert!(recv_blocking(claimed).await);
sweep.abort();
assert!(sweep.await.unwrap_err().is_cancelled());
let external_root_lock = fs::OpenOptions::new()
.read(true)
.write(true)
.open(root.join(ROOT_WRITE_LOCK_FILE))
.unwrap();
let external_database_lock = fs::OpenOptions::new()
.read(true)
.write(true)
.open(database_gc_lock_path(&db_path))
.unwrap();
assert!(matches!(
fs4::FileExt::try_lock(&external_root_lock),
Err(fs4::TryLockError::WouldBlock)
));
assert!(matches!(
fs4::FileExt::try_lock(&external_database_lock),
Err(fs4::TryLockError::WouldBlock)
));
release_delete.send(()).unwrap();
let done_disconnected = tokio::task::spawn_blocking(move || done.recv().is_err())
.await
.unwrap();
assert!(
done_disconnected,
"the cancelled outer task cannot send done"
);
assert!(fs4::FileExt::try_lock(&external_root_lock).is_ok());
assert!(fs4::FileExt::try_lock(&external_database_lock).is_ok());
drop(external_root_lock);
drop(external_database_lock);
assert!(!store.exists(&orphan).await.unwrap());
let stranded_claims: i64 = rusqlite::Connection::open(&db_path)
.unwrap()
.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
.unwrap();
assert_eq!(
stranded_claims, 1,
"cancellation leaves a fail-closed claim"
);
let recovered = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(recovered.deleted, 0);
let remaining: i64 = rusqlite::Connection::open(&db_path)
.unwrap()
.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
.unwrap();
assert_eq!(remaining, 0, "the next exclusive owner recovers the claim");
}
#[tokio::test]
async fn transactional_orphan_sweep_recovers_stale_claims_fail_closed() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(60));
let bytes = b"republished after a crashed claim".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
let canonical_root = root.canonicalize().unwrap();
let root_key = blob_root_key(&canonical_root);
let absent_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
let former_probe_seed = "1111111111111111111111111111111111111111111111111111111111111111";
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES (?1, ?2, 1), (?1, ?3, 1), (?1, ?4, 1)",
rusqlite::params![
root_key,
content_ref.as_str(),
absent_ref,
former_probe_seed
],
)
.unwrap();
}
assert_eq!(store.put(bytes).await.unwrap(), content_ref);
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(result.deleted, 0);
assert_eq!(result.grace_period_skipped, 1);
assert!(store.exists(&content_ref).await.unwrap());
let remaining: i64 = backend
.pool()
.writer()
.unwrap()
.conn()
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?1",
[blob_root_key(&canonical_root)],
|row| row.get(0),
)
.unwrap();
assert_eq!(remaining, 0, "the next sweep recovers stale claims");
}
#[tokio::test]
async fn transactional_orphan_sweep_recovers_claims_after_root_relocation() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let old_root = dir.path().join("old-blobs");
let bytes = b"claim must follow a relocated blob root".to_vec();
let content_ref = {
let old_store = FsBlobStore::new(old_root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(60));
old_store.put(bytes).await.unwrap()
};
let old_root_key = blob_root_key(&old_root.canonicalize().unwrap());
backend
.pool()
.writer()
.unwrap()
.conn()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES (?1, ?2, 1)",
rusqlite::params![old_root_key, content_ref.as_str()],
)
.unwrap();
let new_root = dir.path().join("relocated-blobs");
std::fs::rename(&old_root, &new_root).unwrap();
let relocated_store = FsBlobStore::new(new_root, 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(60));
let result = relocated_store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(
result.deleted, 0,
"a fresh relocated blob remains protected"
);
assert_eq!(result.grace_period_skipped, 1);
assert!(relocated_store.exists(&content_ref).await.unwrap());
let remaining: i64 = backend
.pool()
.writer()
.unwrap()
.conn()
.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
.unwrap();
assert_eq!(
remaining, 0,
"exclusive database sweep ownership makes every pre-existing claim abandoned, \
even when its old path-derived root key no longer matches"
);
}
#[tokio::test]
async fn transactional_orphan_sweep_recovers_claims_copied_by_database_restore() {
let dir = tempfile::tempdir().unwrap();
let source_path = dir.path().join("source.db");
let restored_path = dir.path().join("restored.db");
let bytes = b"claim copied in an online database backup".to_vec();
let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
{
let source = crate::StorageBackend::sqlite_for_test(&source_path).unwrap();
let mut writer = source.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
writer
.conn()
.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('source-root-before-backup', ?1, 1)",
[content_ref.as_str()],
)
.unwrap();
writer
.conn()
.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)")
.unwrap();
}
std::fs::copy(&source_path, &restored_path).unwrap();
let restored = crate::StorageBackend::sqlite_for_test(&restored_path).unwrap();
let restored_root = dir.path().join("restored-blobs");
let store = FsBlobStore::new(restored_root, 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(60));
assert_eq!(store.put(bytes).await.unwrap(), content_ref);
let result = store
.transactional_orphan_sweep(restored.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(result.deleted, 0);
assert_eq!(result.grace_period_skipped, 1);
let remaining: i64 = restored
.pool()
.writer()
.unwrap()
.conn()
.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
.unwrap();
assert_eq!(remaining, 0, "restored claims are abandoned ownership");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn transactional_orphan_sweep_bounds_each_durable_claim_batch() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let candidate_count = BLOB_GC_CLAIM_BATCH_SIZE * 2 + 1;
for index in 0..candidate_count {
store
.put(format!("bounded claim candidate {index}").into_bytes())
.await
.unwrap();
}
let canonical_root = root.canonicalize().unwrap();
let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
let sweep = {
let store = store.clone();
let sql = backend.sql();
tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
};
assert!(
recv_blocking(claimed).await,
"the first bounded claim batch must commit before deletion"
);
let active_claims: i64 = rusqlite::Connection::open(&db_path)
.unwrap()
.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
.unwrap();
assert!(active_claims > 0);
assert!(
active_claims <= BLOB_GC_CLAIM_BATCH_SIZE as i64,
"one transaction may expose at most {BLOB_GC_CLAIM_BATCH_SIZE} claim rows; \
observed {active_claims}"
);
release_delete.send(()).unwrap();
let result = sweep.await.unwrap().unwrap();
assert_eq!(result.deleted, candidate_count as u64);
let remaining: i64 = rusqlite::Connection::open(&db_path)
.unwrap()
.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
.unwrap();
assert_eq!(remaining, 0);
}
#[tokio::test]
async fn abandoned_claim_recovery_deletes_at_most_one_batch_per_writer_hold() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
crate::run_migrations(writer.conn_mut()).unwrap();
let tx = writer.conn_mut().transaction().unwrap();
for index in 0..(BLOB_GC_CLAIM_BATCH_SIZE + 1) {
tx.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES ('abandoned-root', ?1, 1)",
[format!("{index:064x}")],
)
.unwrap();
}
tx.commit().unwrap();
}
let released = release_abandoned_blob_gc_claim_batch(backend.sql().as_ref())
.await
.unwrap();
assert_eq!(released, BLOB_GC_CLAIM_BATCH_SIZE as u64);
let remaining: i64 = backend
.pool()
.writer()
.unwrap()
.conn()
.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
.unwrap();
assert_eq!(remaining, 1);
}
#[tokio::test]
async fn transactional_orphan_sweep_refuses_corrupt_liveness_and_claim_evidence() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO);
let orphan = store
.put(b"must survive corrupt evidence".to_vec())
.await
.unwrap();
let conn = rusqlite::Connection::open(&db_path).unwrap();
conn.execute_batch("PRAGMA ignore_check_constraints = ON")
.unwrap();
conn.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('corrupt-live', 'entity', 'content', 'not-a-content-ref', 1)",
[],
)
.unwrap();
conn.execute_batch("PRAGMA ignore_check_constraints = OFF")
.unwrap();
let live_error = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.expect_err("corrupt live evidence must fail closed");
assert!(matches!(live_error, StorageError::InvalidInput { .. }));
assert!(
store.exists(&orphan).await.unwrap(),
"no file may be removed after corrupt live evidence"
);
conn.execute(
"DELETE FROM attachments WHERE record_uuid = 'corrupt-live'",
[],
)
.unwrap();
let root_key = blob_root_key(&root.canonicalize().unwrap());
conn.execute(
"INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
VALUES (?1, 'also-not-a-content-ref', 1)",
[root_key.as_str()],
)
.unwrap();
let claim_error = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.expect_err("corrupt durable claim evidence must fail closed");
assert!(matches!(claim_error, StorageError::InvalidInput { .. }));
assert!(
store.exists(&orphan).await.unwrap(),
"no file may be removed after corrupt claim evidence"
);
let remaining: i64 = conn
.query_row(
"SELECT COUNT(*) FROM blob_gc_claims \
WHERE root_key = ?1 AND content_ref = 'also-not-a-content-ref'",
[root_key.as_str()],
|row| row.get(0),
)
.unwrap();
assert_eq!(
remaining, 1,
"corrupt claim evidence is not silently erased"
);
let probe_residue: i64 = conn
.query_row(
"SELECT (SELECT COUNT(*) FROM blob_gc_claims \
WHERE root_key GLOB '__fence_probe-*') \
+ (SELECT COUNT(*) FROM attachments \
WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
probe_residue, 0,
"invalid evidence must abort before the functional fence probe"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn transactional_orphan_sweep_republishes_deduplicated_external_put() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO),
);
let payload = b"existing orphan republished during sweep".to_vec();
let orphan = store.put(payload.clone()).await.unwrap();
let canonical_root = root.canonicalize().unwrap();
let (marked, release, _done) = sync_hook::install(&canonical_root);
let sweep = {
let store = store.clone();
let sql = backend.sql();
tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
};
assert!(
recv_blocking(marked).await,
"sweep must finish its liveness mark"
);
let external_lock = fs::OpenOptions::new()
.read(true)
.write(true)
.open(root.join(ROOT_WRITE_LOCK_FILE))
.unwrap();
assert!(
matches!(
fs4::FileExt::try_lock(&external_lock),
Err(fs4::TryLockError::WouldBlock)
),
"the sweep must exclude a publisher using an independently opened root lock"
);
let (started_tx, started_rx) = std::sync::mpsc::channel();
let republished = {
let root = root.clone();
tokio::task::spawn_blocking(move || {
let _ = started_tx.send(());
put_blocking(&root, 0, payload)
})
};
assert!(recv_blocking(started_rx).await, "blob put must start");
release.send(()).unwrap();
let sweep_result = sweep.await.unwrap().unwrap();
let republished = republished.await.unwrap().unwrap();
assert_eq!(sweep_result.deleted, 1);
assert_eq!(republished, orphan);
assert!(
store.exists(&republished).await.unwrap(),
"a deduplicated put concurrent with the sweep must not return a deleted reference"
);
}
#[tokio::test]
async fn transactional_orphan_sweep_uses_all_attachment_refs_as_live() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO);
let live = store.put(b"live".to_vec()).await.unwrap();
let soft_deleted = store.put(b"soft deleted".to_vec()).await.unwrap();
let orphan = store.put(b"orphan".to_vec()).await.unwrap();
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute_batch(
"INSERT INTO entities \
(id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
VALUES ('live', 'local', 'document', 'live', '[]', 1, 1, NULL), \
('deleted', 'local', 'document', 'deleted', '[]', 1, 1, 2);",
)
.unwrap();
writer
.conn()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('live', 'entity', 'content', ?1, 1), \
('deleted', 'entity', 'content', ?2, 1)",
rusqlite::params![live.as_str(), soft_deleted.as_str()],
)
.unwrap();
}
let dry_run = store
.transactional_orphan_sweep(backend.sql().as_ref(), true)
.await
.unwrap();
assert_eq!(dry_run.would_delete, 1);
assert_eq!(dry_run.deleted, 0);
assert!(store.exists(&soft_deleted).await.unwrap());
assert!(store.exists(&orphan).await.unwrap());
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(result.scanned, 3);
assert_eq!(result.deleted, 1);
assert!(store.exists(&live).await.unwrap());
assert!(
store.exists(&soft_deleted).await.unwrap(),
"soft delete retains attachment rows and their blobs"
);
assert!(!store.exists(&orphan).await.unwrap());
}
#[tokio::test]
async fn transactional_orphan_sweep_protects_a_freshly_published_blob_before_its_reference_commits(
) {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
let blob = store
.put(b"published, reference not yet committed".to_vec())
.await
.unwrap();
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(result.deleted, 0, "the blob must survive: {result:?}");
assert_eq!(
result.would_delete, 0,
"not treated as a deletable orphan: {result:?}"
);
assert_eq!(
result.grace_period_skipped, 1,
"must be reported as grace-protected rather than silently ignored: {result:?}"
);
assert!(
store.exists(&blob).await.unwrap(),
"a blob still inside its publish grace period must survive the sweep"
);
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO entities \
(id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
[],
)
.unwrap();
writer
.conn()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('e1', 'entity', 'content', ?1, 1)",
[blob.as_str()],
)
.unwrap();
}
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(result.deleted, 0);
assert!(store.exists(&blob).await.unwrap());
}
#[tokio::test]
async fn put_republishing_an_aged_orphan_restarts_its_grace_clock_before_the_reference_commits()
{
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(60));
let bytes = b"old orphan re-published".to_vec();
let first = store.put(bytes.clone()).await.unwrap();
let path = shard_path(store.root(), &first);
let old_mtime = SystemTime::now() - Duration::from_secs(3600);
fs::OpenOptions::new()
.write(true)
.open(&path)
.unwrap()
.set_modified(old_mtime)
.unwrap();
let second = store.put(bytes).await.unwrap();
assert_eq!(first, second);
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(
result.deleted, 0,
"a dedup-republished blob must survive a sweep landing before its reference \
commits: {result:?}"
);
assert_eq!(
result.grace_period_skipped, 1,
"must be reported as grace-protected, not silently ignored: {result:?}"
);
assert!(store.exists(&first).await.unwrap());
{
let writer = backend.pool().writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO entities \
(id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
[],
)
.unwrap();
writer
.conn()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES ('e1', 'entity', 'content', ?1, 1)",
[first.as_str()],
)
.unwrap();
}
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(result.deleted, 0);
assert!(
store.exists(&first).await.unwrap(),
"the blob must stay live once its reference has committed"
);
}
#[tokio::test]
async fn put_dedup_mtime_refresh_has_no_observable_effect_under_zero_grace_period() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::ZERO);
let bytes = b"zero grace dedup refresh".to_vec();
let first = store.put(bytes.clone()).await.unwrap();
let second = store.put(bytes.clone()).await.unwrap();
assert_eq!(first, second);
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(
result.deleted, 1,
"a zero grace period must still delete an unreferenced blob even after a dedup \
put refreshed its mtime: {result:?}"
);
assert!(!store.exists(&first).await.unwrap());
}
#[tokio::test]
async fn transactional_orphan_sweep_still_removes_orphans_older_than_the_grace_period() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let store = FsBlobStore::new(dir.path().join("blobs"), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(60));
let orphan = store
.put(b"actually orphaned, published long ago".to_vec())
.await
.unwrap();
let path = shard_path(store.root(), &orphan);
let old_mtime = SystemTime::now() - Duration::from_secs(3600);
fs::OpenOptions::new()
.write(true)
.open(&path)
.unwrap()
.set_modified(old_mtime)
.unwrap();
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(
result.deleted, 1,
"an orphan older than the grace period must still be swept: {result:?}"
);
assert_eq!(result.grace_period_skipped, 0);
assert!(!store.exists(&orphan).await.unwrap());
}
#[test]
fn resolve_blob_root_prefers_env_var() {
let _guard = ENV_LOCK.lock().unwrap();
std::env::set_var("KHIVE_BLOB_ROOT", "/tmp/env-override-root");
let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
std::env::remove_var("KHIVE_BLOB_ROOT");
assert_eq!(resolved.unwrap(), PathBuf::from("/tmp/env-override-root"));
}
#[test]
fn resolve_blob_root_prefers_config_over_default() {
let _guard = ENV_LOCK.lock().unwrap();
std::env::remove_var("KHIVE_BLOB_ROOT");
let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
assert_eq!(resolved.unwrap(), PathBuf::from("/cfg/root"));
}
#[test]
fn resolve_blob_root_defaults_beside_db_dir() {
let _guard = ENV_LOCK.lock().unwrap();
std::env::remove_var("KHIVE_BLOB_ROOT");
let resolved = resolve_blob_root(Some(Path::new("/db/dir")), None);
assert_eq!(resolved.unwrap(), PathBuf::from("/db/dir/blobs"));
}
#[test]
fn resolve_blob_root_errors_with_no_env_config_or_db_dir() {
let _guard = ENV_LOCK.lock().unwrap();
std::env::remove_var("KHIVE_BLOB_ROOT");
let resolved = resolve_blob_root(None, None);
assert!(resolved.is_err());
}
static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn transactional_orphan_sweep_walk_ignores_a_leaf_swapped_for_an_outside_symlink_mid_scan(
) {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend =
std::sync::Arc::new(crate::StorageBackend::sqlite_for_test(&db_path).unwrap());
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = std::sync::Arc::new(
FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(3600)),
);
let real = store
.put(b"real freshly-published blob".to_vec())
.await
.unwrap();
let real_path = shard_path(&root, &real);
let outside_dir = dir.path().join("outside");
fs::create_dir_all(&outside_dir).unwrap();
let decoy_path = outside_dir.join(real.as_str());
fs::write(&decoy_path, b"outside decoy, must never be observed").unwrap();
let ancient = SystemTime::now() - Duration::from_secs(7200);
fs::OpenOptions::new()
.write(true)
.open(&decoy_path)
.unwrap()
.set_modified(ancient)
.unwrap();
let (reached, release) = walk_leaf_sync_hook::install(&root);
let sweep = {
let store = store.clone();
let sql = backend.sql();
tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), true).await })
};
assert!(
recv_blocking(reached).await,
"sweep walk must reach the leaf classification pause"
);
fs::remove_file(&real_path).unwrap();
std::os::unix::fs::symlink(&decoy_path, &real_path).unwrap();
release.send(()).unwrap();
let result = sweep.await.unwrap().unwrap();
assert_eq!(
result.would_delete, 0,
"an outside decoy's stale mtime must never make an in-root candidate \
eligible for deletion: {result:?}"
);
assert_eq!(
result.grace_period_skipped, 0,
"the swapped leaf is a symlink; `openat(..., O_NOFOLLOW)` refuses it, so it \
must be dropped from candidates entirely rather than counted (real or \
outside) at all: {result:?}"
);
assert_eq!(
result.scanned, 0,
"the symlinked leaf must never be scanned as a candidate: {result:?}"
);
fs::remove_file(&real_path).unwrap();
fs::write(&real_path, b"real freshly-published blob").unwrap();
let control = store
.transactional_orphan_sweep(backend.sql().as_ref(), true)
.await
.unwrap();
assert_eq!(
control.grace_period_skipped, 1,
"control: the un-replaced root must still classify the real candidate as \
grace-protected: {control:?}"
);
assert_eq!(control.would_delete, 0);
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn transactional_orphan_sweep_walk_still_finds_a_real_orphan_past_its_grace_period() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("khive.db");
let backend = crate::StorageBackend::sqlite_for_test(&db_path).unwrap();
{
let mut writer = backend.pool().writer().unwrap();
prepare_completed_v21_gc_fixture(writer.conn_mut());
}
let root = dir.path().join("blobs");
let store = FsBlobStore::new(root.clone(), 0)
.unwrap()
.with_orphan_sweep_grace(Duration::from_secs(60));
let orphan = store.put(b"aged real orphan".to_vec()).await.unwrap();
let path = shard_path(&root, &orphan);
let ancient = SystemTime::now() - Duration::from_secs(3600);
fs::OpenOptions::new()
.write(true)
.open(&path)
.unwrap()
.set_modified(ancient)
.unwrap();
let result = store
.transactional_orphan_sweep(backend.sql().as_ref(), false)
.await
.unwrap();
assert_eq!(
result.deleted, 1,
"a real orphan older than the grace period must still be swept: {result:?}"
);
assert!(!store.exists(&orphan).await.unwrap());
}
}