Skip to main content

khive_db/stores/
blob.rs

1//! Filesystem-backed `BlobStore` — content-addressed, BLAKE3-sharded on disk.
2//!
3//! Layout: `<root>/<hex[0..2]>/<hex[2..4]>/<hex>`, plus a root-local advisory
4//! lock file. The two-level shard is identical in shape to git's loose-object
5//! store, so a root holding millions of blobs never puts more than a few
6//! thousand entries in one directory. Writes are atomic-publish (khive#292):
7//! bytes land in a temporary entry in the SAME shard directory as the final
8//! path (guaranteeing same-filesystem rename), the written length is checked
9//! against the input length, then an atomic rename publishes the entry —
10//! crash-safe (a crash mid-write leaves an orphaned temp file, never a
11//! partially-committed blob).
12
13use std::collections::HashMap;
14use std::fs;
15use std::io::{Read, Write};
16use std::path::{Path, PathBuf};
17use std::sync::{Arc, Mutex as StdMutex, OnceLock};
18use std::time::{Duration, SystemTime};
19
20use async_trait::async_trait;
21
22use khive_storage::blob::{
23    BlobOrphanSweepConfig, BlobOrphanSweepResult, BlobStore, ContentRef, MAX_BLOB_WHOLE_BYTES,
24};
25use khive_storage::error::StorageError;
26use khive_storage::types::{SqlRow, SqlStatement, SqlValue, StorageResult};
27use khive_storage::{AtomicUnitOp, SqlAccess, StorageCapability};
28
29use crate::error::SqliteError;
30use uuid::Uuid;
31
32const ROOT_WRITE_LOCK_FILE: &str = ".khive-blob-write.lock";
33const DATABASE_GC_LOCK_SUFFIX: &str = ".khive-blob-gc.lock";
34/// Maximum candidates represented by one claim transaction and its matching
35/// physical-delete/cleanup cycle. This bounds JSON binding, returned rows,
36/// claim-table pages dirtied per transaction, and exclusive-writer hold time.
37const BLOB_GC_CLAIM_BATCH_SIZE: usize = 128;
38
39fn map_io_err(e: std::io::Error, op: &'static str) -> StorageError {
40    StorageError::driver(StorageCapability::Blob, op, e)
41}
42
43#[cfg(any(test, not(unix)))]
44fn shard_path(root: &Path, content_ref: &ContentRef) -> PathBuf {
45    let hex = content_ref.as_str();
46    root.join(&hex[0..2]).join(&hex[2..4]).join(hex)
47}
48
49#[cfg(unix)]
50fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
51    open_dir_no_follow(root)
52}
53
54#[cfg(windows)]
55fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
56    use std::fs::OpenOptions;
57    use std::os::windows::fs::OpenOptionsExt;
58
59    const FILE_SHARE_READ: u32 = 0x1;
60    const FILE_SHARE_WRITE: u32 = 0x2;
61    const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
62    const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
63
64    let handle = OpenOptions::new()
65        .read(true)
66        // Retain the handle for the store's lifetime. Omitting
67        // FILE_SHARE_DELETE prevents the root's final component from being
68        // renamed or removed while this store can still issue operations.
69        .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
70        .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
71        .open(root)?;
72    let file_type = handle.metadata()?.file_type();
73    if file_type.is_symlink() || !file_type.is_dir() {
74        return Err(std::io::Error::new(
75            std::io::ErrorKind::InvalidInput,
76            format!(
77                "blob store root is not a directory or is a reparse point: {}",
78                root.display()
79            ),
80        ));
81    }
82    Ok(handle)
83}
84
85#[cfg(not(any(unix, windows)))]
86fn open_blob_root_handle(root: &Path) -> std::io::Result<std::fs::File> {
87    let handle = std::fs::File::open(root)?;
88    if !handle.metadata()?.is_dir() {
89        return Err(std::io::Error::new(
90            std::io::ErrorKind::InvalidInput,
91            format!("blob store root is not a directory: {}", root.display()),
92        ));
93    }
94    Ok(handle)
95}
96
97/// Fail closed when the canonical root spelling no longer names the same
98/// directory handle retained at construction. The comparison is only an
99/// invariant check: filesystem authority for the operation itself remains
100/// the retained handle, so a replacement after this check cannot redirect a
101/// handle-relative walk.
102#[cfg(unix)]
103fn verify_blob_root_identity(root: &Path, root_handle: &std::fs::File) -> std::io::Result<()> {
104    use std::os::unix::fs::MetadataExt;
105
106    let current = open_dir_no_follow(root).map_err(|error| {
107        std::io::Error::new(
108            std::io::ErrorKind::InvalidInput,
109            format!(
110                "blob store root is no longer reachable as its initialization-time directory ({}): {error}",
111                root.display()
112            ),
113        )
114    })?;
115    let expected = root_handle.metadata()?;
116    let current = current.metadata()?;
117    if expected.dev() != current.dev() || expected.ino() != current.ino() {
118        return Err(std::io::Error::new(
119            std::io::ErrorKind::InvalidInput,
120            format!(
121                "blob store root no longer names its initialization-time directory: {}",
122                root.display()
123            ),
124        ));
125    }
126    Ok(())
127}
128
129#[cfg(not(unix))]
130fn verify_blob_root_identity(root: &Path, _root_handle: &std::fs::File) -> std::io::Result<()> {
131    // `root` is canonicalized before the retained handle opens. Re-resolving
132    // the spelling catches an ancestor redirect on Windows; the retained
133    // no-share-delete handle separately pins the root's final component, and
134    // the Windows leaf helpers compare their resolved target handles to that
135    // initialization-time root handle before use.
136    let current = root.canonicalize().map_err(|error| {
137        std::io::Error::new(
138            std::io::ErrorKind::InvalidInput,
139            format!(
140                "blob store root is no longer reachable as its initialization-time directory ({}): {error}",
141                root.display()
142            ),
143        )
144    })?;
145    if current != root {
146        return Err(std::io::Error::new(
147            std::io::ErrorKind::InvalidInput,
148            format!(
149                "blob store root no longer names its initialization-time directory: {}",
150                root.display()
151            ),
152        ));
153    }
154    Ok(())
155}
156
157/// Unlink one blob's shard-relative file using `O_NOFOLLOW`-verified
158/// descriptor traversal instead of a plain path-based delete.
159///
160/// A path-based `fs::remove_file(shard_path(root, content_ref))` resolves
161/// every path component through the kernel exactly like any other path
162/// lookup. If either shard-directory level (`root/<hex[0..2]>` or
163/// `root/<hex[0..2]>/<hex[2..4]>`) has been replaced with a symlink —
164/// through a misconfigured root, a shared/writable parent directory, or a
165/// race with another process — that lookup follows it and can unlink a file
166/// entirely outside the blob root. Each shard component is instead opened
167/// relative to the previous, already-verified descriptor with
168/// `O_DIRECTORY | O_NOFOLLOW`, so a symlink planted at either level is
169/// refused (`ELOOP`) rather than followed, and the final `unlinkat` runs
170/// relative to the verified leaf descriptor rather than a re-resolved path
171/// string. Same fd-pinned idiom as `khive-db`'s walpin sidecar writes
172/// (`walpin.rs`) and `khive-vamana`'s external-id sidecar
173/// (`external_ids.rs`) use for the same TOCTOU hazard.
174#[cfg(unix)]
175fn unlink_blob_shard_file_no_follow(
176    root: &Path,
177    root_handle: &std::fs::File,
178    content_ref: &ContentRef,
179) -> std::io::Result<()> {
180    use std::os::unix::io::AsRawFd;
181
182    verify_blob_root_identity(root, root_handle)?;
183    let hex = content_ref.as_str();
184    let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
185    let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
186    // `unlinkat` removes this directory entry itself rather than following a
187    // leaf symlink, matching the original delete semantics.
188    unlink_entry_at(shard2_dir.as_raw_fd(), hex)
189}
190
191#[cfg(windows)]
192fn unlink_blob_shard_file_no_follow(
193    root: &Path,
194    root_handle: &std::fs::File,
195    content_ref: &ContentRef,
196) -> std::io::Result<()> {
197    // Windows equivalent of the Unix arm's fd-pinned walk. The Unix arm's
198    // guarantee is: the root's ancestors are resolved exactly once (at
199    // `open(root)`), every later step is relative to an already-verified
200    // open descriptor, and the final unlink acts on a descriptor, never on
201    // a re-resolved path string. This arm reproduces each property from
202    // handle semantics:
203    //
204    // 1. No-follow verification BY HANDLE: each directory level is opened
205    //    with `FILE_FLAG_OPEN_REPARSE_POINT` (plus `FILE_FLAG_BACKUP_SEMANTICS`,
206    //    which is what permits opening a directory handle at all), so a
207    //    junction or symlink planted at that level yields a handle to the
208    //    reparse point itself rather than to its target, and the
209    //    handle-derived metadata (`File::metadata`, which queries the handle,
210    //    not a re-resolved path) exposes it for refusal.
211    // 2. Pinning: the directory handles are opened WITHOUT
212    //    `FILE_SHARE_DELETE`. Deleting or renaming a directory requires an
213    //    open with `DELETE` access, which fails with a sharing violation
214    //    while these handles are held, so no checked component can be
215    //    swapped for the duration of the call.
216    // 3. Deletion BY HANDLE with a handle-anchored identity check. A
217    //    path-based `remove_file` here would re-resolve the full path from
218    //    the volume root, so a reparse point swapped at an UNPINNED ancestor
219    //    of `root` (which this function cannot pin — it does not own them)
220    //    could redirect the delete outside the blob root even while all
221    //    three pins hold. Instead the target file itself is opened with
222    //    `DELETE` access and `FILE_FLAG_OPEN_REPARSE_POINT` (a symlink leaf
223    //    opens as the link entry, matching `unlinkat` semantics), its TRUE
224    //    resolved path is read back from the handle with
225    //    `GetFinalPathNameByHandleW`, and the delete proceeds only if that
226    //    path equals the root pin's own handle-final path extended by the
227    //    verified shard components. The root pin's final path is a property
228    //    of the already-open handle — an ancestor swapped after the pin
229    //    opened cannot change it, while it does change (and thereby betrays)
230    //    the file handle's resolution. The delete itself is
231    //    `SetFileInformationByHandle(FileDispositionInfo)` on the verified
232    //    handle: no path is ever re-resolved between check and use.
233    //
234    // An ancestor reparse point already in place BEFORE the root pin opens
235    // resolves identically for the pin and the file and is accepted — the
236    // same exposure the Unix arm accepts for a symlinked ancestor at
237    // `open(root)` time; that is the operator's configured deployment, not
238    // a check-to-use window.
239    use std::fs::OpenOptions;
240    use std::os::windows::ffi::OsStringExt;
241    use std::os::windows::fs::OpenOptionsExt;
242    use std::os::windows::io::AsRawHandle;
243    use windows_sys::Win32::Storage::FileSystem::{
244        FileDispositionInfo, GetFinalPathNameByHandleW, SetFileInformationByHandle,
245        FILE_DISPOSITION_INFO,
246    };
247
248    const FILE_SHARE_READ: u32 = 0x1;
249    const FILE_SHARE_WRITE: u32 = 0x2;
250    const DELETE: u32 = 0x0001_0000;
251    const FILE_READ_ATTRIBUTES: u32 = 0x80;
252    const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
253    const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
254    /// `FILE_NAME_NORMALIZED | VOLUME_NAME_DOS` — both zero; named for the
255    /// contract (normalized on-disk case, drive-letter form) rather than
256    /// passing a bare 0.
257    const FINAL_PATH_FLAGS: u32 = 0x0;
258
259    fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
260        let dir = OpenOptions::new()
261            .read(true)
262            .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
263            .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
264            .open(path)?;
265        let file_type = dir.metadata()?.file_type();
266        if file_type.is_symlink() || !file_type.is_dir() {
267            return Err(std::io::Error::new(
268                std::io::ErrorKind::InvalidInput,
269                format!(
270                    "refusing to unlink blob shard file through non-directory or \
271                     reparse-point path component: {}",
272                    path.display()
273                ),
274            ));
275        }
276        Ok(dir)
277    }
278
279    /// The handle's true, fully resolved path (`GetFinalPathNameByHandleW`,
280    /// normalized `\\?\`-prefixed DOS form). A handle property: later
281    /// changes to any directory the original path traversed cannot alter it.
282    fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<std::path::PathBuf> {
283        let handle = file.as_raw_handle();
284        let mut buf: Vec<u16> = vec![0; 512];
285        loop {
286            let len = unsafe {
287                GetFinalPathNameByHandleW(
288                    handle as _,
289                    buf.as_mut_ptr(),
290                    buf.len() as u32,
291                    FINAL_PATH_FLAGS,
292                )
293            };
294            if len == 0 {
295                return Err(std::io::Error::last_os_error());
296            }
297            let len = len as usize;
298            if len <= buf.len() {
299                buf.truncate(len);
300                return Ok(std::path::PathBuf::from(std::ffi::OsString::from_wide(
301                    &buf,
302                )));
303            }
304            // Returned length is the required buffer size (in wide chars,
305            // including the terminator) when the buffer was too small.
306            buf.resize(len, 0);
307        }
308    }
309
310    verify_blob_root_identity(root, root_handle)?;
311    let hex = content_ref.as_str();
312    let shard1 = root.join(&hex[0..2]);
313    let shard2 = shard1.join(&hex[2..4]);
314    let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
315    let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
316
317    let expected = final_path_by_handle(root_handle)?
318        .join(&hex[0..2])
319        .join(&hex[2..4])
320        .join(hex);
321
322    // Sharing READ|WRITE but NOT DELETE: renaming or deleting a file requires
323    // an open with `DELETE` access, which fails with a sharing violation while
324    // this handle is held, so the file whose final path is validated below is
325    // the same file the disposition call deletes — the leaf is pinned exactly
326    // like the directory components above.
327    let target = OpenOptions::new()
328        .access_mode(DELETE | FILE_READ_ATTRIBUTES)
329        .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
330        .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
331        .open(shard2.join(hex))?;
332
333    let resolved = final_path_by_handle(&target)?;
334    if resolved != expected {
335        return Err(std::io::Error::new(
336            std::io::ErrorKind::InvalidInput,
337            format!(
338                "refusing blob delete: handle resolved outside the verified blob root \
339                 (expected {}, resolved {})",
340                expected.display(),
341                resolved.display()
342            ),
343        ));
344    }
345
346    let disposition = FILE_DISPOSITION_INFO { DeleteFile: true };
347    let ok = unsafe {
348        SetFileInformationByHandle(
349            target.as_raw_handle() as _,
350            FileDispositionInfo,
351            std::ptr::from_ref(&disposition).cast(),
352            std::mem::size_of::<FILE_DISPOSITION_INFO>() as u32,
353        )
354    };
355    if ok == 0 {
356        return Err(std::io::Error::last_os_error());
357    }
358    Ok(())
359}
360
361#[cfg(not(any(unix, windows)))]
362fn unlink_blob_shard_file_no_follow(
363    root: &Path,
364    root_handle: &std::fs::File,
365    content_ref: &ContentRef,
366) -> std::io::Result<()> {
367    // Neither `openat`/`O_NOFOLLOW` nor Windows directory-handle pinning is
368    // available on this tier (which no release artifact targets), so this
369    // checks each shard path component with `symlink_metadata` before the
370    // delete. `std::fs::symlink_metadata` reports symlinks through
371    // `file_type().is_symlink()` without following them, so a link planted
372    // at the root or either shard level is refused rather than walked into
373    // by the final `remove_file`.
374    //
375    // Residual limitation, accepted for this descriptorless platform tier:
376    // nothing here holds an open, referentially-verified handle on the
377    // checked directories between this check and the `remove_file` call
378    // below, so a component could still be swapped for a link in that
379    // window (TOCTOU).
380    verify_blob_root_identity(root, root_handle)?;
381    let hex = content_ref.as_str();
382    let shard1 = root.join(&hex[0..2]);
383    let shard2 = shard1.join(&hex[2..4]);
384    for component in [root, shard1.as_path(), shard2.as_path()] {
385        let metadata = fs::symlink_metadata(component)?;
386        if metadata.file_type().is_symlink() {
387            return Err(std::io::Error::new(
388                std::io::ErrorKind::InvalidInput,
389                format!(
390                    "refusing to unlink blob shard file through symlinked path component: {}",
391                    component.display()
392                ),
393            ));
394        }
395    }
396    fs::remove_file(shard2.join(hex))
397}
398
399/// Open one blob leaf without following any root/shard/leaf symlink and
400/// return the single handle that must remain the authority for metadata,
401/// bounded bytes, and digest verification.
402#[cfg(unix)]
403fn open_blob_shard_file_no_follow(
404    root: &Path,
405    root_handle: &std::fs::File,
406    content_ref: &ContentRef,
407) -> std::io::Result<std::fs::File> {
408    verify_blob_root_identity(root, root_handle)?;
409    open_blob_shard_file_at_no_follow(root_handle, content_ref, libc::O_RDONLY)
410}
411
412#[cfg(windows)]
413fn open_blob_shard_file_no_follow(
414    root: &Path,
415    root_handle: &std::fs::File,
416    content_ref: &ContentRef,
417) -> std::io::Result<std::fs::File> {
418    use std::fs::OpenOptions;
419    use std::os::windows::ffi::OsStringExt;
420    use std::os::windows::fs::OpenOptionsExt;
421    use std::os::windows::io::AsRawHandle;
422    use windows_sys::Win32::Storage::FileSystem::GetFinalPathNameByHandleW;
423
424    const FILE_SHARE_READ: u32 = 0x1;
425    const FILE_SHARE_WRITE: u32 = 0x2;
426    const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
427    const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
428    const FINAL_PATH_FLAGS: u32 = 0x0;
429
430    fn open_dir_pinned_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
431        let dir = OpenOptions::new()
432            .read(true)
433            // Omitting FILE_SHARE_DELETE pins this checked component against
434            // rename/removal until the leaf is open and verified.
435            .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
436            .custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
437            .open(path)?;
438        let file_type = dir.metadata()?.file_type();
439        if file_type.is_symlink() || !file_type.is_dir() {
440            return Err(std::io::Error::new(
441                std::io::ErrorKind::InvalidInput,
442                format!(
443                    "refusing to read a blob through non-directory or reparse-point component: {}",
444                    path.display()
445                ),
446            ));
447        }
448        Ok(dir)
449    }
450
451    fn final_path_by_handle(file: &std::fs::File) -> std::io::Result<PathBuf> {
452        let mut buf: Vec<u16> = vec![0; 512];
453        loop {
454            let len = unsafe {
455                GetFinalPathNameByHandleW(
456                    file.as_raw_handle() as _,
457                    buf.as_mut_ptr(),
458                    buf.len() as u32,
459                    FINAL_PATH_FLAGS,
460                )
461            };
462            if len == 0 {
463                return Err(std::io::Error::last_os_error());
464            }
465            let len = len as usize;
466            if len <= buf.len() {
467                buf.truncate(len);
468                return Ok(PathBuf::from(std::ffi::OsString::from_wide(&buf)));
469            }
470            buf.resize(len, 0);
471        }
472    }
473
474    verify_blob_root_identity(root, root_handle)?;
475    let hex = content_ref.as_str();
476    let shard1 = root.join(&hex[0..2]);
477    let shard2 = shard1.join(&hex[2..4]);
478    let _shard1_pin = open_dir_pinned_no_follow(&shard1)?;
479    let _shard2_pin = open_dir_pinned_no_follow(&shard2)?;
480    let expected = final_path_by_handle(root_handle)?
481        .join(&hex[0..2])
482        .join(&hex[2..4])
483        .join(hex);
484
485    let target = OpenOptions::new()
486        .read(true)
487        .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE)
488        .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
489        .open(shard2.join(hex))?;
490    let file_type = target.metadata()?.file_type();
491    if file_type.is_symlink() || !file_type.is_file() {
492        return Err(std::io::Error::new(
493            std::io::ErrorKind::InvalidInput,
494            "refusing to read a blob leaf that is not a regular file",
495        ));
496    }
497    let resolved = final_path_by_handle(&target)?;
498    if resolved != expected {
499        return Err(std::io::Error::new(
500            std::io::ErrorKind::InvalidInput,
501            format!(
502                "refusing blob read: handle resolved outside the verified blob root (expected {}, resolved {})",
503                expected.display(),
504                resolved.display()
505            ),
506        ));
507    }
508    Ok(target)
509}
510
511#[cfg(not(any(unix, windows)))]
512fn open_blob_shard_file_no_follow(
513    root: &Path,
514    root_handle: &std::fs::File,
515    _content_ref: &ContentRef,
516) -> std::io::Result<std::fs::File> {
517    verify_blob_root_identity(root, root_handle)?;
518    Err(std::io::Error::new(
519        std::io::ErrorKind::Unsupported,
520        "bounded verified blob reads require handle-relative no-follow file APIs",
521    ))
522}
523
524#[cfg(unix)]
525fn open_dir_no_follow(path: &Path) -> std::io::Result<std::fs::File> {
526    use std::os::unix::ffi::OsStrExt;
527    use std::os::unix::io::FromRawFd;
528
529    let c_path = std::ffi::CString::new(path.as_os_str().as_bytes())
530        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
531    // SAFETY: `c_path` is NUL-terminated for the call; a successful fd is
532    // uniquely owned and wrapped immediately below.
533    let fd = unsafe {
534        libc::open(
535            c_path.as_ptr(),
536            libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
537        )
538    };
539    if fd < 0 {
540        return Err(std::io::Error::last_os_error());
541    }
542    // SAFETY: `fd` was just returned by the successful `open` above and is
543    // uniquely owned by this `File`, which closes it exactly once on drop.
544    Ok(unsafe { std::fs::File::from_raw_fd(fd) })
545}
546
547#[cfg(unix)]
548fn openat_dir_no_follow(
549    parent_fd: std::os::unix::io::RawFd,
550    name: &str,
551) -> std::io::Result<std::fs::File> {
552    use std::os::unix::io::FromRawFd;
553
554    let c_name = std::ffi::CString::new(name)
555        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
556    // SAFETY: `c_name` is NUL-terminated; `parent_fd` is a live, open
557    // directory descriptor for the duration of this call. A successful fd
558    // is uniquely owned and wrapped immediately below.
559    let fd = unsafe {
560        libc::openat(
561            parent_fd,
562            c_name.as_ptr(),
563            libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
564        )
565    };
566    if fd < 0 {
567        return Err(std::io::Error::last_os_error());
568    }
569    // SAFETY: `fd` was just returned by the successful `openat` above and is
570    // uniquely owned by this `File`, which closes it exactly once on drop.
571    Ok(unsafe { std::fs::File::from_raw_fd(fd) })
572}
573
574#[cfg(unix)]
575fn openat_regular_file_no_follow(
576    parent_fd: std::os::unix::io::RawFd,
577    name: &str,
578    access_flags: libc::c_int,
579) -> std::io::Result<std::fs::File> {
580    use std::os::unix::io::FromRawFd;
581
582    let c_name = std::ffi::CString::new(name)
583        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
584    // O_NONBLOCK prevents a planted FIFO/device-like entry from hanging the
585    // worker before handle metadata can reject it. It is inert for regular
586    // files.
587    let fd = unsafe {
588        libc::openat(
589            parent_fd,
590            c_name.as_ptr(),
591            access_flags | libc::O_NOFOLLOW | libc::O_CLOEXEC | libc::O_NONBLOCK,
592        )
593    };
594    if fd < 0 {
595        return Err(std::io::Error::last_os_error());
596    }
597    // SAFETY: `fd` was newly returned and is transferred exactly once.
598    let file = unsafe { std::fs::File::from_raw_fd(fd) };
599    if !file.metadata()?.file_type().is_file() {
600        return Err(std::io::Error::new(
601            std::io::ErrorKind::InvalidInput,
602            format!("refusing non-regular blob store entry: {name}"),
603        ));
604    }
605    Ok(file)
606}
607
608#[cfg(unix)]
609fn open_blob_shard_file_at_no_follow(
610    root_handle: &std::fs::File,
611    content_ref: &ContentRef,
612    access_flags: libc::c_int,
613) -> std::io::Result<std::fs::File> {
614    use std::os::unix::io::AsRawFd;
615
616    let hex = content_ref.as_str();
617    let shard1_dir = openat_dir_no_follow(root_handle.as_raw_fd(), &hex[0..2])?;
618    let shard2_dir = openat_dir_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])?;
619    openat_regular_file_no_follow(shard2_dir.as_raw_fd(), hex, access_flags)
620}
621
622#[cfg(unix)]
623fn open_or_create_dir_at_no_follow(
624    parent_fd: std::os::unix::io::RawFd,
625    name: &str,
626) -> std::io::Result<std::fs::File> {
627    let c_name = std::ffi::CString::new(name)
628        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
629    match openat_dir_no_follow(parent_fd, name) {
630        Ok(dir) => Ok(dir),
631        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
632            // SAFETY: `parent_fd` is live and `c_name` is NUL-terminated.
633            let rc = unsafe { libc::mkdirat(parent_fd, c_name.as_ptr(), 0o777) };
634            if rc != 0 {
635                let error = std::io::Error::last_os_error();
636                if error.kind() != std::io::ErrorKind::AlreadyExists {
637                    return Err(error);
638                }
639            }
640            // Always validate the resulting entry on a descriptor. Another
641            // process may have won the creation race with a non-directory or
642            // a symlink.
643            openat_dir_no_follow(parent_fd, name)
644        }
645        Err(error) => Err(error),
646    }
647}
648
649#[cfg(unix)]
650fn create_regular_file_at_no_follow(
651    parent_fd: std::os::unix::io::RawFd,
652    name: &str,
653    mode: libc::mode_t,
654) -> std::io::Result<std::fs::File> {
655    use std::os::unix::io::FromRawFd;
656
657    let c_name = std::ffi::CString::new(name)
658        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
659    // SAFETY: `parent_fd` is live, `c_name` is NUL-terminated, and the mode
660    // argument is supplied because O_CREAT is present.
661    let fd = unsafe {
662        libc::openat(
663            parent_fd,
664            c_name.as_ptr(),
665            libc::O_RDWR | libc::O_CREAT | libc::O_EXCL | libc::O_NOFOLLOW | libc::O_CLOEXEC,
666            mode as libc::c_uint,
667        )
668    };
669    if fd < 0 {
670        return Err(std::io::Error::last_os_error());
671    }
672    // SAFETY: `fd` was newly returned and is transferred exactly once.
673    Ok(unsafe { std::fs::File::from_raw_fd(fd) })
674}
675
676#[cfg(unix)]
677fn unlink_entry_at(parent_fd: std::os::unix::io::RawFd, name: &str) -> std::io::Result<()> {
678    let c_name = std::ffi::CString::new(name)
679        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
680    // SAFETY: `parent_fd` is live and `c_name` is NUL-terminated.
681    let rc = unsafe { libc::unlinkat(parent_fd, c_name.as_ptr(), 0) };
682    if rc != 0 {
683        return Err(std::io::Error::last_os_error());
684    }
685    Ok(())
686}
687
688#[cfg(unix)]
689fn rename_entry_at(
690    parent_fd: std::os::unix::io::RawFd,
691    from: &str,
692    to: &str,
693) -> std::io::Result<()> {
694    let c_from = std::ffi::CString::new(from)
695        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
696    let c_to = std::ffi::CString::new(to)
697        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
698    // SAFETY: both names are NUL-terminated and relative to the same live
699    // shard-directory handle.
700    let rc = unsafe { libc::renameat(parent_fd, c_from.as_ptr(), parent_fd, c_to.as_ptr()) };
701    if rc != 0 {
702        return Err(std::io::Error::last_os_error());
703    }
704    Ok(())
705}
706
707#[cfg(unix)]
708fn available_space_at(root_handle: &std::fs::File) -> std::io::Result<u64> {
709    use std::os::unix::io::AsRawFd;
710
711    let mut stat: libc::statvfs = unsafe { std::mem::zeroed() };
712    // SAFETY: `root_handle` is live and `stat` is a valid output buffer.
713    let rc = unsafe { libc::fstatvfs(root_handle.as_raw_fd(), &mut stat) };
714    if rc != 0 {
715        return Err(std::io::Error::last_os_error());
716    }
717    // `f_bavail`'s width is platform-dependent (u64 on Linux glibc, u32 on
718    // macOS), so the widening is a no-op on some targets and required on
719    // others; the lint only sees one target at a time.
720    #[allow(clippy::useless_conversion)]
721    Ok(stat.f_frsize.saturating_mul(u64::from(stat.f_bavail)))
722}
723
724#[cfg(unix)]
725fn acquire_root_write_lock_at(root_handle: &std::fs::File) -> StorageResult<std::fs::File> {
726    use std::os::unix::io::{AsRawFd, FromRawFd};
727
728    let c_name = std::ffi::CString::new(ROOT_WRITE_LOCK_FILE)
729        .expect("the static blob root lock name contains no NUL");
730    // SAFETY: the root handle is live, the static name is NUL-terminated,
731    // and the mode is present because O_CREAT is used.
732    let fd = unsafe {
733        libc::openat(
734            root_handle.as_raw_fd(),
735            c_name.as_ptr(),
736            libc::O_RDWR | libc::O_CREAT | libc::O_NOFOLLOW | libc::O_CLOEXEC,
737            0o666,
738        )
739    };
740    if fd < 0 {
741        return Err(map_io_err(
742            std::io::Error::last_os_error(),
743            "root_write_lock_open",
744        ));
745    }
746    // SAFETY: `fd` was newly returned and is transferred exactly once.
747    let lock_file = unsafe { std::fs::File::from_raw_fd(fd) };
748    if !lock_file
749        .metadata()
750        .map_err(|e| map_io_err(e, "root_write_lock_metadata"))?
751        .file_type()
752        .is_file()
753    {
754        return Err(map_io_err(
755            std::io::Error::new(
756                std::io::ErrorKind::InvalidInput,
757                "blob root write lock is not a regular file",
758            ),
759            "root_write_lock_open",
760        ));
761    }
762    fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
763    Ok(lock_file)
764}
765
766/// Collision-resistant diagnostic identity for one canonical blob root.
767///
768/// Paths are not required to be UTF-8. Hash the platform-native path bytes
769/// rather than a lossy display spelling. Recovery does not depend on this
770/// mutable identity: exclusive database-scoped sweep ownership makes every
771/// pre-existing claim abandoned before a new sweep starts, including claims
772/// copied by backup or left under an earlier root spelling.
773fn blob_root_key(root: &Path) -> String {
774    #[cfg(unix)]
775    let bytes = {
776        use std::os::unix::ffi::OsStrExt;
777        root.as_os_str().as_bytes().to_vec()
778    };
779    #[cfg(windows)]
780    let bytes = {
781        use std::os::windows::ffi::OsStrExt;
782        root.as_os_str()
783            .encode_wide()
784            .flat_map(u16::to_le_bytes)
785            .collect::<Vec<_>>()
786    };
787    #[cfg(not(any(unix, windows)))]
788    let bytes = root.to_string_lossy().as_bytes().to_vec();
789    blake3::hash(&bytes).to_hex().to_string()
790}
791
792/// Resolve the blob store root directory.
793///
794/// Precedence (khive#292, SPEC-gate ruling): `KHIVE_BLOB_ROOT` env var >
795/// caller-supplied `config_root` (resolved from `khive.toml` by a layer above
796/// this crate — `khive-db` cannot parse TOML itself without an upward
797/// dependency) > beside the database directory (`<db_dir>/blobs`). Errors
798/// when none apply — an in-memory backend with no override and no env var has
799/// no directory to default beside.
800pub fn resolve_blob_root(
801    db_dir: Option<&Path>,
802    config_root: Option<&Path>,
803) -> Result<PathBuf, SqliteError> {
804    if let Ok(env_root) = std::env::var("KHIVE_BLOB_ROOT") {
805        if !env_root.trim().is_empty() {
806            return Ok(PathBuf::from(env_root));
807        }
808    }
809    if let Some(root) = config_root {
810        return Ok(root.to_path_buf());
811    }
812    if let Some(dir) = db_dir {
813        return Ok(dir.join("blobs"));
814    }
815    Err(SqliteError::InvalidData(
816        "cannot resolve a blob store root: no KHIVE_BLOB_ROOT env var, no configured \
817         root, and the database has no on-disk directory to default beside (in-memory \
818         backend)"
819            .to_string(),
820    ))
821}
822
823/// Whether writing `required_write_bytes` more bytes to a volume currently
824/// reporting `available` free bytes would leave it below `floor_bytes`.
825///
826/// Pure and filesystem-independent on purpose: the
827/// exact boundary this guards — `available == floor_bytes + 1` must still
828/// refuse a 2-byte write, because a floor-only check (`available <
829/// floor_bytes`) does not account for the pending write's own size — is unit
830/// tested directly against this function rather than against the real
831/// filesystem's `fs4::available_space`, which fluctuates under concurrent
832/// build/agent activity on a shared machine and made an earlier
833/// exact-boundary integration test flaky. `saturating_sub` avoids underflow
834/// when `required_write_bytes` exceeds `available` outright — that case
835/// still correctly refuses for any nonzero floor.
836fn crosses_floor(available: u64, required_write_bytes: u64, floor_bytes: u64) -> bool {
837    available.saturating_sub(required_write_bytes) < floor_bytes
838}
839
840#[cfg(any(test, not(unix)))]
841fn put_blocking_with_space_probe<F>(
842    root: &Path,
843    floor_bytes: u64,
844    bytes: Vec<u8>,
845    available_space: F,
846) -> StorageResult<ContentRef>
847where
848    F: FnOnce(&Path) -> std::io::Result<u64>,
849{
850    let digest = blake3::hash(&bytes);
851    let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
852    let target = shard_path(root, &content_ref);
853
854    // Content-addressed: identical bytes already on disk means this put is a
855    // no-op (BlobStore::put's documented dedup contract) — skip the floor
856    // check and the write entirely. The existing file's mtime is still
857    // refreshed to now: a prior orphan re-published through this path
858    // restarts its publish-grace clock exactly as a fresh write would,
859    // rather than keeping a stale mtime that lets the orphan sweep delete it
860    // out from under the caller's follow-up attachment write (khive#1313). The
861    // caller already holds both the async and file-based publish advisory
862    // locks for the duration of this call, so the refresh is serialized
863    // against a concurrent sweep the same way an ordinary write is.
864    if target.exists() {
865        let file = fs::OpenOptions::new()
866            .write(true)
867            .open(&target)
868            .map_err(|e| map_io_err(e, "put_touch_open"))?;
869        file.set_modified(SystemTime::now())
870            .map_err(|e| map_io_err(e, "put_touch_mtime"))?;
871        return Ok(content_ref);
872    }
873
874    let required_write_bytes = bytes.len() as u64;
875    let available = available_space(root).map_err(|e| map_io_err(e, "put_check_space"))?;
876    if crosses_floor(available, required_write_bytes, floor_bytes) {
877        return Err(StorageError::CapacityFloor {
878            capability: StorageCapability::Blob,
879            volume: root.display().to_string(),
880            available_bytes: available,
881            floor_bytes,
882        });
883    }
884
885    let shard_dir = target
886        .parent()
887        .expect("shard_path always nests under two directory levels");
888    fs::create_dir_all(shard_dir).map_err(|e| map_io_err(e, "put_mkdir"))?;
889
890    let mut tmp = tempfile::Builder::new()
891        .prefix(".tmp-")
892        .tempfile_in(shard_dir)
893        .map_err(|e| map_io_err(e, "put_tempfile"))?;
894    tmp.write_all(&bytes)
895        .map_err(|e| map_io_err(e, "put_write"))?;
896    tmp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
897    tmp.as_file()
898        .sync_all()
899        .map_err(|e| map_io_err(e, "put_fsync"))?;
900
901    let written_len = tmp
902        .as_file()
903        .metadata()
904        .map_err(|e| map_io_err(e, "put_verify"))?
905        .len();
906    if written_len != bytes.len() as u64 {
907        return Err(map_io_err(
908            std::io::Error::other(format!(
909                "temp file length {written_len} does not match {} written bytes",
910                bytes.len()
911            )),
912            "put_verify",
913        ));
914    }
915
916    tmp.persist(&target)
917        .map_err(|e| map_io_err(e.error, "put_persist"))?;
918
919    Ok(content_ref)
920}
921
922#[cfg(any(test, not(unix)))]
923fn put_blocking(root: &Path, floor_bytes: u64, bytes: Vec<u8>) -> StorageResult<ContentRef> {
924    let _root_write_guard = acquire_root_write_lock(root)?;
925    put_blocking_with_space_probe(root, floor_bytes, bytes, |path| fs4::available_space(path))
926}
927
928fn acquire_root_write_lock_anchored(
929    root: &Path,
930    root_handle: &std::fs::File,
931) -> StorageResult<std::fs::File> {
932    verify_blob_root_identity(root, root_handle)
933        .map_err(|e| map_io_err(e, "root_write_lock_identity"))?;
934    #[cfg(unix)]
935    {
936        acquire_root_write_lock_at(root_handle)
937    }
938    #[cfg(not(unix))]
939    {
940        acquire_root_write_lock(root)
941    }
942}
943
944#[cfg(unix)]
945fn put_blocking_from_root_handle(
946    root: &Path,
947    root_handle: &std::fs::File,
948    floor_bytes: u64,
949    bytes: Vec<u8>,
950) -> StorageResult<ContentRef> {
951    use std::os::unix::io::AsRawFd;
952
953    let _root_write_guard = acquire_root_write_lock_anchored(root, root_handle)?;
954    let digest = blake3::hash(&bytes);
955    let content_ref = ContentRef::from_digest_bytes(digest.as_bytes());
956
957    // Preserve the dedup/republish contract while opening the existing leaf
958    // relative to the retained root. A missing shard level and a missing leaf
959    // are both the ordinary publish path; every other traversal failure is a
960    // hard refusal.
961    match open_blob_shard_file_at_no_follow(root_handle, &content_ref, libc::O_WRONLY) {
962        Ok(file) => {
963            file.set_modified(SystemTime::now())
964                .map_err(|e| map_io_err(e, "put_touch_mtime"))?;
965            return Ok(content_ref);
966        }
967        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
968        Err(error) => return Err(map_io_err(error, "put_touch_open")),
969    }
970
971    let required_write_bytes = bytes.len() as u64;
972    let available =
973        available_space_at(root_handle).map_err(|e| map_io_err(e, "put_check_space"))?;
974    if crosses_floor(available, required_write_bytes, floor_bytes) {
975        return Err(StorageError::CapacityFloor {
976            capability: StorageCapability::Blob,
977            volume: root.display().to_string(),
978            available_bytes: available,
979            floor_bytes,
980        });
981    }
982
983    let hex = content_ref.as_str();
984    let shard1_dir = open_or_create_dir_at_no_follow(root_handle.as_raw_fd(), &hex[0..2])
985        .map_err(|e| map_io_err(e, "put_mkdir"))?;
986    let shard2_dir = open_or_create_dir_at_no_follow(shard1_dir.as_raw_fd(), &hex[2..4])
987        .map_err(|e| map_io_err(e, "put_mkdir"))?;
988
989    let temp_name = format!(".tmp-{}", Uuid::new_v4());
990    let mut temp = create_regular_file_at_no_follow(shard2_dir.as_raw_fd(), &temp_name, 0o600)
991        .map_err(|e| map_io_err(e, "put_tempfile"))?;
992    let write_result = (|| -> StorageResult<()> {
993        temp.write_all(&bytes)
994            .map_err(|e| map_io_err(e, "put_write"))?;
995        temp.flush().map_err(|e| map_io_err(e, "put_flush"))?;
996        temp.sync_all().map_err(|e| map_io_err(e, "put_fsync"))?;
997
998        let written_len = temp
999            .metadata()
1000            .map_err(|e| map_io_err(e, "put_verify"))?
1001            .len();
1002        if written_len != bytes.len() as u64 {
1003            return Err(map_io_err(
1004                std::io::Error::other(format!(
1005                    "temp file length {written_len} does not match {} written bytes",
1006                    bytes.len()
1007                )),
1008                "put_verify",
1009            ));
1010        }
1011        Ok(())
1012    })();
1013    drop(temp);
1014    if let Err(error) = write_result {
1015        let _ = unlink_entry_at(shard2_dir.as_raw_fd(), &temp_name);
1016        return Err(error);
1017    }
1018
1019    if let Err(error) = rename_entry_at(shard2_dir.as_raw_fd(), &temp_name, hex) {
1020        let _ = unlink_entry_at(shard2_dir.as_raw_fd(), &temp_name);
1021        return Err(map_io_err(error, "put_persist"));
1022    }
1023    Ok(content_ref)
1024}
1025
1026#[cfg(not(unix))]
1027fn put_blocking_from_root_handle(
1028    root: &Path,
1029    root_handle: &std::fs::File,
1030    floor_bytes: u64,
1031    bytes: Vec<u8>,
1032) -> StorageResult<ContentRef> {
1033    verify_blob_root_identity(root, root_handle).map_err(|e| map_io_err(e, "put_root_identity"))?;
1034    put_blocking(root, floor_bytes, bytes)
1035}
1036
1037#[cfg(any(test, not(unix)))]
1038fn acquire_root_write_lock(root: &Path) -> StorageResult<fs::File> {
1039    let lock_file = fs::OpenOptions::new()
1040        .read(true)
1041        .write(true)
1042        .create(true)
1043        .truncate(false)
1044        .open(root.join(ROOT_WRITE_LOCK_FILE))
1045        .map_err(|e| map_io_err(e, "root_write_lock_open"))?;
1046    fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "root_write_lock_acquire"))?;
1047    Ok(lock_file)
1048}
1049
1050fn database_gc_lock_path(database_path: &Path) -> PathBuf {
1051    let mut lock_path = database_path.as_os_str().to_os_string();
1052    lock_path.push(DATABASE_GC_LOCK_SUFFIX);
1053    PathBuf::from(lock_path)
1054}
1055
1056fn acquire_database_gc_lock(database_path: Option<&Path>) -> StorageResult<Option<fs::File>> {
1057    let Some(database_path) = database_path else {
1058        // An in-memory database cannot be shared by another process. The
1059        // process-local lock below is therefore the complete owner fence.
1060        return Ok(None);
1061    };
1062    let lock_path = database_gc_lock_path(database_path);
1063    let lock_file = fs::OpenOptions::new()
1064        .read(true)
1065        .write(true)
1066        .create(true)
1067        .truncate(false)
1068        .open(&lock_path)
1069        .map_err(|e| map_io_err(e, "database_gc_lock_open"))?;
1070    fs4::FileExt::lock(&lock_file).map_err(|e| map_io_err(e, "database_gc_lock_acquire"))?;
1071    Ok(Some(lock_file))
1072}
1073
1074#[cfg(all(unix, target_os = "macos"))]
1075fn errno_location() -> *mut libc::c_int {
1076    // SAFETY: `__error` returns this thread's errno cell; obtaining the
1077    // pointer has no preconditions beyond running on a thread, which every
1078    // caller here does.
1079    unsafe { libc::__error() }
1080}
1081
1082#[cfg(all(unix, not(target_os = "macos")))]
1083fn errno_location() -> *mut libc::c_int {
1084    // SAFETY: see the macOS arm above; `__errno_location` is the Linux/glibc
1085    // equivalent thread-local errno accessor.
1086    unsafe { libc::__errno_location() }
1087}
1088
1089/// Zero `errno` on the current thread. `readdir` never clears `errno` itself
1090/// on success, so this must run immediately before each call for the
1091/// NULL-return ambiguity below to be resolvable afterward.
1092#[cfg(unix)]
1093fn clear_errno() {
1094    // SAFETY: `errno_location` returns a valid, live thread-local `c_int`
1095    // cell for the duration of this call.
1096    unsafe { *errno_location() = 0 };
1097}
1098
1099#[cfg(unix)]
1100fn current_errno() -> libc::c_int {
1101    // SAFETY: see `clear_errno`.
1102    unsafe { *errno_location() }
1103}
1104
1105/// List non-dot entry names of an open directory descriptor via `fdopendir`
1106/// on an INDEPENDENTLY reopened fd for the same directory (`openat(dir_fd,
1107/// ".", O_NOFOLLOW)`, never `dup(dir_fd)`) — the caller's descriptor stays
1108/// owned by the caller and, just as importantly, keeps its own read
1109/// position. `dup` shares the underlying open file description, and with it
1110/// the directory-stream position, with the fd it was duplicated from; a
1111/// persisted, repeatedly-listed handle (like `FsBlobStore::root_handle`)
1112/// would silently read as empty on its second call once a `dup`'d
1113/// `fdopendir`/`readdir` pass had already driven that shared position to
1114/// EOF. `openat(..., ".", ...)` yields a genuinely new open file
1115/// description, so this listing never perturbs the position of `dir_fd`
1116/// itself. `.`/`..` are excluded so a handle-relative walk built on this can
1117/// never step to a directory's parent or re-enter itself.
1118///
1119/// `readdir`'s NULL return is POSIX-ambiguous between end-of-stream and a
1120/// read error, distinguished only by `errno`: EOF leaves it unchanged (`0`,
1121/// since this loop always clears it first), an error sets it non-zero. This
1122/// walk clears `errno` before every `readdir` call and checks it on a NULL
1123/// return, so a mid-stream read error is reported as `Err` instead of
1124/// silently truncating the listing as if the stream had simply ended — the
1125/// distinction `transactional_orphan_sweep`'s candidate walk depends on to
1126/// avoid treating a partial, error-truncated scan as a complete one. The
1127/// `DIR*` stream is closed exactly once on every return path (`Ok`, the
1128/// mid-stream error path, and the pre-loop `fdopendir` failure).
1129#[cfg(unix)]
1130fn read_dir_names_no_follow(dir_fd: std::os::unix::io::RawFd) -> std::io::Result<Vec<String>> {
1131    use std::os::unix::io::IntoRawFd;
1132
1133    let reopened = openat_dir_no_follow(dir_fd, ".")?;
1134    // SAFETY: `reopened` was just opened above and is uniquely owned; its
1135    // raw fd is handed to `fdopendir`, which takes ownership of it on
1136    // success (and closes it via `closedir` below).
1137    let owned_fd = reopened.into_raw_fd();
1138    // SAFETY: `owned_fd` is valid and uniquely owned; `fdopendir` takes
1139    // ownership of it on success.
1140    let dirp = unsafe { libc::fdopendir(owned_fd) };
1141    if dirp.is_null() {
1142        let err = std::io::Error::last_os_error();
1143        // SAFETY: `owned_fd` is still owned by us since `fdopendir` failed.
1144        unsafe { libc::close(owned_fd) };
1145        return Err(err);
1146    }
1147    let mut names = Vec::new();
1148    loop {
1149        // `readdir` leaves `errno` untouched on EOF; clearing it here is
1150        // what makes that observable as distinct from an error below.
1151        clear_errno();
1152        // SAFETY: `dirp` is a valid, open `DIR*` for this whole loop.
1153        let entry = unsafe { libc::readdir(dirp) };
1154        if entry.is_null() {
1155            if current_errno() != 0 {
1156                let err = std::io::Error::last_os_error();
1157                // SAFETY: `dirp` is still open and owned by this call; this
1158                // is the one closedir on the error path, matching the one
1159                // closedir on the `Ok` path below.
1160                unsafe { libc::closedir(dirp) };
1161                return Err(err);
1162            }
1163            break;
1164        }
1165        // SAFETY: `d_name` is NUL-terminated, so its first byte is always
1166        // in bounds; `.`/`..` (and any other dot-leading entry, e.g. an
1167        // in-flight `.tmp-*` file or the root write-lock file) are rejected
1168        // on this raw byte before any allocation happens for them.
1169        let first = unsafe { *(*entry).d_name.as_ptr() };
1170        if first == b'.' as libc::c_char {
1171            continue;
1172        }
1173        // SAFETY: `entry` is valid until the next `readdir`/`closedir`
1174        // call; the name is copied out before either.
1175        let name = unsafe { std::ffi::CStr::from_ptr((*entry).d_name.as_ptr()) }
1176            .to_string_lossy()
1177            .into_owned();
1178        names.push(name);
1179    }
1180    // SAFETY: `dirp` was successfully opened above and not yet closed.
1181    unsafe { libc::closedir(dirp) };
1182    Ok(names)
1183}
1184
1185/// Enumerate orphan-sweep candidates and their mtimes entirely relative to
1186/// the retained `root_handle` — no path is ever re-resolved from `root` for
1187/// this walk, so a concurrent replacement of the root or a shard directory
1188/// cannot redirect candidate enumeration or the mtime read used for
1189/// `within_publish_grace` classification outside the retained root.
1190///
1191/// Each shard level is opened with `openat(..., O_NOFOLLOW)` relative to the
1192/// previous, already-verified descriptor (same helpers `put`/`get`/`delete`
1193/// use), so a symlink planted at either shard level is refused rather than
1194/// followed and its contents skipped, never swept. A candidate's mtime is
1195/// read via `fstat` on the handle returned by opening the leaf with
1196/// `openat(..., O_NOFOLLOW)`, never via `fs::metadata` on a path. If the
1197/// leaf cannot be opened as a verified regular file it is dropped entirely
1198/// (nothing to protect — it is not a delete candidate); if it opens but its
1199/// metadata/mtime cannot be read, it is kept as a candidate with an unknown
1200/// mtime, which `within_publish_grace` treats as protected — the safe
1201/// direction for a sweep that only ever destroys data.
1202#[cfg(unix)]
1203fn walk_blob_files_from_root_handle(
1204    root_handle: &std::fs::File,
1205    // Only used under `#[cfg(test)]`, to key `walk_leaf_sync_hook` by this
1206    // walk's canonical root; underscore-prefixed so non-test builds don't
1207    // warn on the unused parameter.
1208    _root: &Path,
1209) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
1210    use std::os::unix::io::AsRawFd;
1211
1212    let mut out = Vec::new();
1213    for l1_name in read_dir_names_no_follow(root_handle.as_raw_fd())? {
1214        let l1_dir = match openat_dir_no_follow(root_handle.as_raw_fd(), &l1_name) {
1215            Ok(dir) => dir,
1216            Err(_) => continue,
1217        };
1218        for l2_name in read_dir_names_no_follow(l1_dir.as_raw_fd())? {
1219            let l2_dir = match openat_dir_no_follow(l1_dir.as_raw_fd(), &l2_name) {
1220                Ok(dir) => dir,
1221                Err(_) => continue,
1222            };
1223            for leaf_name in read_dir_names_no_follow(l2_dir.as_raw_fd())? {
1224                // Non-hex names never round-trip through `ContentRef`;
1225                // orphan_sweep only ever acts on names that do.
1226                let Ok(content_ref) = ContentRef::from_hex(leaf_name.clone()) else {
1227                    continue;
1228                };
1229                #[cfg(all(test, unix))]
1230                if let Some(hook) = walk_leaf_sync_hook::take(_root) {
1231                    let _ = hook.reached.send(());
1232                    let _ = hook.release.recv();
1233                }
1234                let file = match openat_regular_file_no_follow(
1235                    l2_dir.as_raw_fd(),
1236                    &leaf_name,
1237                    libc::O_RDONLY,
1238                ) {
1239                    Ok(file) => file,
1240                    Err(_) => continue,
1241                };
1242                let mtime = file.metadata().ok().and_then(|meta| meta.modified().ok());
1243                out.push((content_ref, mtime));
1244            }
1245        }
1246    }
1247    Ok(out)
1248}
1249
1250/// Non-Unix tier has no descriptor-relative directory-listing API in this
1251/// codebase (see the equivalent fail-closed note on `unlink_blob_shard_file_no_follow`'s
1252/// `not(any(unix, windows))` arm). Rather than fall back to path-based
1253/// `fs::read_dir`/`fs::metadata` reads — which reintroduces exactly the
1254/// TOCTOU this function exists to close — candidate enumeration refuses
1255/// outright, and `transactional_orphan_sweep` surfaces that as a sweep
1256/// failure instead of classifying anything.
1257#[cfg(not(unix))]
1258fn walk_blob_files_from_root_handle(
1259    _root_handle: &std::fs::File,
1260    _root: &Path,
1261) -> std::io::Result<Vec<(ContentRef, Option<SystemTime>)>> {
1262    Err(std::io::Error::new(
1263        std::io::ErrorKind::Unsupported,
1264        "orphan sweep candidate enumeration requires descriptor-relative directory \
1265         reads, available only on unix in this release; refusing to classify via \
1266         path-based reads",
1267    ))
1268}
1269
1270/// Whether a candidate file is still inside its publish grace period and must
1271/// be left alone regardless of liveness.
1272///
1273/// `put`'s two-step client protocol (bytes land first, a *later* attachment write
1274/// commits the `content_ref`) means a blob can be physically on disk with
1275/// zero live references for a window entirely outside this store's control —
1276/// the referencing write simply hasn't happened yet. A file whose mtime is
1277/// younger than `grace_period` is therefore treated as not-yet-orphaned: an
1278/// unreadable mtime (removed mid-scan, clock weirdness) is treated the same
1279/// way (age unknown -> protect it), the safe direction for a sweep that only
1280/// ever destroys data.
1281fn within_publish_grace(
1282    mtime: Option<SystemTime>,
1283    now: SystemTime,
1284    grace_period: Duration,
1285) -> bool {
1286    let age = mtime.and_then(|mtime| now.duration_since(mtime).ok());
1287    match age {
1288        Some(age) => age < grace_period,
1289        None => true,
1290    }
1291}
1292
1293#[derive(Debug)]
1294struct PreparedTransactionalSweep {
1295    result: BlobOrphanSweepResult,
1296    candidates: Vec<(ContentRef, bool)>,
1297}
1298
1299/// Perform every filesystem-dependent part of candidate classification before
1300/// SQLite's writer transaction opens.
1301fn prepare_transactional_sweep(
1302    files: Vec<(ContentRef, Option<SystemTime>)>,
1303    grace_period: Duration,
1304) -> PreparedTransactionalSweep {
1305    let now = SystemTime::now();
1306    let mut result = BlobOrphanSweepResult::default();
1307    let mut candidates = Vec::with_capacity(files.len());
1308    for (content_ref, mtime) in files {
1309        result.scanned += 1;
1310        let within_grace = within_publish_grace(mtime, now, grace_period);
1311        candidates.push((content_ref, within_grace));
1312    }
1313    PreparedTransactionalSweep { result, candidates }
1314}
1315
1316#[derive(Debug)]
1317struct BlobGcBatchRows {
1318    grace_period_skipped: u64,
1319    would_delete: u64,
1320    claimed_rows: Vec<SqlRow>,
1321}
1322
1323fn required_nonnegative_count(
1324    value: Option<SqlValue>,
1325    operation: &'static str,
1326) -> StorageResult<u64> {
1327    match value {
1328        Some(SqlValue::Integer(value)) if value >= 0 => Ok(value as u64),
1329        other => Err(StorageError::Internal(format!(
1330            "{operation} returned an invalid count: {other:?}"
1331        ))),
1332    }
1333}
1334
1335fn invalid_content_ref(message: String) -> StorageError {
1336    StorageError::InvalidInput {
1337        capability: StorageCapability::Blob,
1338        operation: "transactional_orphan_sweep".into(),
1339        message,
1340    }
1341}
1342
1343/// Whether this database carries the complete V21 attachment-only GC fencing
1344/// set and durable completed cutover marker.
1345///
1346/// `transactional_orphan_sweep` is reachable from any `SqlAccess` a caller
1347/// hands it, including a `StorageBackend` constructed directly (e.g.
1348/// `StorageBackend::memory()`/`sqlite()` used without `prepare_core_schema`)
1349/// that never ran core migrations. The triggers are the fence that keeps a
1350/// concurrent attachment write from resurrecting a claimed digest in the
1351/// released-writer window, so a database missing any element of the set
1352/// cannot satisfy the fail-closed guarantee the
1353/// [`BlobStore::transactional_orphan_sweep`] contract requires; the sweep
1354/// refuses with [`StorageError::Unsupported`] rather than degrading to
1355/// unfenced deletion.
1356async fn blob_gc_fencing_complete(sql: &dyn SqlAccess) -> StorageResult<bool> {
1357    let mut reader = sql.reader().await?;
1358    let present = required_nonnegative_count(
1359        reader
1360            .query_scalar(SqlStatement {
1361                sql: "SELECT COUNT(*) FROM sqlite_master \
1362                      WHERE (type = 'table' AND name IN ( \
1363                                 'blob_gc_claims', 'attachments', \
1364                                 'attachment_cutover_state')) \
1365                         OR (type = 'index' AND name IN ( \
1366                             'idx_blob_gc_claims_content_ref', \
1367                             'idx_attachments_content_ref')) \
1368                         OR (type = 'trigger' AND name IN ( \
1369                             'attachments_reject_claimed_blob_insert', \
1370                             'attachments_reject_claimed_blob_update'))"
1371                    .to_string(),
1372                params: vec![],
1373                label: Some("blob_gc_fencing_complete".to_string()),
1374            })
1375            .await?,
1376        "blob_gc_fencing_complete",
1377    )?;
1378    if present != 7 {
1379        return Ok(false);
1380    }
1381
1382    let legacy_objects = required_nonnegative_count(
1383        reader
1384            .query_scalar(SqlStatement {
1385                sql: "SELECT \
1386                        (SELECT COUNT(*) FROM pragma_table_info('entities') \
1387                         WHERE name = 'content_ref') \
1388                      + (SELECT COUNT(*) FROM sqlite_master \
1389                         WHERE (type = 'index' AND name = 'idx_entities_content_ref') \
1390                            OR (type = 'trigger' AND name IN ( \
1391                                'entities_reject_claimed_blob_insert', \
1392                                'entities_reject_claimed_blob_update')))"
1393                    .to_string(),
1394                params: vec![],
1395                label: Some("blob_gc_legacy_fencing_absent".to_string()),
1396            })
1397            .await?,
1398        "blob_gc_legacy_fencing_absent",
1399    )?;
1400    if legacy_objects != 0 {
1401        return Ok(false);
1402    }
1403
1404    // The sweep is admitted only for the EXACT completed V21 epoch
1405    // (ADR-160 Phase 4a: "report-only and destructive sweeps only for an
1406    // exact completed V21 epoch … ahead-of-V21 epochs return typed
1407    // Unsupported"). A ledger above V21 — whether ahead of this binary or a
1408    // migration this same binary applied — belongs to a schema epoch whose
1409    // attachment/liveness semantics this gate never validated, so it fails
1410    // closed; the author of a future migration extends the gate in the same
1411    // change that proves the new epoch's liveness set, never by default.
1412    // General schema validation (`attachment_cutover_status`) deliberately
1413    // keeps accepting later versions on top of a completed cutover: that is
1414    // the normal serving course, and the exact-epoch rule is scoped to
1415    // destructive GC admission.
1416    //
1417    // "Exact completed V21 epoch" means the WHOLE canonical ledger, not a
1418    // V21 terminal row: `version` is the table's PRIMARY KEY, so
1419    // COUNT(*) = 21 ∧ MIN = 1 ∧ MAX = 21 holds if and only if the ledger is
1420    // exactly the contiguous set {1..21}. A ledger that merely retains a V21
1421    // row at MAX(version) = 21 while earlier rows are missing is an
1422    // incomplete migration history whose physical schema this gate never
1423    // validated — it must fail closed, same as ahead-of-V21. Name-level
1424    // canonicality of the below-terminal rows stays boot's job
1425    // (`validate_applied_migration_ledger`); this predicate enforces the
1426    // structural contiguity a destructive sweep's admission rests on, plus
1427    // the named V21 row itself.
1428    let complete = required_nonnegative_count(
1429        reader
1430            .query_scalar(SqlStatement {
1431                sql: "SELECT COUNT(*) FROM attachment_cutover_state AS cutover \
1432                      WHERE cutover.singleton = 1 \
1433                        AND cutover.state = 'complete' \
1434                        AND cutover.completed_at IS NOT NULL \
1435                        AND (SELECT COUNT(*) FROM _schema_migrations \
1436                             WHERE version = ?1 \
1437                               AND name = 'attachments_first_class') = 1 \
1438                        AND (SELECT COUNT(*) FROM _schema_migrations) = ?1 \
1439                        AND (SELECT MIN(version) FROM _schema_migrations) = 1 \
1440                        AND (SELECT MAX(version) FROM _schema_migrations) = ?1"
1441                    .to_string(),
1442                params: vec![SqlValue::Integer(i64::from(
1443                    crate::migrations::ATTACHMENT_CUTOVER_VERSION,
1444                ))],
1445                label: Some("blob_gc_cutover_complete".to_string()),
1446            })
1447            .await?,
1448        "blob_gc_cutover_complete",
1449    )?;
1450    Ok(complete == 1)
1451}
1452
1453fn unsupported_blob_gc_epoch() -> StorageError {
1454    StorageError::Unsupported {
1455        capability: StorageCapability::Blob,
1456        operation: "transactional_orphan_sweep".into(),
1457        message: "transactional blob GC requires a complete V21 attachment cutover with \
1458                  the attachment claim-fencing set; refusing both report-only and \
1459                  destructive sweep in this database epoch"
1460            .into(),
1461    }
1462}
1463
1464/// The sentinel digest the fence probe claims. All zeros is canonical-form
1465/// valid (64 lowercase hex) and unreachable as a real BLAKE3 digest for any
1466/// stored object in practice; probe rows never survive the probe transaction.
1467const BLOB_GC_FENCE_PROBE_REF: &str =
1468    "0000000000000000000000000000000000000000000000000000000000000000";
1469
1470/// The RAISE(ABORT) message shared by the V20 and V21 fencing triggers. The probe
1471/// requires the rejection to be OUR fence, not an incidental failure.
1472const BLOB_GC_FENCE_TRIGGER_MESSAGE: &str = "content_ref is reserved by an active blob sweep";
1473
1474/// Prove the V21 fence actually fences, not merely that objects with the
1475/// right NAMES exist in `sqlite_master`. Same-named no-op triggers (or a
1476/// rewritten trigger body) would pass the name census while letting a
1477/// claimed `content_ref` become live in the released-writer window, so the
1478/// gate exercises the fence: inside one writer transaction it claims the
1479/// all-zero sentinel AND a second random digest, and attempts the attachment
1480/// INSERT and attachment UPDATE the triggers must reject for EACH claimed
1481/// digest — with the second digest's arms using a different attachment shape
1482/// (substrate `note`, role `evidence`), so a trigger rewrite restricted to
1483/// one digest, substrate, or role fails an arm instead of passing a
1484/// fixed-sentinel census. Every arm must fail with the triggers' own RAISE
1485/// message, and every probe row is deleted before the unit commits. Any
1486/// other outcome — a write accepted, or rejected for a different reason —
1487/// refuses the sweep with [`StorageError::Unsupported`].
1488async fn blob_gc_fence_probe(sql: &dyn SqlAccess) -> StorageResult<()> {
1489    let run = Uuid::new_v4().simple().to_string();
1490    blob_gc_fence_probe_with_ids(
1491        sql,
1492        format!("__blob-gc-fence-probe-insert-{run}__"),
1493        format!("__blob-gc-fence-probe-update-{run}__"),
1494        format!("__blob-gc-fence-probe-insert2-{run}__"),
1495        format!("__blob-gc-fence-probe-update2-{run}__"),
1496        format!("__fence_probe-{run}__"),
1497    )
1498    .await
1499}
1500
1501/// Probe body with explicit row ids so tests can force an id collision.
1502/// Production callers go through [`blob_gc_fence_probe`], which mints
1503/// per-run random ids; the guard below still refuses to run — touching
1504/// nothing — if any minted id already names a row.
1505async fn blob_gc_fence_probe_with_ids(
1506    sql: &dyn SqlAccess,
1507    insert_id: String,
1508    update_id: String,
1509    insert2_id: String,
1510    update2_id: String,
1511    claim_key: String,
1512) -> StorageResult<()> {
1513    fn fence_rejection(result: Result<u64, StorageError>) -> Result<bool, String> {
1514        match result {
1515            Ok(_) => Ok(false),
1516            Err(error) => {
1517                let text = error.to_string();
1518                if text.contains(BLOB_GC_FENCE_TRIGGER_MESSAGE) {
1519                    Ok(true)
1520                } else {
1521                    Err(text)
1522                }
1523            }
1524        }
1525    }
1526
1527    fn required_seed(value: Option<SqlValue>) -> StorageResult<String> {
1528        match value {
1529            Some(SqlValue::Text(seed)) => Ok(seed),
1530            _ => Err(StorageError::Unsupported {
1531                capability: StorageCapability::Blob,
1532                operation: "transactional_orphan_sweep".into(),
1533                message: "the blob GC fence probe could not select an unclaimed \
1534                          canonical seed; refusing deletion so a later sweep can retry"
1535                    .into(),
1536            }),
1537        }
1538    }
1539
1540    let op: AtomicUnitOp = Box::new(move |writer| {
1541        Box::pin(async move {
1542            // Ownership guard: the cleanup below deletes these ids
1543            // unconditionally, so the probe may only proceed when it can
1544            // prove every id is unclaimed in EVERY table cleanup touches.
1545            let preexisting = writer
1546                .query_row(SqlStatement {
1547                    sql: "SELECT (SELECT COUNT(*) FROM attachments \
1548                                   WHERE record_uuid IN (?1, ?2, ?3, ?4)) \
1549                              + (SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?5)"
1550                        .to_string(),
1551                    params: vec![
1552                        SqlValue::Text(insert_id.clone()),
1553                        SqlValue::Text(update_id.clone()),
1554                        SqlValue::Text(insert2_id.clone()),
1555                        SqlValue::Text(update2_id.clone()),
1556                        SqlValue::Text(claim_key.clone()),
1557                    ],
1558                    label: Some("blob_gc_fence_probe_ownership_guard".to_string()),
1559                })
1560                .await?
1561                .and_then(|row| row.columns.first().map(|c| c.value.clone()));
1562            match preexisting {
1563                Some(SqlValue::Integer(0)) => {}
1564                Some(SqlValue::Integer(_)) => {
1565                    return Err(StorageError::Unsupported {
1566                        capability: StorageCapability::Blob,
1567                        operation: "transactional_orphan_sweep".into(),
1568                        message: "the blob GC fence probe's row ids collide with existing \
1569                                  rows; refusing to probe rather than delete data the \
1570                                  probe does not own"
1571                            .into(),
1572                    });
1573                }
1574                _ => {
1575                    return Err(StorageError::Internal(
1576                        "blob GC fence probe ownership guard returned no count".into(),
1577                    ));
1578                }
1579            }
1580
1581            // The UPDATE arm needs a valid, initially unclaimed reference.
1582            // Select it under this same writer transaction instead of using a
1583            // fixed sentinel that a recoverable abandoned claim could fence
1584            // forever. Eight fresh candidates keep collision handling bounded;
1585            // no candidate means a safe, retryable refusal.
1586            let seed_ref = writer
1587                .query_row(SqlStatement {
1588                    sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1589                              SELECT 1, lower(hex(randomblob(32))) \
1590                              UNION ALL \
1591                              SELECT attempt + 1, lower(hex(randomblob(32))) \
1592                              FROM candidates WHERE attempt < 8 \
1593                          ) \
1594                          SELECT candidate.content_ref FROM candidates AS candidate \
1595                          WHERE candidate.content_ref <> ?1 \
1596                            AND NOT EXISTS ( \
1597                                SELECT 1 FROM blob_gc_claims \
1598                                WHERE content_ref = candidate.content_ref \
1599                            ) \
1600                          LIMIT 1"
1601                        .to_string(),
1602                    params: vec![SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string())],
1603                    label: Some("blob_gc_fence_probe_select_seed".to_string()),
1604                })
1605                .await?
1606                .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1607            let seed_ref = required_seed(seed_ref)?;
1608
1609            let seed2_ref = writer
1610                .query_row(SqlStatement {
1611                    sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1612                              SELECT 1, lower(hex(randomblob(32))) \
1613                              UNION ALL \
1614                              SELECT attempt + 1, lower(hex(randomblob(32))) \
1615                              FROM candidates WHERE attempt < 8 \
1616                          ) \
1617                          SELECT candidate.content_ref FROM candidates AS candidate \
1618                          WHERE candidate.content_ref NOT IN (?1, ?2) \
1619                            AND NOT EXISTS ( \
1620                                SELECT 1 FROM blob_gc_claims \
1621                                WHERE content_ref = candidate.content_ref \
1622                            ) \
1623                          LIMIT 1"
1624                        .to_string(),
1625                    params: vec![
1626                        SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1627                        SqlValue::Text(seed_ref.clone()),
1628                    ],
1629                    label: Some("blob_gc_fence_probe_select_seed2".to_string()),
1630                })
1631                .await?
1632                .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1633            let seed2_ref = required_seed(seed2_ref)?;
1634
1635            // The second CLAIMED digest. A trigger rewrite conditioned on the
1636            // fixed all-zero sentinel passes that sentinel's arms; this digest
1637            // is random per run, so such a rewrite fails the arms below.
1638            let probe2_ref = writer
1639                .query_row(SqlStatement {
1640                    sql: "WITH RECURSIVE candidates(attempt, content_ref) AS ( \
1641                              SELECT 1, lower(hex(randomblob(32))) \
1642                              UNION ALL \
1643                              SELECT attempt + 1, lower(hex(randomblob(32))) \
1644                              FROM candidates WHERE attempt < 8 \
1645                          ) \
1646                          SELECT candidate.content_ref FROM candidates AS candidate \
1647                          WHERE candidate.content_ref NOT IN (?1, ?2, ?3) \
1648                            AND NOT EXISTS ( \
1649                                SELECT 1 FROM blob_gc_claims \
1650                                WHERE content_ref = candidate.content_ref \
1651                            ) \
1652                          LIMIT 1"
1653                        .to_string(),
1654                    params: vec![
1655                        SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1656                        SqlValue::Text(seed_ref.clone()),
1657                        SqlValue::Text(seed2_ref.clone()),
1658                    ],
1659                    label: Some("blob_gc_fence_probe_select_probe2".to_string()),
1660                })
1661                .await?
1662                .and_then(|row| row.columns.first().map(|column| column.value.clone()));
1663            let probe2_ref = required_seed(probe2_ref)?;
1664
1665            writer
1666                .execute(SqlStatement {
1667                    sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
1668                          VALUES (?1, ?2, 0), (?1, ?3, 0)"
1669                        .to_string(),
1670                    params: vec![
1671                        SqlValue::Text(claim_key.clone()),
1672                        SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1673                        SqlValue::Text(probe2_ref.clone()),
1674                    ],
1675                    label: Some("blob_gc_fence_probe_claim".to_string()),
1676                })
1677                .await?;
1678
1679            let insert_attempt = writer
1680                .execute(SqlStatement {
1681                    sql: "INSERT INTO attachments \
1682                          (record_uuid, substrate, role, content_ref, created_at) \
1683                          VALUES (?1, 'entity', 'content', ?2, 0)"
1684                        .to_string(),
1685                    params: vec![
1686                        SqlValue::Text(insert_id.clone()),
1687                        SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1688                    ],
1689                    label: Some("blob_gc_fence_probe_insert_arm".to_string()),
1690                })
1691                .await;
1692            let insert_fenced = fence_rejection(insert_attempt);
1693
1694            writer
1695                .execute(SqlStatement {
1696                    sql: "INSERT INTO attachments \
1697                          (record_uuid, substrate, role, content_ref, created_at) \
1698                          VALUES (?1, 'entity', 'content', ?2, 0)"
1699                        .to_string(),
1700                    params: vec![SqlValue::Text(update_id.clone()), SqlValue::Text(seed_ref)],
1701                    label: Some("blob_gc_fence_probe_update_arm_seed".to_string()),
1702                })
1703                .await?;
1704            let update_attempt = writer
1705                .execute(SqlStatement {
1706                    sql: "UPDATE attachments SET content_ref = ?1 \
1707                          WHERE record_uuid = ?2 AND role = 'content'"
1708                        .to_string(),
1709                    params: vec![
1710                        SqlValue::Text(BLOB_GC_FENCE_PROBE_REF.to_string()),
1711                        SqlValue::Text(update_id.clone()),
1712                    ],
1713                    label: Some("blob_gc_fence_probe_update_arm".to_string()),
1714                })
1715                .await;
1716            let update_fenced = fence_rejection(update_attempt);
1717
1718            let insert2_attempt = writer
1719                .execute(SqlStatement {
1720                    sql: "INSERT INTO attachments \
1721                          (record_uuid, substrate, role, content_ref, created_at) \
1722                          VALUES (?1, 'note', 'evidence', ?2, 0)"
1723                        .to_string(),
1724                    params: vec![
1725                        SqlValue::Text(insert2_id.clone()),
1726                        SqlValue::Text(probe2_ref.clone()),
1727                    ],
1728                    label: Some("blob_gc_fence_probe_insert2_arm".to_string()),
1729                })
1730                .await;
1731            let insert2_fenced = fence_rejection(insert2_attempt);
1732
1733            writer
1734                .execute(SqlStatement {
1735                    sql: "INSERT INTO attachments \
1736                          (record_uuid, substrate, role, content_ref, created_at) \
1737                          VALUES (?1, 'note', 'evidence', ?2, 0)"
1738                        .to_string(),
1739                    params: vec![
1740                        SqlValue::Text(update2_id.clone()),
1741                        SqlValue::Text(seed2_ref),
1742                    ],
1743                    label: Some("blob_gc_fence_probe_update2_arm_seed".to_string()),
1744                })
1745                .await?;
1746            let update2_attempt = writer
1747                .execute(SqlStatement {
1748                    sql: "UPDATE attachments SET content_ref = ?1 \
1749                          WHERE record_uuid = ?2 AND role = 'evidence'"
1750                        .to_string(),
1751                    params: vec![
1752                        SqlValue::Text(probe2_ref),
1753                        SqlValue::Text(update2_id.clone()),
1754                    ],
1755                    label: Some("blob_gc_fence_probe_update2_arm".to_string()),
1756                })
1757                .await;
1758            let update2_fenced = fence_rejection(update2_attempt);
1759
1760            // Remove every probe row before this unit commits, including an
1761            // attachment row a dead fence let through.
1762            writer
1763                .execute(SqlStatement {
1764                    sql: "DELETE FROM attachments WHERE record_uuid IN (?1, ?2, ?3, ?4)"
1765                        .to_string(),
1766                    params: vec![
1767                        SqlValue::Text(insert_id.clone()),
1768                        SqlValue::Text(update_id.clone()),
1769                        SqlValue::Text(insert2_id.clone()),
1770                        SqlValue::Text(update2_id.clone()),
1771                    ],
1772                    label: Some("blob_gc_fence_probe_cleanup_attachments".to_string()),
1773                })
1774                .await?;
1775            writer
1776                .execute(SqlStatement {
1777                    sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
1778                    params: vec![SqlValue::Text(claim_key)],
1779                    label: Some("blob_gc_fence_probe_cleanup_claim".to_string()),
1780                })
1781                .await?;
1782
1783            Ok(
1784                Box::new((insert_fenced, update_fenced, insert2_fenced, update2_fenced))
1785                    as Box<dyn std::any::Any + Send>,
1786            )
1787        })
1788    });
1789    let outcome = sql.atomic_unit(op).await?;
1790    let (insert_fenced, update_fenced, insert2_fenced, update2_fenced) = *outcome
1791        .downcast::<(
1792            Result<bool, String>,
1793            Result<bool, String>,
1794            Result<bool, String>,
1795            Result<bool, String>,
1796        )>()
1797        .map_err(|_| {
1798            StorageError::Internal("blob GC fence probe returned an unexpected outcome type".into())
1799        })?;
1800    let arm_verdict = |arm: &str, fenced: Result<bool, String>| -> StorageResult<()> {
1801        match fenced {
1802            Ok(true) => Ok(()),
1803            Ok(false) => Err(StorageError::Unsupported {
1804                capability: StorageCapability::Blob,
1805                operation: "transactional_orphan_sweep".into(),
1806                message: format!(
1807                    "the V21 fencing triggers exist by name but did not reject a claimed \
1808                     content_ref on the attachment {arm} path; refusing unfenced deletion"
1809                ),
1810            }),
1811            Err(other) => Err(StorageError::Unsupported {
1812                capability: StorageCapability::Blob,
1813                operation: "transactional_orphan_sweep".into(),
1814                message: format!(
1815                    "the blob GC fence probe could not verify the attachment {arm} fence \
1816                     (unexpected rejection: {other}); refusing unfenced deletion"
1817                ),
1818            }),
1819        }
1820    };
1821    arm_verdict("INSERT", insert_fenced)?;
1822    arm_verdict("UPDATE", update_fenced)?;
1823    arm_verdict("second-digest INSERT", insert2_fenced)?;
1824    arm_verdict("second-digest UPDATE", update2_fenced)
1825}
1826
1827async fn validate_blob_gc_evidence(sql: &dyn SqlAccess) -> StorageResult<()> {
1828    // These full-table integrity probes are statement-scoped reads. Keep them
1829    // off the single writer; only their one-row result is materialized. The
1830    // database sweep owner excludes another claim producer, and each bounded
1831    // claim unit anti-joins the then-current live rows under its writer lock.
1832    let mut reader = sql.reader().await?;
1833    // length() and GLOB both stop at an embedded NUL, so a value of 64 hex
1834    // characters followed by a NUL and arbitrary bytes passes them while
1835    // failing the exact-equality liveness anti-join. The byte-length arm
1836    // closes that class: chars = 64 AND bytes = 64 * the encoding's bytes
1837    // per ASCII character forces a NUL-free canonical value. CAST(TEXT AS
1838    // BLOB) yields the database text encoding's bytes (1 per hex char in
1839    // UTF-8, 2 in UTF-16), so the width is derived from the same database
1840    // rather than assumed, and an unrecognizable answer fails closed.
1841    let canonical_bytes = match reader
1842        .query_row(SqlStatement {
1843            sql: "SELECT length(CAST('x' AS BLOB))".to_string(),
1844            params: vec![],
1845            label: Some("blob_gc_validate_encoding_width".to_string()),
1846        })
1847        .await?
1848        .and_then(|row| row.columns.first().map(|column| column.value.clone()))
1849    {
1850        Some(SqlValue::Integer(width)) if (1..=4).contains(&width) => width * 64,
1851        other => {
1852            return Err(invalid_content_ref(format!(
1853                "the text-encoding width probe returned {other:?}; refusing GC validation"
1854            )));
1855        }
1856    };
1857    let invalid_claim = reader
1858        .query_row(SqlStatement {
1859            sql: "SELECT content_ref FROM blob_gc_claims \
1860                  WHERE typeof(content_ref) <> 'text' \
1861                     OR length(content_ref) <> 64 \
1862                     OR length(CAST(content_ref AS BLOB)) <> ?1 \
1863                     OR content_ref GLOB '*[^0-9a-f]*' \
1864                  LIMIT 1"
1865                .to_string(),
1866            params: vec![SqlValue::Integer(canonical_bytes)],
1867            label: Some("blob_gc_validate_existing_claims".to_string()),
1868        })
1869        .await?;
1870    if invalid_claim.is_some() {
1871        return Err(invalid_content_ref(
1872            "blob_gc_claims.content_ref contained a non-canonical value".into(),
1873        ));
1874    }
1875
1876    let invalid_live = reader
1877        .query_row(SqlStatement {
1878            sql: "SELECT content_ref FROM attachments \
1879                  WHERE typeof(content_ref) <> 'text' \
1880                      OR length(content_ref) <> 64 \
1881                      OR length(CAST(content_ref AS BLOB)) <> ?1 \
1882                      OR content_ref GLOB '*[^0-9a-f]*' \
1883                  LIMIT 1"
1884                .to_string(),
1885            params: vec![SqlValue::Integer(canonical_bytes)],
1886            label: Some("blob_gc_validate_live_refs".to_string()),
1887        })
1888        .await?;
1889    if invalid_live.is_some() {
1890        return Err(invalid_content_ref(
1891            "attachments.content_ref contained a non-canonical value".into(),
1892        ));
1893    }
1894    Ok(())
1895}
1896
1897async fn release_abandoned_blob_gc_claim_batch(sql: &dyn SqlAccess) -> StorageResult<u64> {
1898    let op: AtomicUnitOp = Box::new(move |writer| {
1899        Box::pin(async move {
1900            let released = writer
1901                .execute(SqlStatement {
1902                    sql: "DELETE FROM blob_gc_claims \
1903                          WHERE rowid IN ( \
1904                            SELECT rowid FROM blob_gc_claims \
1905                            ORDER BY rowid LIMIT ?1 \
1906                          )"
1907                    .to_string(),
1908                    params: vec![SqlValue::Integer(BLOB_GC_CLAIM_BATCH_SIZE as i64)],
1909                    label: Some("blob_gc_release_abandoned_claim_batch".to_string()),
1910                })
1911                .await?;
1912            Ok(Box::new(released) as Box<dyn std::any::Any + Send>)
1913        })
1914    });
1915    let released = sql.atomic_unit(op).await?;
1916    released.downcast::<u64>().map(|count| *count).map_err(|_| {
1917        StorageError::Internal(
1918            "transactional orphan sweep returned an unexpected recovery count type".into(),
1919        )
1920    })
1921}
1922
1923async fn claim_blob_gc_batch(
1924    sql: &dyn SqlAccess,
1925    root_key: String,
1926    candidates: &[(ContentRef, bool)],
1927    dry_run: bool,
1928) -> StorageResult<BlobGcBatchRows> {
1929    debug_assert!(candidates.len() <= BLOB_GC_CLAIM_BATCH_SIZE);
1930    let eligible_refs = candidates
1931        .iter()
1932        .filter(|(_, within_grace)| !within_grace)
1933        .map(|(content_ref, _)| content_ref.to_string())
1934        .collect::<Vec<_>>();
1935    let grace_refs = candidates
1936        .iter()
1937        .filter(|(_, within_grace)| *within_grace)
1938        .map(|(content_ref, _)| content_ref.to_string())
1939        .collect::<Vec<_>>();
1940    let eligible_json = serde_json::to_string(&eligible_refs).map_err(|error| {
1941        StorageError::Internal(format!(
1942            "failed to prepare blob GC eligible candidate batch: {error}"
1943        ))
1944    })?;
1945    let grace_json = serde_json::to_string(&grace_refs).map_err(|error| {
1946        StorageError::Internal(format!(
1947            "failed to prepare blob GC grace candidate batch: {error}"
1948        ))
1949    })?;
1950    let claimed_at = chrono::Utc::now().timestamp_micros();
1951    let op: AtomicUnitOp = Box::new(move |writer| {
1952        Box::pin(async move {
1953            let grace_period_skipped = required_nonnegative_count(
1954                writer
1955                    .query_scalar(SqlStatement {
1956                        sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
1957                              WHERE NOT EXISTS ( \
1958                                SELECT 1 FROM attachments \
1959                                WHERE content_ref = candidate.value \
1960                              )"
1961                        .to_string(),
1962                        params: vec![SqlValue::Text(grace_json)],
1963                        label: Some("blob_gc_count_grace_candidates_batch".to_string()),
1964                    })
1965                    .await?,
1966                "blob_gc_count_grace_candidates_batch",
1967            )?;
1968
1969            if dry_run {
1970                let would_delete = required_nonnegative_count(
1971                    writer
1972                        .query_scalar(SqlStatement {
1973                            sql: "SELECT COUNT(*) FROM json_each(?1) AS candidate \
1974                                  WHERE NOT EXISTS ( \
1975                                    SELECT 1 FROM attachments \
1976                                    WHERE content_ref = candidate.value \
1977                                  )"
1978                            .to_string(),
1979                            params: vec![SqlValue::Text(eligible_json)],
1980                            label: Some("blob_gc_count_dry_run_candidates_batch".to_string()),
1981                        })
1982                        .await?,
1983                    "blob_gc_count_dry_run_candidates_batch",
1984                )?;
1985                return Ok(Box::new(BlobGcBatchRows {
1986                    grace_period_skipped,
1987                    would_delete,
1988                    claimed_rows: Vec::new(),
1989                }) as Box<dyn std::any::Any + Send>);
1990            }
1991
1992            writer
1993                .execute(SqlStatement {
1994                    sql: "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
1995                          SELECT ?1, candidate.value, ?3 \
1996                          FROM json_each(?2) AS candidate \
1997                          WHERE NOT EXISTS ( \
1998                            SELECT 1 FROM attachments \
1999                            WHERE content_ref = candidate.value \
2000                          )"
2001                    .to_string(),
2002                    params: vec![
2003                        SqlValue::Text(root_key.clone()),
2004                        SqlValue::Text(eligible_json),
2005                        SqlValue::Integer(claimed_at),
2006                    ],
2007                    label: Some("blob_gc_claim_candidate_batch".to_string()),
2008                })
2009                .await?;
2010
2011            let claimed_rows = writer
2012                .query_all(SqlStatement {
2013                    sql: "SELECT content_ref FROM blob_gc_claims \
2014                          WHERE root_key = ?1 ORDER BY content_ref"
2015                        .to_string(),
2016                    params: vec![SqlValue::Text(root_key)],
2017                    label: Some("blob_gc_claimed_candidate_batch".to_string()),
2018                })
2019                .await?;
2020            Ok(Box::new(BlobGcBatchRows {
2021                grace_period_skipped,
2022                would_delete: claimed_rows.len() as u64,
2023                claimed_rows,
2024            }) as Box<dyn std::any::Any + Send>)
2025        })
2026    });
2027    let rows = sql.atomic_unit(op).await?;
2028    rows.downcast::<BlobGcBatchRows>()
2029        .map(|rows| *rows)
2030        .map_err(|_| {
2031            StorageError::Internal(
2032                "transactional orphan sweep returned an unexpected batch-row type".into(),
2033            )
2034        })
2035}
2036
2037fn parse_blob_gc_claim_rows(rows: Vec<SqlRow>) -> StorageResult<Vec<ContentRef>> {
2038    let mut claimed = Vec::with_capacity(rows.len());
2039    for row in rows {
2040        let raw = match row.get("content_ref") {
2041            Some(SqlValue::Text(raw)) => raw.clone(),
2042            _ => {
2043                return Err(invalid_content_ref(
2044                    "blob_gc_claims.content_ref contained a non-text value".into(),
2045                ));
2046            }
2047        };
2048        claimed.push(ContentRef::from_hex(raw).map_err(invalid_content_ref)?);
2049    }
2050    Ok(claimed)
2051}
2052
2053async fn release_blob_gc_batch(sql: &dyn SqlAccess, root_key: String) -> StorageResult<()> {
2054    let cleanup: AtomicUnitOp = Box::new(move |writer| {
2055        Box::pin(async move {
2056            writer
2057                .execute(SqlStatement {
2058                    sql: "DELETE FROM blob_gc_claims WHERE root_key = ?1".to_string(),
2059                    params: vec![SqlValue::Text(root_key)],
2060                    label: Some("blob_gc_release_claim_batch".to_string()),
2061                })
2062                .await?;
2063            Ok(Box::new(()) as Box<dyn std::any::Any + Send>)
2064        })
2065    });
2066    sql.atomic_unit(cleanup).await?;
2067    Ok(())
2068}
2069
2070/// Process-wide database owner fence for transactional blob sweeps.
2071///
2072/// Claims live in the database and their attachment triggers are database-global,
2073/// so a root-only lock is insufficient: two differently configured roots for
2074/// one database must not recover each other's live claims. File-backed pools
2075/// additionally take [`acquire_database_gc_lock`] for cross-process exclusion.
2076type SweepLockMap = HashMap<Option<PathBuf>, Arc<DatabaseGcProcessLock>>;
2077
2078#[derive(Debug, Default)]
2079struct DatabaseGcProcessLock {
2080    held: StdMutex<bool>,
2081    released: std::sync::Condvar,
2082    #[cfg(test)]
2083    waiters: std::sync::atomic::AtomicUsize,
2084}
2085
2086impl DatabaseGcProcessLock {
2087    fn acquire(self: &Arc<Self>) -> DatabaseGcProcessGuard {
2088        let mut held = self
2089            .held
2090            .lock()
2091            .unwrap_or_else(std::sync::PoisonError::into_inner);
2092        while *held {
2093            #[cfg(test)]
2094            self.waiters
2095                .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2096            held = self
2097                .released
2098                .wait(held)
2099                .unwrap_or_else(std::sync::PoisonError::into_inner);
2100            #[cfg(test)]
2101            self.waiters
2102                .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
2103        }
2104        *held = true;
2105        DatabaseGcProcessGuard {
2106            lock: Arc::clone(self),
2107        }
2108    }
2109
2110    fn try_acquire(self: &Arc<Self>) -> Option<DatabaseGcProcessGuard> {
2111        let mut held = self
2112            .held
2113            .lock()
2114            .unwrap_or_else(std::sync::PoisonError::into_inner);
2115        if *held {
2116            return None;
2117        }
2118        *held = true;
2119        Some(DatabaseGcProcessGuard {
2120            lock: Arc::clone(self),
2121        })
2122    }
2123}
2124
2125#[derive(Debug)]
2126struct DatabaseGcProcessGuard {
2127    lock: Arc<DatabaseGcProcessLock>,
2128}
2129
2130impl Drop for DatabaseGcProcessGuard {
2131    fn drop(&mut self) {
2132        let mut held = self
2133            .lock
2134            .held
2135            .lock()
2136            .unwrap_or_else(std::sync::PoisonError::into_inner);
2137        debug_assert!(*held, "database GC process owner released twice");
2138        *held = false;
2139        self.lock.released.notify_one();
2140    }
2141}
2142
2143fn database_sweep_locks() -> &'static StdMutex<SweepLockMap> {
2144    static REGISTRY: OnceLock<StdMutex<SweepLockMap>> = OnceLock::new();
2145    REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2146}
2147
2148fn sweep_lock_for_database(database_path: Option<&Path>) -> Arc<DatabaseGcProcessLock> {
2149    let key = database_path.map(Path::to_path_buf);
2150    let mut locks = database_sweep_locks()
2151        .lock()
2152        .unwrap_or_else(std::sync::PoisonError::into_inner);
2153    locks
2154        .entry(key)
2155        .or_insert_with(|| Arc::new(DatabaseGcProcessLock::default()))
2156        .clone()
2157}
2158
2159#[cfg(test)]
2160pub(crate) fn database_gc_waiter_count(database_path: Option<&Path>) -> usize {
2161    sweep_lock_for_database(database_path)
2162        .waiters
2163        .load(std::sync::atomic::Ordering::SeqCst)
2164}
2165
2166/// Exclusive canonical-database ownership shared by transactional blob GC and
2167/// the boot-gated V21 attachment cutover.
2168///
2169/// The process-local mutex is acquired first and retained while the advisory
2170/// file lock is acquired on a blocking thread. Moving the owned mutex guard
2171/// into that closure makes cancellation safe: dropping the outer future cannot
2172/// release process ownership while a blocking advisory acquisition continues.
2173pub struct DatabaseGcOwnerGuard {
2174    _process_guard: DatabaseGcProcessGuard,
2175    _advisory_guard: Option<fs::File>,
2176    database_path: Option<PathBuf>,
2177}
2178
2179pub(crate) fn acquire_database_gc_owner_for_path_blocking(
2180    database_path: Option<PathBuf>,
2181) -> StorageResult<DatabaseGcOwnerGuard> {
2182    let process_guard = sweep_lock_for_database(database_path.as_deref()).acquire();
2183    let advisory_guard = acquire_database_gc_lock(database_path.as_deref())?;
2184    Ok(DatabaseGcOwnerGuard {
2185        _process_guard: process_guard,
2186        _advisory_guard: advisory_guard,
2187        database_path,
2188    })
2189}
2190
2191/// Try to acquire canonical database-GC ownership without waiting.
2192///
2193/// This is the fail-closed bridge for the legacy raw-connection migration API:
2194/// a caller may already hold an opaque pooled writer guard, so waiting here
2195/// could invert the canonical owner-before-writer order used by sweeps. The
2196/// production backend boot path uses the blocking helper before writer
2197/// checkout instead.
2198pub(crate) fn try_acquire_database_gc_owner_for_path(
2199    database_path: PathBuf,
2200) -> StorageResult<DatabaseGcOwnerGuard> {
2201    let process_guard = sweep_lock_for_database(Some(&database_path))
2202        .try_acquire()
2203        .ok_or_else(|| {
2204            StorageError::Internal(format!(
2205                "database GC owner for {} is already held; retry schema migration through the \
2206                 coordinated backend boot path",
2207                database_path.display()
2208            ))
2209        })?;
2210    let lock_path = database_gc_lock_path(&database_path);
2211    let advisory_guard = fs::OpenOptions::new()
2212        .read(true)
2213        .write(true)
2214        .create(true)
2215        .truncate(false)
2216        .open(&lock_path)
2217        .map_err(|error| map_io_err(error, "database_gc_lock_open"))?;
2218    fs4::FileExt::try_lock(&advisory_guard)
2219        .map_err(|error| map_io_err(error.into(), "database_gc_lock_try_acquire"))?;
2220    Ok(DatabaseGcOwnerGuard {
2221        _process_guard: process_guard,
2222        _advisory_guard: Some(advisory_guard),
2223        database_path: Some(database_path),
2224    })
2225}
2226
2227impl std::fmt::Debug for DatabaseGcOwnerGuard {
2228    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2229        formatter
2230            .debug_struct("DatabaseGcOwnerGuard")
2231            .field("database_path", &self.database_path)
2232            .finish_non_exhaustive()
2233    }
2234}
2235
2236impl DatabaseGcOwnerGuard {
2237    /// Canonical database path that keys this owner, or `None` for an
2238    /// in-memory database whose process-local mutex is the complete fence.
2239    pub fn database_path(&self) -> Option<&Path> {
2240        self.database_path.as_deref()
2241    }
2242}
2243
2244/// Acquire the canonical database owner used by both V21 boot cutover and
2245/// transactional blob sweep. Callers must retain the returned guard across
2246/// every stage that must exclude the other protocol.
2247pub async fn acquire_database_gc_owner(sql: &dyn SqlAccess) -> StorageResult<DatabaseGcOwnerGuard> {
2248    let database_path = sql.database_path();
2249    tokio::task::spawn_blocking(move || acquire_database_gc_owner_for_path_blocking(database_path))
2250        .await
2251        .map_err(|error| {
2252            StorageError::driver(StorageCapability::Blob, "acquire_database_gc_owner", error)
2253        })?
2254}
2255
2256/// Process-wide registry of per-canonical-root write locks.
2257///
2258/// A `Mutex` field scoped to one `FsBlobStore` instance does NOT serialize
2259/// writes across independently constructed stores for the same root — and
2260/// callers construct fresh stores for the same root routinely
2261/// (`StorageBackend::blob_store` builds a new `FsBlobStore` on every call).
2262/// Keying a shared `Arc<tokio::sync::Mutex<()>>` by
2263/// the filesystem's own canonical path closes that gap: every `FsBlobStore`
2264/// for the same root, however many separate `new` calls produced them,
2265/// resolves to the exact same lock.
2266fn root_write_locks() -> &'static StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>> {
2267    static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<()>>>>> =
2268        OnceLock::new();
2269    REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2270}
2271
2272/// Look up (or create) the shared write lock for `root`'s canonical path.
2273///
2274/// `root` must already exist when this is called — `FsBlobStore::new`
2275/// creates it first, and `Path::canonicalize` requires the path to exist.
2276/// The lookup-or-insert happens under the registry's own (synchronous, very
2277/// briefly held) lock, so two `FsBlobStore::new` calls racing for the same
2278/// root cannot each install a different `Arc` and defeat the sharing this
2279/// exists for.
2280fn write_lock_for_root(root: &Path) -> std::io::Result<Arc<tokio::sync::Mutex<()>>> {
2281    let canonical = root.canonicalize()?;
2282    let mut locks = root_write_locks()
2283        .lock()
2284        .unwrap_or_else(std::sync::PoisonError::into_inner);
2285    Ok(locks
2286        .entry(canonical)
2287        .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
2288        .clone())
2289}
2290
2291/// A `BlobStore` backed by a BLAKE3-sharded directory tree.
2292#[derive(Debug)]
2293pub struct FsBlobStore {
2294    root: PathBuf,
2295    /// Initialization-time filesystem authority for `root`. Every blob
2296    /// descriptor walk starts from this retained handle rather than
2297    /// re-opening the root path and trusting its mutable ancestors again.
2298    root_handle: Arc<fs::File>,
2299    floor_bytes: u64,
2300    /// Shared per-canonical-root guard (see `write_lock_for_root`) that
2301    /// serializes the check-then-publish critical section of `put`: without
2302    /// this, two puts (whether on the same
2303    /// `FsBlobStore` instance or two independently constructed ones for the
2304    /// same root) can each observe the same pre-write `available_space`
2305    /// snapshot, each pass their own write-size-aware floor check against
2306    /// it, and then both write, jointly pushing the volume under the floor.
2307    /// `put` acquires this as an OWNED guard (`lock_owned`) and MOVES it
2308    /// into the `spawn_blocking` closure rather than borrowing it across the
2309    /// closure's `.await` — cancelling/dropping the outer `put` future then
2310    /// cannot release the guard before the underlying blocking write (which
2311    /// keeps running on its own thread regardless of the outer future's
2312    /// fate) actually finishes. A per-root async mutex is adequate at this
2313    /// write rate. The blocking write also takes a root-local advisory file
2314    /// lock to coordinate with publishers and transactional sweeps in other
2315    /// processes.
2316    write_lock: Arc<tokio::sync::Mutex<()>>,
2317    /// How long a blob with zero live references is left alone before an
2318    /// orphan sweep will delete it — see `within_publish_grace`. Bounds the
2319    /// window between `put` (bytes land, lock released) and the later,
2320    /// separate attachment write that commits a `content_ref` to it; it does not
2321    /// close that window entirely; see `within_publish_grace` and
2322    /// `transactional_orphan_sweep`'s doc comment for the residual exposure.
2323    orphan_sweep_grace: Duration,
2324}
2325
2326impl FsBlobStore {
2327    /// Default fail-closed free-space floor (khive#292 SPEC-gate ruling):
2328    /// 100 GB. Config-overridable via the `floor_bytes` constructor argument.
2329    pub const DEFAULT_FLOOR_BYTES: u64 = 100_000_000_000;
2330
2331    /// Default orphan-sweep publish grace period: 1 hour. Generous on
2332    /// purpose — it only needs to outlast the gap between a client's `put`
2333    /// call returning and its follow-up attachment write landing, not any
2334    /// steady-state condition.
2335    pub const DEFAULT_ORPHAN_SWEEP_GRACE: Duration = Duration::from_secs(3600);
2336
2337    /// Create a store rooted at `root`, creating the directory if absent.
2338    pub fn new(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
2339        fs::create_dir_all(&root)?;
2340        Self::open_existing(root, floor_bytes)
2341    }
2342
2343    /// Open a store rooted at an existing directory without creating any
2344    /// filesystem entry. Used by snapshot runtimes so boot can retain blob
2345    /// reads while remaining side-effect free; mutation is fenced by the
2346    /// runtime's read-only wrapper.
2347    pub fn open_existing(root: PathBuf, floor_bytes: u64) -> Result<Self, SqliteError> {
2348        // Preserve existing support for configured symlink spellings without
2349        // carrying that mutable indirection into handle-relative reads. The
2350        // bounded reader deliberately opens its stored root with NOFOLLOW.
2351        let root = root.canonicalize()?;
2352        let metadata = fs::metadata(&root)?;
2353        if !metadata.is_dir() {
2354            return Err(SqliteError::InvalidData(format!(
2355                "blob store root is not a directory: {}",
2356                root.display()
2357            )));
2358        }
2359        let root_handle = Arc::new(open_blob_root_handle(&root)?);
2360        let write_lock = write_lock_for_root(&root)?;
2361        verify_blob_root_identity(&root, &root_handle)?;
2362        Ok(Self {
2363            root,
2364            root_handle,
2365            floor_bytes,
2366            write_lock,
2367            orphan_sweep_grace: Self::DEFAULT_ORPHAN_SWEEP_GRACE,
2368        })
2369    }
2370
2371    /// Override the orphan-sweep publish grace period (default: one hour —
2372    /// see `DEFAULT_ORPHAN_SWEEP_GRACE`).
2373    pub fn with_orphan_sweep_grace(mut self, grace_period: Duration) -> Self {
2374        self.orphan_sweep_grace = grace_period;
2375        self
2376    }
2377
2378    /// The resolved root directory this store writes under.
2379    pub fn root(&self) -> &Path {
2380        &self.root
2381    }
2382}
2383
2384#[async_trait]
2385impl BlobStore for FsBlobStore {
2386    async fn put(&self, bytes: Vec<u8>) -> StorageResult<ContentRef> {
2387        // OWNED guard, MOVED into the blocking closure below: a guard merely
2388        // borrowed here and held in this
2389        // async fn's own stack frame would be released the instant the
2390        // *outer* `put` future is cancelled or dropped, while an
2391        // already-started `spawn_blocking` closure keeps running on its own
2392        // thread regardless — letting a second `put` pass its check against
2393        // an unprotected in-flight write. Moving the owned guard into the
2394        // closure ties its lifetime to the blocking work itself, not to
2395        // whether anyone is still awaiting this future.
2396        let owned_guard = self.write_lock.clone().lock_owned().await;
2397        let root = self.root.clone();
2398        let root_handle = Arc::clone(&self.root_handle);
2399        let floor_bytes = self.floor_bytes;
2400        // `sync_hook::take` (added for PR #922) is the
2401        // test-only seam that lets regression tests observe/control
2402        // exactly when this call is inside the guarded section, replacing
2403        // fixed-sleep/fixed-duration-poll timing assumptions with
2404        // deterministic, event-driven synchronization. `#[cfg(test)]`-
2405        // gated end to end -- zero effect on non-test builds.
2406        #[cfg(test)]
2407        let hook = sync_hook::take(&root);
2408        tokio::task::spawn_blocking(move || {
2409            // The guard lives in this inner block so it is dropped BEFORE
2410            // the test hook's `done` signal fires below -- a test that
2411            // waits on `done` and then immediately asserts the lock is
2412            // free needs that ordering to hold exactly, not "usually".
2413            #[cfg_attr(not(test), allow(clippy::let_and_return))]
2414            let result = {
2415                let _owned_guard = owned_guard;
2416                #[cfg(test)]
2417                if let Some(h) = &hook {
2418                    let _ = h.reached.send(());
2419                    let _ = h.release.recv();
2420                }
2421                put_blocking_from_root_handle(&root, &root_handle, floor_bytes, bytes)
2422            };
2423            #[cfg(test)]
2424            if let Some(h) = &hook {
2425                let _ = h.done.send(());
2426            }
2427            result
2428        })
2429        .await
2430        .map_err(|e| StorageError::driver(StorageCapability::Blob, "put", e))?
2431    }
2432
2433    async fn get_bounded_verified(
2434        &self,
2435        content_ref: &ContentRef,
2436        max_bytes: u64,
2437    ) -> StorageResult<Vec<u8>> {
2438        // Argument validation is deliberately outside `spawn_blocking`: an
2439        // invalid portable-envelope request must fail before any backend work
2440        // is scheduled or the filesystem is touched (ADR-160 D2).
2441        if max_bytes > MAX_BLOB_WHOLE_BYTES {
2442            return Err(StorageError::InvalidInput {
2443                capability: StorageCapability::Blob,
2444                operation: "get_bounded_verified".into(),
2445                message: format!(
2446                    "max_bytes {max_bytes} exceeds the {MAX_BLOB_WHOLE_BYTES}-byte portable whole-buffer envelope"
2447                ),
2448            });
2449        }
2450
2451        let root = self.root.clone();
2452        let root_handle = Arc::clone(&self.root_handle);
2453        let content_ref = content_ref.clone();
2454        #[cfg(test)]
2455        let read_hook = bounded_read_sync_hook::take(&root);
2456        tokio::task::spawn_blocking(move || {
2457            // Open exactly once. Both metadata and bytes below come from this
2458            // no-follow handle, so replacing the path after open cannot
2459            // switch the integrity authority underneath the read.
2460            let mut file = open_blob_shard_file_no_follow(&root, &root_handle, &content_ref)
2461                .map_err(|e| {
2462                    if e.kind() == std::io::ErrorKind::NotFound {
2463                        StorageError::NotFound {
2464                            capability: StorageCapability::Blob,
2465                            resource: "blob",
2466                            key: content_ref.to_string(),
2467                        }
2468                    } else if e.kind() == std::io::ErrorKind::Unsupported {
2469                        StorageError::Unsupported {
2470                            capability: StorageCapability::Blob,
2471                            operation: "get_bounded_verified".into(),
2472                            message: e.to_string(),
2473                        }
2474                    } else {
2475                        map_io_err(e, "get_bounded_verified.open")
2476                    }
2477                })?;
2478
2479            let metadata_bytes = file
2480                .metadata()
2481                .map_err(|e| map_io_err(e, "get_bounded_verified.metadata"))?
2482                .len();
2483            #[cfg(test)]
2484            if let Some(hook) = &read_hook {
2485                let _ = hook.reached.send(());
2486                let _ = hook.release.recv();
2487            }
2488            if metadata_bytes > max_bytes {
2489                return Err(StorageError::BlobTooLarge {
2490                    content_ref,
2491                    max_bytes,
2492                    observed_at_least: metadata_bytes,
2493                });
2494            }
2495
2496            // Read at most one sentinel byte beyond the caller's limit. The
2497            // +1 is safe because max_bytes was already bounded to 64 MiB.
2498            let mut bytes = Vec::with_capacity(metadata_bytes as usize);
2499            (&mut file)
2500                .take(max_bytes + 1)
2501                .read_to_end(&mut bytes)
2502                .map_err(|e| map_io_err(e, "get_bounded_verified.read"))?;
2503            let actual_bytes = bytes.len() as u64;
2504            if actual_bytes > max_bytes {
2505                return Err(StorageError::BlobTooLarge {
2506                    content_ref,
2507                    max_bytes,
2508                    observed_at_least: actual_bytes,
2509                });
2510            }
2511            if metadata_bytes != actual_bytes {
2512                return Err(StorageError::BlobSizeMismatch {
2513                    content_ref,
2514                    metadata_bytes,
2515                    actual_bytes,
2516                });
2517            }
2518
2519            let actual = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
2520            if actual != content_ref {
2521                return Err(StorageError::BlobDigestMismatch {
2522                    expected: content_ref,
2523                    actual,
2524                });
2525            }
2526            Ok(bytes)
2527        })
2528        .await
2529        .map_err(|e| StorageError::driver(StorageCapability::Blob, "get_bounded_verified", e))?
2530    }
2531
2532    async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool> {
2533        let root = self.root.clone();
2534        let root_handle = Arc::clone(&self.root_handle);
2535        let content_ref = content_ref.clone();
2536        tokio::task::spawn_blocking(move || {
2537            match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2538                Ok(_) => Ok(true),
2539                Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
2540                Err(error) => Err(map_io_err(error, "exists")),
2541            }
2542        })
2543        .await
2544        .map_err(|e| StorageError::driver(StorageCapability::Blob, "exists", e))?
2545    }
2546
2547    async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>> {
2548        let root = self.root.clone();
2549        let root_handle = Arc::clone(&self.root_handle);
2550        let content_ref = content_ref.clone();
2551        tokio::task::spawn_blocking(move || {
2552            match open_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2553                Ok(file) => file
2554                    .metadata()
2555                    .map(|metadata| Some(metadata.len()))
2556                    .map_err(|error| map_io_err(error, "size")),
2557                Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
2558                Err(error) => Err(map_io_err(error, "size")),
2559            }
2560        })
2561        .await
2562        .map_err(|e| StorageError::driver(StorageCapability::Blob, "size", e))?
2563    }
2564
2565    async fn delete(&self, content_ref: &ContentRef) -> StorageResult<bool> {
2566        let root = self.root.clone();
2567        let root_handle = Arc::clone(&self.root_handle);
2568        let content_ref = content_ref.clone();
2569        tokio::task::spawn_blocking(move || {
2570            match unlink_blob_shard_file_no_follow(&root, &root_handle, &content_ref) {
2571                Ok(()) => Ok(true),
2572                Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
2573                Err(e) => Err(map_io_err(e, "delete")),
2574            }
2575        })
2576        .await
2577        .map_err(|e| StorageError::driver(StorageCapability::Blob, "delete", e))?
2578    }
2579
2580    // Disabled for this compatibility release (Phase4a, ADR-111 §8 amended
2581    // 2026-08-21): this API has no `SqlAccess` capability of its own, so it
2582    // cannot prove the completed-V21 epoch `transactional_orphan_sweep`
2583    // requires before any destructive path runs. A caller-assembled
2584    // `live_refs` snapshot could otherwise delete an object a V20 SQL query
2585    // cannot see as live (e.g. a moodboard FANN network), bypassing the
2586    // epoch gate entirely. Report-only and destructive calls both refuse,
2587    // matching the trait default, until a snapshot API can carry its own
2588    // epoch proof.
2589    async fn orphan_sweep(
2590        &self,
2591        config: &BlobOrphanSweepConfig,
2592    ) -> StorageResult<BlobOrphanSweepResult> {
2593        let _ = config;
2594        Err(StorageError::Unsupported {
2595            capability: StorageCapability::Blob,
2596            operation: "orphan_sweep".into(),
2597            message: "caller-snapshot orphan_sweep is disabled in this compatibility release; \
2598                      it cannot prove a completed V21 attachment epoch, use \
2599                      transactional_orphan_sweep instead"
2600                .into(),
2601        })
2602    }
2603
2604    // `put` and the attachment write that later commits a `content_ref` to its
2605    // result are two separate steps of the client protocol -- the write
2606    // lock this method takes only serializes it against a concurrent `put`,
2607    // it is not held across the caller's own gap between finishing `put` and
2608    // issuing that follow-up attachment write. A blob can therefore be fully on
2609    // disk with zero live references purely because its referencing write
2610    // hasn't landed yet, not because it is actually orphaned.
2611    // `within_publish_grace` (via `orphan_sweep_grace`) is what protects that
2612    // window: a file younger than the grace period is left alone regardless
2613    // of liveness. Residual assumption: a client that waits longer than the
2614    // grace period between `put` returning and its attachment write committing
2615    // is still exposed to this method deleting the blob out from under it --
2616    // callers with an unusually slow publish path should widen the grace
2617    // period (`FsBlobStore::with_orphan_sweep_grace`) accordingly.
2618    //
2619    // Cross-resource ordering (#1850): database/root ownership and filesystem
2620    // walk/metadata happen before SQL. Bounded SQL-only units recover abandoned
2621    // rows and commit at most 128 fresh claims whose attachment triggers fence new
2622    // live references; physical deletion happens after each COMMIT; a second
2623    // bounded SQL-only unit releases that batch. Database/root owners span the
2624    // destructive phases, but SQLite's single writer never spans external I/O.
2625    async fn transactional_orphan_sweep(
2626        &self,
2627        sql: &dyn SqlAccess,
2628        dry_run: bool,
2629    ) -> StorageResult<BlobOrphanSweepResult> {
2630        // This compatibility gate deliberately precedes every database/root
2631        // owner wait, filesystem walk, and abandoned-claim cleanup. V20 and
2632        // staged V21 cannot represent the complete attachment liveness set,
2633        // so even report-only sweeps must refuse without observable mutation.
2634        if !blob_gc_fencing_complete(sql).await? {
2635            return Err(unsupported_blob_gc_epoch());
2636        }
2637
2638        // Claims and their attachment triggers are database-global. Serialize the
2639        // whole cross-resource protocol by database before taking the root
2640        // locks, so differently configured roots cannot recover one another's
2641        // active claim batches. The OS lock is the crash-detecting owner:
2642        // acquiring it proves that every row left in this database is
2643        // abandoned, including rows copied by backup or left before a root
2644        // relocation.
2645        let database_path = sql.database_path();
2646        let lock_database_path = database_path.clone();
2647        #[cfg(test)]
2648        let hook_database_path = database_path.clone();
2649        let (database_guard, database_file_guard) = tokio::task::spawn_blocking(move || {
2650            let process_guard = sweep_lock_for_database(lock_database_path.as_deref()).acquire();
2651            let file_guard = acquire_database_gc_lock(lock_database_path.as_deref())?;
2652            #[cfg(test)]
2653            if let Some(hook) = db_ownership_sync_hook::take(hook_database_path.as_deref()) {
2654                let _ = hook.reached.send(());
2655                let _ = hook.release.recv();
2656            }
2657            Ok::<_, StorageError>((process_guard, file_guard))
2658        })
2659        .await
2660        .map_err(|e| {
2661            StorageError::driver(
2662                StorageCapability::Blob,
2663                "transactional_orphan_sweep_lock",
2664                e,
2665            )
2666        })??;
2667
2668        // Recheck immediately once database ownership -- the process-local
2669        // guard plus the cross-process advisory lock -- is held, and before
2670        // ever waiting on the root guard/lock or walking the filesystem, so
2671        // the documented pre-lock refusal contract holds even if the epoch
2672        // regressed between the read-only preflight above and ownership.
2673        if !blob_gc_fencing_complete(sql).await? {
2674            return Err(unsupported_blob_gc_epoch());
2675        }
2676
2677        let root_guard = self.write_lock.clone().lock_owned().await;
2678        let root = self.root.clone();
2679        let root_handle = Arc::clone(&self.root_handle);
2680        let scan_root = root.clone();
2681        let scan_root_handle = Arc::clone(&root_handle);
2682        let grace_period = self.orphan_sweep_grace;
2683        let (write_guards, canonical_root, prepared) = tokio::task::spawn_blocking(move || {
2684            verify_blob_root_identity(&scan_root, &scan_root_handle)
2685                .map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
2686            // `self.root` was canonicalized at construction. Preserve that
2687            // spelling instead of re-resolving mutable ancestors here.
2688            let canonical_root = scan_root;
2689            let root_write_guard =
2690                acquire_root_write_lock_anchored(&canonical_root, &scan_root_handle)?;
2691            let candidates = walk_blob_files_from_root_handle(&scan_root_handle, &canonical_root)
2692                .map_err(|e| map_io_err(e, "transactional_orphan_sweep_walk"))?;
2693            verify_blob_root_identity(&canonical_root, &scan_root_handle)
2694                .map_err(|e| map_io_err(e, "transactional_orphan_sweep_root"))?;
2695            let prepared = prepare_transactional_sweep(candidates, grace_period);
2696            Ok::<_, StorageError>((
2697                (
2698                    database_guard,
2699                    database_file_guard,
2700                    root_guard,
2701                    root_write_guard,
2702                ),
2703                canonical_root,
2704                prepared,
2705            ))
2706        })
2707        .await
2708        .map_err(|e| {
2709            StorageError::driver(
2710                StorageCapability::Blob,
2711                "transactional_orphan_sweep_walk",
2712                e,
2713            )
2714        })??;
2715        let root_key = blob_root_key(&canonical_root);
2716        validate_blob_gc_evidence(sql).await?;
2717        blob_gc_fence_probe(sql).await?;
2718        if !dry_run {
2719            loop {
2720                let released = release_abandoned_blob_gc_claim_batch(sql).await?;
2721                if released < BLOB_GC_CLAIM_BATCH_SIZE as u64 {
2722                    break;
2723                }
2724            }
2725        }
2726
2727        let mut write_guards = write_guards;
2728        let mut result = prepared.result;
2729        let mut delete_error = None;
2730        #[cfg(test)]
2731        let mut hook: Option<sync_hook::Hook> = None;
2732        #[cfg(not(test))]
2733        let mut hook: Option<()> = None;
2734        #[cfg(test)]
2735        let mut hook_paused = false;
2736
2737        // Every unit below has a strict cardinality bound. The database owner
2738        // and root locks span the sequence, while each claim transaction and
2739        // cleanup transaction commits before filesystem work or the next
2740        // batch. SQLite can therefore checkpoint/reuse claim-table pages
2741        // between batches instead of receiving one orphan-population-sized
2742        // transaction.
2743        for candidates in prepared.candidates.chunks(BLOB_GC_CLAIM_BATCH_SIZE) {
2744            let batch = claim_blob_gc_batch(sql, root_key.clone(), candidates, dry_run).await?;
2745            result.grace_period_skipped += batch.grace_period_skipped;
2746            result.would_delete += batch.would_delete;
2747            if dry_run {
2748                continue;
2749            }
2750
2751            let claimed_refs = parse_blob_gc_claim_rows(batch.claimed_rows)?;
2752            if claimed_refs.is_empty() {
2753                continue;
2754            }
2755
2756            #[cfg(test)]
2757            if hook.is_none() {
2758                hook = sync_hook::take(&root);
2759            }
2760            #[cfg(test)]
2761            let pause_hook = hook.is_some() && !hook_paused;
2762            #[cfg(test)]
2763            if pause_hook {
2764                hook_paused = true;
2765            }
2766
2767            let delete_root = canonical_root.clone();
2768            let delete_root_handle = Arc::clone(&root_handle);
2769            let (returned_guards, deleted, batch_delete_error, returned_hook) =
2770                tokio::task::spawn_blocking(move || {
2771                    #[cfg(test)]
2772                    if pause_hook {
2773                        if let Some(hook) = &hook {
2774                            let _ = hook.reached.send(());
2775                            let _ = hook.release.recv();
2776                        }
2777                    }
2778
2779                    let mut deleted = 0_u64;
2780                    let mut first_error = None;
2781                    for content_ref in claimed_refs {
2782                        match unlink_blob_shard_file_no_follow(
2783                            &delete_root,
2784                            &delete_root_handle,
2785                            &content_ref,
2786                        ) {
2787                            Ok(()) => deleted += 1,
2788                            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
2789                            Err(error) => {
2790                                first_error =
2791                                    Some(map_io_err(error, "transactional_orphan_sweep_delete"));
2792                                break;
2793                            }
2794                        }
2795                    }
2796                    // Field order is load-bearing under cancellation: a
2797                    // discarded blocking-task result drops tuple fields from
2798                    // left to right, so both owner guards release before the
2799                    // test hook's `done` sender disconnects.
2800                    (write_guards, deleted, first_error, hook)
2801                })
2802                .await
2803                .map_err(|error| {
2804                    StorageError::driver(
2805                        StorageCapability::Blob,
2806                        "transactional_orphan_sweep_delete",
2807                        error,
2808                    )
2809                })?;
2810            write_guards = returned_guards;
2811            hook = returned_hook;
2812            result.deleted += deleted;
2813
2814            // Release only this bounded batch after its physical phase. On
2815            // cancellation or process death before this commit, the claims
2816            // remain fail-closed and the next exclusive database owner
2817            // reevaluates them rather than resuming deletion blindly.
2818            release_blob_gc_batch(sql, root_key.clone()).await?;
2819            if batch_delete_error.is_some() {
2820                delete_error = batch_delete_error;
2821                break;
2822            }
2823        }
2824
2825        drop(write_guards);
2826        #[cfg(test)]
2827        if let Some(hook) = hook {
2828            let _ = hook.done.send(());
2829        }
2830        #[cfg(not(test))]
2831        let _ = hook;
2832        if let Some(error) = delete_error {
2833            return Err(error);
2834        }
2835        Ok(result)
2836    }
2837}
2838
2839/// Test-only synchronization seam into blob write-lock-guarded critical
2840/// sections (added for PR #922 and reused by the transactional sweep).
2841///
2842/// The prior regression tests proved mutual exclusion and cancellation-
2843/// safety with a fixed sleep before racing/aborting and a fixed-duration
2844/// poll loop waiting for the lock to free -- timing-dependent, and the poll
2845/// loop actually failed once in a required-suite run (a flaky
2846/// gate, not a real regression). This seam replaces both edges of the race
2847/// with event-driven coordination: a one-shot hook, queued per canonical
2848/// root, signals `reached` the instant execution is inside the guarded
2849/// closure (the owned guard already moved in) and blocks there until the
2850/// test sends `release`; `done` fires only after the guard has actually
2851/// been dropped (see `put`'s inner-block scoping of `_owned_guard`).
2852/// `#[cfg(test)]`-gated end to end -- zero effect on non-test builds.
2853#[cfg(test)]
2854mod sync_hook {
2855    use std::collections::{HashMap, VecDeque};
2856    use std::path::{Path, PathBuf};
2857    use std::sync::mpsc::{Receiver, Sender};
2858    use std::sync::{Mutex as StdMutex, OnceLock};
2859
2860    pub(super) struct Hook {
2861        pub(super) reached: Sender<()>,
2862        pub(super) release: Receiver<()>,
2863        pub(super) done: Sender<()>,
2864    }
2865
2866    fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
2867        static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
2868        REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2869    }
2870
2871    /// Queue a one-shot hook for the next instrumented operation against
2872    /// `root`'s canonical path. Consumed exactly once, FIFO.
2873    pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>, Receiver<()>) {
2874        let canonical = root
2875            .canonicalize()
2876            .expect("root must exist before installing a sync_hook");
2877        let (reached_tx, reached_rx) = std::sync::mpsc::channel();
2878        let (release_tx, release_rx) = std::sync::mpsc::channel();
2879        let (done_tx, done_rx) = std::sync::mpsc::channel();
2880        registry()
2881            .lock()
2882            .unwrap_or_else(std::sync::PoisonError::into_inner)
2883            .entry(canonical)
2884            .or_default()
2885            .push_back(Hook {
2886                reached: reached_tx,
2887                release: release_rx,
2888                done: done_tx,
2889            });
2890        (reached_rx, release_tx, done_rx)
2891    }
2892
2893    /// Pop the next queued hook for `root`'s canonical path, if any (`None`
2894    /// for every ordinary, non-instrumented test -- `put` runs completely
2895    /// unaffected). `root` need not be pre-canonicalized by the caller --
2896    /// both `install` and `take` canonicalize, matching how
2897    /// `write_lock_for_root` keys the shared lock registry.
2898    pub(super) fn take(root: &Path) -> Option<Hook> {
2899        let canonical = root.canonicalize().ok()?;
2900        registry()
2901            .lock()
2902            .unwrap_or_else(std::sync::PoisonError::into_inner)
2903            .get_mut(&canonical)
2904            .and_then(VecDeque::pop_front)
2905    }
2906}
2907
2908/// Test-only pause fired the next time the orphan-sweep walk is about to
2909/// open a leaf candidate file in `walk_blob_files_from_root_handle`. Lets a
2910/// test replace that exact on-disk leaf entry with a symlink to an outside
2911/// decoy between candidate discovery and the handle-relative open/fstat that
2912/// classifies it, proving the classification stays anchored to the opened
2913/// handle rather than re-resolving a path.
2914///
2915/// Keyed by the walking store's canonical root path (`entry`/`take` mirror
2916/// `sync_hook` above), not a single process-wide slot: a process-wide slot
2917/// is stealable by ANY other sweep walk that reaches a valid leaf while
2918/// running concurrently in the same test binary (the crate's default
2919/// parallel test runner interleaves `#[tokio::test]` functions), which made
2920/// the swap regression test flaky under that interleaving. Each test in this
2921/// module sweeps its own tempdir root, so keying by canonical root gives
2922/// each test's hook install/take pair exclusive use of its own slot.
2923#[cfg(all(test, unix))]
2924mod walk_leaf_sync_hook {
2925    use std::collections::HashMap;
2926    use std::path::{Path, PathBuf};
2927    use std::sync::mpsc::{Receiver, Sender};
2928    use std::sync::{Mutex as StdMutex, OnceLock};
2929
2930    pub(super) struct Hook {
2931        pub(super) reached: Sender<()>,
2932        pub(super) release: Receiver<()>,
2933    }
2934
2935    fn registry() -> &'static StdMutex<HashMap<PathBuf, Hook>> {
2936        static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, Hook>>> = OnceLock::new();
2937        REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2938    }
2939
2940    pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
2941        let canonical = root
2942            .canonicalize()
2943            .expect("root must exist before installing a walk_leaf_sync_hook");
2944        let (reached_tx, reached_rx) = std::sync::mpsc::channel();
2945        let (release_tx, release_rx) = std::sync::mpsc::channel();
2946        registry()
2947            .lock()
2948            .unwrap_or_else(std::sync::PoisonError::into_inner)
2949            .insert(
2950                canonical,
2951                Hook {
2952                    reached: reached_tx,
2953                    release: release_rx,
2954                },
2955            );
2956        (reached_rx, release_tx)
2957    }
2958
2959    /// `root` need not be pre-canonicalized by the caller -- both `install`
2960    /// and `take` canonicalize.
2961    pub(super) fn take(root: &Path) -> Option<Hook> {
2962        let canonical = root.canonicalize().ok()?;
2963        registry()
2964            .lock()
2965            .unwrap_or_else(std::sync::PoisonError::into_inner)
2966            .remove(&canonical)
2967    }
2968}
2969
2970/// Test-only pause after a bounded read has opened its authoritative handle
2971/// and captured handle metadata, but before its first byte read. This makes
2972/// append/truncate/path-replacement races deterministic without timing sleeps
2973/// and is deliberately separate from the put/GC lock hook above.
2974#[cfg(test)]
2975mod bounded_read_sync_hook {
2976    use std::collections::{HashMap, VecDeque};
2977    use std::path::{Path, PathBuf};
2978    use std::sync::mpsc::{Receiver, Sender};
2979    use std::sync::{Mutex as StdMutex, OnceLock};
2980
2981    pub(super) struct Hook {
2982        pub(super) reached: Sender<()>,
2983        pub(super) release: Receiver<()>,
2984    }
2985
2986    fn registry() -> &'static StdMutex<HashMap<PathBuf, VecDeque<Hook>>> {
2987        static REGISTRY: OnceLock<StdMutex<HashMap<PathBuf, VecDeque<Hook>>>> = OnceLock::new();
2988        REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
2989    }
2990
2991    pub(super) fn install(root: &Path) -> (Receiver<()>, Sender<()>) {
2992        let canonical = root
2993            .canonicalize()
2994            .expect("root must exist before installing a bounded-read hook");
2995        let (reached_tx, reached_rx) = std::sync::mpsc::channel();
2996        let (release_tx, release_rx) = std::sync::mpsc::channel();
2997        registry()
2998            .lock()
2999            .unwrap_or_else(std::sync::PoisonError::into_inner)
3000            .entry(canonical)
3001            .or_default()
3002            .push_back(Hook {
3003                reached: reached_tx,
3004                release: release_rx,
3005            });
3006        (reached_rx, release_tx)
3007    }
3008
3009    pub(super) fn take(root: &Path) -> Option<Hook> {
3010        let canonical = root.canonicalize().ok()?;
3011        registry()
3012            .lock()
3013            .unwrap_or_else(std::sync::PoisonError::into_inner)
3014            .get_mut(&canonical)
3015            .and_then(VecDeque::pop_front)
3016    }
3017}
3018
3019/// Test-only pause on `transactional_orphan_sweep`'s database-ownership
3020/// blocking task, right after the cross-process advisory lock is acquired
3021/// and before the epoch recheck that immediately follows it. Lets a test
3022/// mutate the database strictly between the read-only preflight and the
3023/// recheck, making the recheck's ordering relative to the root guard/lock
3024/// deterministic instead of racing on scheduling.
3025#[cfg(test)]
3026mod db_ownership_sync_hook {
3027    use std::collections::{HashMap, VecDeque};
3028    use std::path::{Path, PathBuf};
3029    use std::sync::mpsc::{Receiver, Sender};
3030    use std::sync::{Mutex as StdMutex, OnceLock};
3031
3032    pub(super) struct Hook {
3033        pub(super) reached: Sender<()>,
3034        pub(super) release: Receiver<()>,
3035    }
3036
3037    fn registry() -> &'static StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>> {
3038        static REGISTRY: OnceLock<StdMutex<HashMap<Option<PathBuf>, VecDeque<Hook>>>> =
3039            OnceLock::new();
3040        REGISTRY.get_or_init(|| StdMutex::new(HashMap::new()))
3041    }
3042
3043    pub(super) fn install(database_path: Option<&Path>) -> (Receiver<()>, Sender<()>) {
3044        let key = database_path.map(Path::to_path_buf);
3045        let (reached_tx, reached_rx) = std::sync::mpsc::channel();
3046        let (release_tx, release_rx) = std::sync::mpsc::channel();
3047        registry()
3048            .lock()
3049            .unwrap_or_else(std::sync::PoisonError::into_inner)
3050            .entry(key)
3051            .or_default()
3052            .push_back(Hook {
3053                reached: reached_tx,
3054                release: release_rx,
3055            });
3056        (reached_rx, release_tx)
3057    }
3058
3059    pub(super) fn take(database_path: Option<&Path>) -> Option<Hook> {
3060        let key = database_path.map(Path::to_path_buf);
3061        registry()
3062            .lock()
3063            .unwrap_or_else(std::sync::PoisonError::into_inner)
3064            .get_mut(&key)
3065            .and_then(VecDeque::pop_front)
3066    }
3067}
3068
3069#[cfg(test)]
3070mod tests {
3071    use super::*;
3072
3073    fn store(floor_bytes: u64) -> (tempfile::TempDir, FsBlobStore) {
3074        let dir = tempfile::tempdir().unwrap();
3075        let root = dir.path().join("blobs");
3076        // Zero orphan-sweep grace period: these tests exercise immediate
3077        // orphan deletion, not the publish-grace window (covered by the
3078        // `orphan_sweep_grace` tests below).
3079        let store = FsBlobStore::new(root, floor_bytes)
3080            .unwrap()
3081            .with_orphan_sweep_grace(Duration::ZERO);
3082        (dir, store)
3083    }
3084
3085    /// Build the exact historical V20 prefix without invoking the V21
3086    /// zero-reference fast path in [`crate::run_migrations`].
3087    fn prepare_v20_gc_fixture(conn: &mut rusqlite::Connection) {
3088        conn.execute_batch(include_str!("../../sql/schema-migrations-table.sql"))
3089            .expect("create migration ledger");
3090        for migration in crate::MIGRATIONS
3091            .iter()
3092            .filter(|migration| migration.version <= 20)
3093        {
3094            let tx = conn.transaction().expect("begin historical migration");
3095            tx.execute_batch(migration.up)
3096                .expect("apply historical migration body");
3097            tx.execute(
3098                "INSERT INTO _schema_migrations (version, name, applied_at) \
3099                 VALUES (?1, ?2, 0)",
3100                rusqlite::params![migration.version, migration.name],
3101            )
3102            .expect("record historical migration");
3103            tx.commit().expect("commit historical migration");
3104        }
3105    }
3106
3107    /// Build the canonical completed V21 schema used by transactional-GC
3108    /// tests. Phase 4b owns the real cutover now, so tests exercise its schema
3109    /// instead of retaining Phase 4a's synthetic future-schema fixture.
3110    fn prepare_completed_v21_gc_fixture(conn: &mut rusqlite::Connection) {
3111        let version = crate::run_migrations(conn).expect("prepare canonical completed V21");
3112        // Pinned to the V21 epoch, NOT to the moving latest version: the GC
3113        // gate admits exactly V21, so the day a V22 migration lands this
3114        // assert fails LOUDLY here — telling that migration's author to give
3115        // these tests a fixture built explicitly through V21 — instead of
3116        // silently handing every GC test a ledger the gate rejects.
3117        assert_eq!(
3118            version,
3119            crate::migrations::ATTACHMENT_CUTOVER_VERSION,
3120            "GC-gate fixtures need a completed-V21 ledger; a later migration chain \
3121             must provide a pinned through-V21 fixture builder for these tests"
3122        );
3123    }
3124
3125    #[tokio::test]
3126    async fn completed_v21_gc_gate_requires_new_indexes_and_absent_legacy_column() {
3127        let dir = tempfile::tempdir().unwrap();
3128        let db_path = dir.path().join("khive.db");
3129        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3130        {
3131            let mut writer = backend.pool().writer().unwrap();
3132            prepare_completed_v21_gc_fixture(writer.conn_mut());
3133        }
3134        assert!(blob_gc_fencing_complete(backend.sql().as_ref())
3135            .await
3136            .unwrap());
3137
3138        {
3139            let writer = backend.pool().writer().unwrap();
3140            writer
3141                .conn()
3142                .execute_batch("DROP INDEX idx_attachments_content_ref")
3143                .unwrap();
3144        }
3145        assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
3146            .await
3147            .unwrap());
3148
3149        {
3150            let writer = backend.pool().writer().unwrap();
3151            writer
3152                .conn()
3153                .execute_batch(
3154                    "CREATE INDEX idx_attachments_content_ref \
3155                         ON attachments(content_ref); \
3156                     ALTER TABLE entities ADD COLUMN content_ref TEXT",
3157                )
3158                .unwrap();
3159        }
3160        assert!(!blob_gc_fencing_complete(backend.sql().as_ref())
3161            .await
3162            .unwrap());
3163    }
3164
3165    /// Table-driven acceptance matrix: starting from a fixture that passes
3166    /// the gate, remove exactly one required table/index/trigger/marker/
3167    /// ledger fact at a time and prove the gate rejects it, no probe or
3168    /// abandoned-claim residue survives, and every candidate file and
3169    /// pre-existing claim is untouched. A wrong implementation that stops
3170    /// requiring any single one of these facts (e.g. the claims content_ref
3171    /// index, the V21 ledger row/name/max, the marker row/`completed_at`, or
3172    /// the INSERT fence alone) would still pass the narrower pre-existing
3173    /// tests but fails here.
3174    #[tokio::test]
3175    async fn completed_v21_gate_acceptance_matrix_rejects_each_removed_fact_independently() {
3176        let cases: &[(&str, &str)] = &[
3177            (
3178                "blob_gc_claims_content_ref_index_dropped",
3179                "DROP INDEX idx_blob_gc_claims_content_ref",
3180            ),
3181            (
3182                "v21_ledger_row_deleted",
3183                "DELETE FROM _schema_migrations WHERE version = 21",
3184            ),
3185            (
3186                "v21_ledger_row_renamed",
3187                "UPDATE _schema_migrations SET name = 'not_attachments_first_class' \
3188                 WHERE version = 21",
3189            ),
3190            // NOTE: the ledger predicate requires the exact contiguous
3191            // canonical ledger {1..21} — the named V21 row, COUNT(*) = 21,
3192            // MIN(version) = 1, MAX(version) = 21 (version is the PRIMARY
3193            // KEY, so together these pin the set exactly). Anything else —
3194            // rows missing below V21, or any row above it — is a schema
3195            // epoch this gate never validated and fails closed. The
3196            // ahead-of-V21 arm ADDS a fact rather than removing one and so
3197            // lives outside this matrix, in
3198            // `gate_rejects_ledger_ahead_of_binary_latest`.
3199            (
3200                "below_v21_ledger_rows_deleted",
3201                "DELETE FROM _schema_migrations WHERE version < 21",
3202            ),
3203            (
3204                "marker_row_deleted",
3205                "DELETE FROM attachment_cutover_state WHERE singleton = 1",
3206            ),
3207            (
3208                // The schema's own compound CHECK constraint already forbids
3209                // `state = 'complete' AND completed_at IS NULL` via a plain
3210                // UPDATE, so this rebuilds the table without that constraint
3211                // to construct the row directly -- proving the gate's own
3212                // `completed_at IS NOT NULL` predicate rejects it too,
3213                // independent of the CHECK constraint's own enforcement.
3214                "marker_completed_at_null_while_state_complete",
3215                "DROP TABLE attachment_cutover_state; \
3216                 CREATE TABLE attachment_cutover_state ( \
3217                     singleton    INTEGER PRIMARY KEY CHECK (singleton = 1), \
3218                     state        TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
3219                     started_at   INTEGER NOT NULL, \
3220                     completed_at INTEGER \
3221                 ) STRICT; \
3222                 INSERT INTO attachment_cutover_state \
3223                     (singleton, state, started_at, completed_at) \
3224                 VALUES (1, 'complete', 21, NULL)",
3225            ),
3226            (
3227                "insert_fence_dropped_alone",
3228                "DROP TRIGGER attachments_reject_claimed_blob_insert",
3229            ),
3230        ];
3231
3232        for (case, mutation_sql) in cases {
3233            let dir = tempfile::tempdir().unwrap();
3234            let db_path = dir.path().join("khive.db");
3235            let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3236            {
3237                let mut writer = backend.pool().writer().unwrap();
3238                prepare_completed_v21_gc_fixture(writer.conn_mut());
3239                writer
3240                    .conn_mut()
3241                    .execute_batch(mutation_sql)
3242                    .unwrap_or_else(|e| panic!("case {case}: failed to apply mutation: {e}"));
3243            }
3244
3245            assert!(
3246                !blob_gc_fencing_complete(backend.sql().as_ref())
3247                    .await
3248                    .unwrap(),
3249                "case {case}: gate must reject with this fact removed"
3250            );
3251
3252            let store = Arc::new(
3253                FsBlobStore::new(dir.path().join("blobs"), 0)
3254                    .unwrap()
3255                    .with_orphan_sweep_grace(Duration::ZERO),
3256            );
3257            let orphan = store
3258                .put(format!("gate matrix orphan for {case}").into_bytes())
3259                .await
3260                .unwrap();
3261            let abandoned_ref = "c".repeat(64);
3262            {
3263                let writer = backend.pool().writer().unwrap();
3264                writer
3265                    .conn()
3266                    .execute(
3267                        "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
3268                         VALUES ('gate-matrix-abandoned', ?1, 1)",
3269                        [abandoned_ref.as_str()],
3270                    )
3271                    .unwrap();
3272            }
3273
3274            let _root_guard = store.write_lock.clone().lock_owned().await;
3275            for dry_run in [true, false] {
3276                let outcome = tokio::time::timeout(
3277                    Duration::from_secs(1),
3278                    store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
3279                )
3280                .await
3281                .unwrap_or_else(|_| panic!("case {case}: refusal must precede the root wait"));
3282                assert!(
3283                    matches!(outcome, Err(StorageError::Unsupported { .. })),
3284                    "case {case} dry_run={dry_run}: expected Unsupported, got {outcome:?}"
3285                );
3286            }
3287
3288            assert!(
3289                store.exists(&orphan).await.unwrap(),
3290                "case {case}: a refused sweep must not delete anything"
3291            );
3292            let reader = backend.pool().reader().unwrap();
3293            let remaining: i64 = reader
3294                .conn()
3295                .query_row(
3296                    "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = 'gate-matrix-abandoned'",
3297                    [],
3298                    |row| row.get(0),
3299                )
3300                .unwrap();
3301            assert_eq!(remaining, 1, "case {case}: refusal must not recover claims");
3302
3303            // The gate refuses before `blob_gc_fence_probe` ever runs in
3304            // every one of these cases, so no probe-shaped row (the fence
3305            // probe's own claim/attachment id patterns) should exist in
3306            // either table for either dry_run mode.
3307            let probe_claims: i64 = reader
3308                .conn()
3309                .query_row(
3310                    "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key GLOB '__fence_probe-*'",
3311                    [],
3312                    |row| row.get(0),
3313                )
3314                .unwrap();
3315            assert_eq!(
3316                probe_claims, 0,
3317                "case {case}: refusal must leave no fence-probe claim residue"
3318            );
3319            let probe_attachments: i64 = reader
3320                .conn()
3321                .query_row(
3322                    "SELECT COUNT(*) FROM attachments \
3323                     WHERE record_uuid GLOB '__blob-gc-fence-probe-*'",
3324                    [],
3325                    |row| row.get(0),
3326                )
3327                .unwrap();
3328            assert_eq!(
3329                probe_attachments, 0,
3330                "case {case}: refusal must leave no fence-probe attachment residue"
3331            );
3332        }
3333    }
3334
3335    /// A ledger whose MAX(version) is above the V21 epoch belongs to a
3336    /// schema epoch this gate never validated; ADR-160 admits destructive GC
3337    /// only for the EXACT completed V21 epoch, so anything ahead fails
3338    /// closed. This is the arm the removal matrix above cannot host (it ADDS
3339    /// a ledger fact), and it separates the exact-epoch predicate from the
3340    /// fail-open `>= 21`: under `>= 21` this fixture passes the gate. The
3341    /// appended row is shaped exactly like a row `run_migrations` records —
3342    /// canonical name style, real timestamp — because "a migration this same
3343    /// binary applied on top" is the case the exact-epoch rule exists for.
3344    /// While the binary's terminal version IS 21, this fixture cannot
3345    /// distinguish `= 21` from the former `BETWEEN 21 AND terminal` — both
3346    /// refuse a V22 row — so the exact-epoch property is enforced by the
3347    /// predicate's text and by the contiguity clause the removal matrix
3348    /// binds (`below_v21_ledger_rows_deleted`); the first appended migration
3349    /// makes the ahead distinction observable.
3350    #[tokio::test]
3351    async fn gate_rejects_ledger_ahead_of_binary_latest() {
3352        let dir = tempfile::tempdir().unwrap();
3353        let db_path = dir.path().join("khive.db");
3354        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3355        {
3356            let mut writer = backend.pool().writer().unwrap();
3357            prepare_completed_v21_gc_fixture(writer.conn_mut());
3358            // Positive control: the untouched fixture passes.
3359        }
3360        assert!(
3361            blob_gc_fencing_complete(backend.sql().as_ref())
3362                .await
3363                .unwrap(),
3364            "control: completed fixture at the binary's latest version must pass"
3365        );
3366
3367        {
3368            let writer = backend.pool().writer().unwrap();
3369            writer
3370                .conn()
3371                .execute(
3372                    "INSERT INTO _schema_migrations (version, name, applied_at) \
3373                     VALUES (?1, 'post_cutover_feature', unixepoch())",
3374                    [i64::from(crate::migrations::latest_schema_version()) + 1],
3375                )
3376                .unwrap();
3377        }
3378        assert!(
3379            !blob_gc_fencing_complete(backend.sql().as_ref())
3380                .await
3381                .unwrap(),
3382            "a ledger ahead of the binary's latest schema version must fail closed"
3383        );
3384    }
3385
3386    /// An incomplete migration history must fail closed even when its
3387    /// terminal row looks right: delete every ledger row below V21 while
3388    /// keeping the V21 row, so the named-row check AND `MAX(version) = 21`
3389    /// both still hold. A predicate reading only those two facts admits this
3390    /// ledger; only the contiguity clause (COUNT/MIN/MAX over the
3391    /// PRIMARY-KEY `version` column) rejects it, so removing that clause
3392    /// turns this test red.
3393    #[tokio::test]
3394    async fn gate_rejects_incomplete_ledger_behind_v21() {
3395        let dir = tempfile::tempdir().unwrap();
3396        let db_path = dir.path().join("khive.db");
3397        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3398        {
3399            let mut writer = backend.pool().writer().unwrap();
3400            prepare_completed_v21_gc_fixture(writer.conn_mut());
3401        }
3402        assert!(
3403            blob_gc_fencing_complete(backend.sql().as_ref())
3404                .await
3405                .unwrap(),
3406            "control: the untouched completed fixture must pass"
3407        );
3408
3409        {
3410            let writer = backend.pool().writer().unwrap();
3411            writer
3412                .conn()
3413                .execute_batch("DELETE FROM _schema_migrations WHERE version < 21")
3414                .unwrap();
3415            // The facts the old predicate read are still intact.
3416            let (v21_named, max_version): (i64, i64) = writer
3417                .conn()
3418                .query_row(
3419                    "SELECT (SELECT COUNT(*) FROM _schema_migrations \
3420                             WHERE version = 21 AND name = 'attachments_first_class'), \
3421                            (SELECT MAX(version) FROM _schema_migrations)",
3422                    [],
3423                    |row| Ok((row.get(0)?, row.get(1)?)),
3424                )
3425                .unwrap();
3426            assert_eq!(
3427                (v21_named, max_version),
3428                (1, 21),
3429                "fixture must keep the named V21 row and MAX(version) = 21 so the \
3430                 rejection can only come from the contiguity clause"
3431            );
3432        }
3433        assert!(
3434            !blob_gc_fencing_complete(backend.sql().as_ref())
3435                .await
3436                .unwrap(),
3437            "an incomplete ledger behind V21 must fail closed despite a valid terminal row"
3438        );
3439    }
3440
3441    /// A read failure while evaluating the completed-marker/ledger predicate
3442    /// (a malformed or partially-migrated `attachment_cutover_state`) must
3443    /// propagate as an error, not be silently treated as "not complete" via
3444    /// some default-to-false path that could theoretically be confused with
3445    /// a permissive read elsewhere. The gate's three `query_scalar` calls all
3446    /// use `?`, so this pins that direction rather than leaving it assumed.
3447    #[tokio::test]
3448    async fn completed_v21_gate_fails_closed_when_marker_read_errors() {
3449        let dir = tempfile::tempdir().unwrap();
3450        let db_path = dir.path().join("khive.db");
3451        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
3452        {
3453            let mut writer = backend.pool().writer().unwrap();
3454            prepare_completed_v21_gc_fixture(writer.conn_mut());
3455            // `ALTER TABLE ... DROP COLUMN` refuses outright here because the
3456            // table's own compound CHECK constraint still references
3457            // `completed_at`; rebuild the table without the column instead
3458            // to get a genuine "no such column" read failure.
3459            writer
3460                .conn_mut()
3461                .execute_batch(
3462                    "DROP TABLE attachment_cutover_state; \
3463                     CREATE TABLE attachment_cutover_state ( \
3464                         singleton  INTEGER PRIMARY KEY CHECK (singleton = 1), \
3465                         state      TEXT NOT NULL CHECK (state IN ('incomplete', 'complete')), \
3466                         started_at INTEGER NOT NULL \
3467                     ) STRICT; \
3468                     INSERT INTO attachment_cutover_state (singleton, state, started_at) \
3469                     VALUES (1, 'complete', 21)",
3470                )
3471                .unwrap();
3472        }
3473
3474        let gate_error = blob_gc_fencing_complete(backend.sql().as_ref())
3475            .await
3476            .expect_err("a marker read error must propagate, not silently resolve to false");
3477        assert!(
3478            !matches!(gate_error, StorageError::Unsupported { .. }),
3479            "a read error is a distinct failure from the typed epoch refusal: {gate_error:?}"
3480        );
3481
3482        let store = Arc::new(
3483            FsBlobStore::new(dir.path().join("blobs"), 0)
3484                .unwrap()
3485                .with_orphan_sweep_grace(Duration::ZERO),
3486        );
3487        let orphan = store
3488            .put(b"marker read error orphan".to_vec())
3489            .await
3490            .unwrap();
3491        let _root_guard = store.write_lock.clone().lock_owned().await;
3492        for dry_run in [true, false] {
3493            let outcome = tokio::time::timeout(
3494                Duration::from_secs(1),
3495                store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
3496            )
3497            .await
3498            .expect("a marker read error must fail before waiting on the root lock");
3499            assert!(
3500                outcome.is_err(),
3501                "dry_run={dry_run}: expected the sweep to fail closed, got {outcome:?}"
3502            );
3503        }
3504        assert!(store.exists(&orphan).await.unwrap());
3505    }
3506
3507    #[test]
3508    fn database_sweep_owner_is_keyed_by_database_not_blob_root() {
3509        let dir = tempfile::tempdir().unwrap();
3510        let database = dir.path().join("khive.db");
3511        let same_database_a = sweep_lock_for_database(Some(&database));
3512        let same_database_b = sweep_lock_for_database(Some(&database));
3513        let other_database = sweep_lock_for_database(Some(&dir.path().join("other.db")));
3514        let mut expected_lock_path = database.as_os_str().to_os_string();
3515        expected_lock_path.push(DATABASE_GC_LOCK_SUFFIX);
3516
3517        assert!(Arc::ptr_eq(&same_database_a, &same_database_b));
3518        assert!(!Arc::ptr_eq(&same_database_a, &other_database));
3519        assert_eq!(
3520            database_gc_lock_path(&database),
3521            PathBuf::from(expected_lock_path)
3522        );
3523    }
3524
3525    #[tokio::test]
3526    async fn database_gc_owner_holds_process_and_advisory_fences_until_drop() {
3527        let dir = tempfile::tempdir().unwrap();
3528        let database = dir.path().join("owner.db");
3529        let backend = crate::StorageBackend::sqlite(&database).unwrap();
3530        let owner = acquire_database_gc_owner(backend.sql().as_ref())
3531            .await
3532            .unwrap();
3533        let canonical_database = owner
3534            .database_path()
3535            .expect("file-backed owner path")
3536            .to_path_buf();
3537
3538        assert!(
3539            sweep_lock_for_database(Some(&canonical_database))
3540                .try_acquire()
3541                .is_none(),
3542            "boot and sweep must share one process-local database owner"
3543        );
3544        let external = fs::OpenOptions::new()
3545            .read(true)
3546            .write(true)
3547            .open(database_gc_lock_path(&canonical_database))
3548            .unwrap();
3549        assert!(
3550            matches!(
3551                fs4::FileExt::try_lock(&external),
3552                Err(fs4::TryLockError::WouldBlock)
3553            ),
3554            "the reusable owner must also retain the cross-process advisory fence"
3555        );
3556
3557        drop(owner);
3558        fs4::FileExt::try_lock(&external).expect("owner drop releases advisory fence");
3559    }
3560
3561    #[cfg(unix)]
3562    #[test]
3563    fn database_gc_lock_path_preserves_non_utf8_identity() {
3564        use std::os::unix::ffi::{OsStrExt, OsStringExt};
3565
3566        let database = PathBuf::from(std::ffi::OsString::from_vec(
3567            b"khive-non-utf8-\xff.db".to_vec(),
3568        ));
3569        let lock_path = database_gc_lock_path(&database);
3570        let mut expected = database.as_os_str().as_bytes().to_vec();
3571        expected.extend_from_slice(DATABASE_GC_LOCK_SUFFIX.as_bytes());
3572        assert_eq!(lock_path.as_os_str().as_bytes(), expected);
3573    }
3574
3575    /// A shard directory replaced by a symlink (an attacker with write access
3576    /// to the blob root, a misconfigured shared parent, or a race between
3577    /// this sweep's directory walk and its physical delete) must be refused,
3578    /// never followed, and an unrelated real shard must keep sweeping
3579    /// normally. This is the fix for the shard-directory symlink-replacement
3580    /// hazard: a plain `fs::remove_file(shard_path(root, content_ref))`
3581    /// would resolve straight through the symlink and unlink whatever file
3582    /// its target names.
3583    #[cfg(unix)]
3584    #[tokio::test]
3585    async fn unlink_blob_shard_refuses_symlinked_shard_dir_and_still_sweeps_real_shard() {
3586        let (dir, store) = store(0);
3587        let root = dir.path().join("blobs");
3588
3589        // Real, non-attacked blob: written through the normal `put` path and
3590        // must still sweep after the fix.
3591        let real = store.put(b"real blob content".to_vec()).await.unwrap();
3592
3593        // Attack setup: a `content_ref` whose first shard directory does not
3594        // exist yet is replaced by a symlink to a directory entirely outside
3595        // the blob root, with a file planted at the exact name a path-based
3596        // delete would target.
3597        let outside = tempfile::tempdir().unwrap();
3598        let victim = outside.path().join("victim.txt");
3599        fs::write(&victim, b"do not delete me").unwrap();
3600
3601        let real_prefix = &real.as_str()[0..2];
3602        let attack_prefix = if real_prefix == "aa" { "bb" } else { "aa" };
3603        let fake_ref = ContentRef::from_hex(format!("{attack_prefix}{}", "0".repeat(62))).unwrap();
3604        let fake_hex = fake_ref.as_str().to_string();
3605
3606        let shard1 = root.join(attack_prefix);
3607        std::os::unix::fs::symlink(outside.path(), &shard1).unwrap();
3608        // Plant the file a naive path-based delete would actually resolve
3609        // to through the symlink: `<outside>/<shard2>/<full-hex>` is never
3610        // reached because the fix refuses at the `shard1` open itself, but
3611        // planting it this deep proves the refusal isn't accidental — there
3612        // was a real target for the naive path to have deleted.
3613        fs::create_dir_all(outside.path().join(&fake_hex[2..4])).unwrap();
3614        fs::write(
3615            outside.path().join(&fake_hex[2..4]).join(&fake_hex),
3616            b"decoy",
3617        )
3618        .unwrap();
3619
3620        let root_handle = open_blob_root_handle(&root).unwrap();
3621        let error = unlink_blob_shard_file_no_follow(&root, &root_handle, &fake_ref).unwrap_err();
3622        // `O_DIRECTORY | O_NOFOLLOW` against a symlink refuses to follow it,
3623        // but the exact errno is platform-dependent: Linux reports `ELOOP`,
3624        // Darwin reports `ENOTDIR` (the symlink itself is not a directory
3625        // once `O_NOFOLLOW` stops it from being resolved). Either way it
3626        // must be an outright refusal, not a successful open/unlink.
3627        assert!(
3628            matches!(
3629                error.raw_os_error(),
3630                Some(libc::ELOOP) | Some(libc::ENOTDIR)
3631            ),
3632            "opening a symlinked shard directory must be refused, not followed; got: {error}"
3633        );
3634        assert!(
3635            victim.exists(),
3636            "the file outside the blob root must never be touched by a refused shard-dir open"
3637        );
3638
3639        // The unrelated real shard, never touched by the attack, must still
3640        // unlink normally after the fix. `orphan_sweep` is disabled in this
3641        // compatibility release (it cannot prove a completed V21 epoch), so
3642        // this exercises the same underlying primitive directly rather than
3643        // going through that disabled API.
3644        unlink_blob_shard_file_no_follow(&root, &root_handle, &real).unwrap();
3645        assert!(!store.exists(&real).await.unwrap());
3646    }
3647
3648    /// Block on `rx.recv()` on a dedicated thread so a `#[tokio::test]`
3649    /// (current-thread runtime) doesn't stall other spawned tasks while
3650    /// waiting on a `sync_hook` signal: the deterministic,
3651    /// event-driven replacement for fixed-sleep / fixed-duration-poll
3652    /// assertions.
3653    async fn recv_blocking(rx: std::sync::mpsc::Receiver<()>) -> bool {
3654        tokio::task::spawn_blocking(move || rx.recv().is_ok())
3655            .await
3656            .expect("recv_blocking thread panicked")
3657    }
3658
3659    #[tokio::test]
3660    async fn put_bounded_get_roundtrip() {
3661        let (_dir, store) = store(0);
3662        let bytes = b"hello blob store".to_vec();
3663        let content_ref = store.put(bytes.clone()).await.unwrap();
3664        let fetched = store
3665            .get_bounded_verified(&content_ref, bytes.len() as u64)
3666            .await
3667            .unwrap();
3668        assert_eq!(fetched, bytes);
3669    }
3670
3671    #[tokio::test]
3672    async fn bounded_verified_get_accepts_exact_and_portable_maximum_limits() {
3673        let (_dir, store) = store(0);
3674        let bytes = b"bounded fs blob".to_vec();
3675        let content_ref = store.put(bytes.clone()).await.unwrap();
3676
3677        assert_eq!(
3678            store
3679                .get_bounded_verified(&content_ref, bytes.len() as u64)
3680                .await
3681                .unwrap(),
3682            bytes
3683        );
3684        assert_eq!(
3685            store
3686                .get_bounded_verified(&content_ref, MAX_BLOB_WHOLE_BYTES)
3687                .await
3688                .unwrap(),
3689            bytes
3690        );
3691    }
3692
3693    #[cfg(unix)]
3694    #[tokio::test]
3695    async fn bounded_verified_get_resolves_a_configured_symlink_root_once() {
3696        use std::os::unix::fs::symlink;
3697
3698        let dir = tempfile::tempdir().unwrap();
3699        let target = dir.path().join("blob-target");
3700        fs::create_dir(&target).unwrap();
3701        let configured = dir.path().join("blob-configured");
3702        symlink(&target, &configured).unwrap();
3703
3704        let store = FsBlobStore::new(configured, 0).unwrap();
3705        assert_eq!(store.root(), target.canonicalize().unwrap());
3706        let bytes = b"symlink-configured root".to_vec();
3707        let content_ref = store.put(bytes.clone()).await.unwrap();
3708        assert_eq!(
3709            store
3710                .get_bounded_verified(&content_ref, bytes.len() as u64)
3711                .await
3712                .unwrap(),
3713            bytes
3714        );
3715    }
3716
3717    #[cfg(unix)]
3718    #[tokio::test]
3719    async fn fs_blob_store_refuses_root_replacement_before_put() {
3720        use std::os::unix::fs::symlink;
3721
3722        let dir = tempfile::tempdir().unwrap();
3723        let root = dir.path().join("blobs");
3724        let store = FsBlobStore::new(root.clone(), 0).unwrap();
3725        let original_root = dir.path().join("blobs-original");
3726        fs::rename(&root, &original_root).unwrap();
3727
3728        let redirected_root = tempfile::tempdir().unwrap();
3729        symlink(redirected_root.path(), &root).unwrap();
3730
3731        let bytes = b"must not land in a replaced root".to_vec();
3732        let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
3733        store
3734            .put(bytes)
3735            .await
3736            .expect_err("a root replaced after construction must be refused");
3737
3738        assert!(
3739            !shard_path(redirected_root.path(), &content_ref).exists(),
3740            "put must not publish into the tree selected by the replacement symlink"
3741        );
3742        assert!(
3743            !shard_path(&original_root, &content_ref).exists(),
3744            "a refused put must not mutate the initialization-time root either"
3745        );
3746    }
3747
3748    #[cfg(unix)]
3749    #[tokio::test]
3750    async fn fs_blob_store_refuses_ancestor_replacement_for_existing_blob_operations() {
3751        use std::os::unix::fs::symlink;
3752
3753        let dir = tempfile::tempdir().unwrap();
3754        let ancestor = dir.path().join("store-parent");
3755        let root = ancestor.join("blobs");
3756        let store = FsBlobStore::new(root.clone(), 0).unwrap();
3757        let bytes = b"same bytes in both trees".to_vec();
3758        let content_ref = store.put(bytes.clone()).await.unwrap();
3759
3760        let original_ancestor = dir.path().join("store-parent-original");
3761        fs::rename(&ancestor, &original_ancestor).unwrap();
3762        let original_blob = shard_path(&original_ancestor.join("blobs"), &content_ref);
3763
3764        let redirected_ancestor = tempfile::tempdir().unwrap();
3765        let redirected_blob = shard_path(&redirected_ancestor.path().join("blobs"), &content_ref);
3766        fs::create_dir_all(redirected_blob.parent().unwrap()).unwrap();
3767        fs::write(&redirected_blob, &bytes).unwrap();
3768        symlink(redirected_ancestor.path(), &ancestor).unwrap();
3769
3770        store
3771            .get_bounded_verified(&content_ref, bytes.len() as u64)
3772            .await
3773            .expect_err("a read through a replaced root ancestor must be refused");
3774        store
3775            .exists(&content_ref)
3776            .await
3777            .expect_err("exists through a replaced root ancestor must be refused");
3778        store
3779            .size(&content_ref)
3780            .await
3781            .expect_err("size through a replaced root ancestor must be refused");
3782        store
3783            .delete(&content_ref)
3784            .await
3785            .expect_err("delete through a replaced root ancestor must be refused");
3786
3787        assert!(
3788            original_blob.exists(),
3789            "refusal must preserve the initialization-time blob"
3790        );
3791        assert!(
3792            redirected_blob.exists(),
3793            "refusal must not read as authority or delete the redirected blob"
3794        );
3795    }
3796
3797    #[tokio::test]
3798    async fn bounded_verified_get_rejects_same_size_digest_corruption() {
3799        let (_dir, store) = store(0);
3800        let expected_bytes = b"expected".to_vec();
3801        let actual_bytes = b"mutated!".to_vec();
3802        let expected = store.put(expected_bytes).await.unwrap();
3803        let actual = ContentRef::from_digest_bytes(blake3::hash(&actual_bytes).as_bytes());
3804        fs::write(shard_path(store.root(), &expected), actual_bytes).unwrap();
3805
3806        let err = store.get_bounded_verified(&expected, 8).await.unwrap_err();
3807        assert!(matches!(
3808            err,
3809            StorageError::BlobDigestMismatch {
3810                expected: ref got_expected,
3811                actual: ref got_actual,
3812            } if got_expected == &expected && got_actual == &actual
3813        ));
3814    }
3815
3816    #[tokio::test]
3817    async fn bounded_verified_get_stops_at_max_plus_one_after_file_growth() {
3818        let (_dir, store) = store(0);
3819        let store = Arc::new(store);
3820        let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
3821        let path = shard_path(store.root(), &content_ref);
3822        let (reached, release) = bounded_read_sync_hook::install(store.root());
3823
3824        let read_store = Arc::clone(&store);
3825        let read_ref = content_ref.clone();
3826        let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
3827        assert!(
3828            recv_blocking(reached).await,
3829            "read must reach the metadata seam"
3830        );
3831        let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
3832        writer.write_all(b"efgh-poison-tail").unwrap();
3833        writer.flush().unwrap();
3834        release.send(()).unwrap();
3835
3836        let err = read.await.unwrap().unwrap_err();
3837        assert!(matches!(
3838            err,
3839            StorageError::BlobTooLarge {
3840                content_ref: ref got,
3841                max_bytes: 4,
3842                observed_at_least: 5,
3843            } if got == &content_ref
3844        ));
3845    }
3846
3847    #[tokio::test]
3848    async fn bounded_verified_get_reports_growth_within_limit_as_size_mismatch() {
3849        let (_dir, store) = store(0);
3850        let store = Arc::new(store);
3851        let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
3852        let path = shard_path(store.root(), &content_ref);
3853        let (reached, release) = bounded_read_sync_hook::install(store.root());
3854
3855        let read_store = Arc::clone(&store);
3856        let read_ref = content_ref.clone();
3857        let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 8).await });
3858        assert!(
3859            recv_blocking(reached).await,
3860            "read must reach the metadata seam"
3861        );
3862        let mut writer = fs::OpenOptions::new().append(true).open(&path).unwrap();
3863        writer.write_all(b"ef").unwrap();
3864        writer.flush().unwrap();
3865        release.send(()).unwrap();
3866
3867        let err = read.await.unwrap().unwrap_err();
3868        assert!(matches!(
3869            err,
3870            StorageError::BlobSizeMismatch {
3871                content_ref: ref got,
3872                metadata_bytes: 4,
3873                actual_bytes: 6,
3874            } if got == &content_ref
3875        ));
3876    }
3877
3878    #[tokio::test]
3879    async fn bounded_verified_get_reports_truncation_as_size_mismatch() {
3880        let (_dir, store) = store(0);
3881        let store = Arc::new(store);
3882        let content_ref = store.put(b"abcd".to_vec()).await.unwrap();
3883        let path = shard_path(store.root(), &content_ref);
3884        let (reached, release) = bounded_read_sync_hook::install(store.root());
3885
3886        let read_store = Arc::clone(&store);
3887        let read_ref = content_ref.clone();
3888        let read = tokio::spawn(async move { read_store.get_bounded_verified(&read_ref, 4).await });
3889        assert!(
3890            recv_blocking(reached).await,
3891            "read must reach the metadata seam"
3892        );
3893        let mut writer = fs::OpenOptions::new()
3894            .write(true)
3895            .truncate(true)
3896            .open(&path)
3897            .unwrap();
3898        writer.write_all(b"abc").unwrap();
3899        writer.flush().unwrap();
3900        release.send(()).unwrap();
3901
3902        let err = read.await.unwrap().unwrap_err();
3903        assert!(matches!(
3904            err,
3905            StorageError::BlobSizeMismatch {
3906                content_ref: ref got,
3907                metadata_bytes: 4,
3908                actual_bytes: 3,
3909            } if got == &content_ref
3910        ));
3911    }
3912
3913    #[cfg(unix)]
3914    #[tokio::test]
3915    async fn bounded_verified_get_keeps_the_opened_inode_when_the_path_is_replaced() {
3916        let (_dir, store) = store(0);
3917        let store = Arc::new(store);
3918        let original = b"original".to_vec();
3919        let replacement = b"replaced".to_vec();
3920        let content_ref = store.put(original.clone()).await.unwrap();
3921        let path = shard_path(store.root(), &content_ref);
3922        let moved_path = path.with_extension("opened-inode");
3923        let max_bytes = original.len() as u64;
3924        let (reached, release) = bounded_read_sync_hook::install(store.root());
3925
3926        let read_store = Arc::clone(&store);
3927        let read_ref = content_ref.clone();
3928        let read =
3929            tokio::spawn(
3930                async move { read_store.get_bounded_verified(&read_ref, max_bytes).await },
3931            );
3932        assert!(
3933            recv_blocking(reached).await,
3934            "read must reach the metadata seam"
3935        );
3936        fs::rename(&path, &moved_path).unwrap();
3937        fs::write(&path, replacement).unwrap();
3938        release.send(()).unwrap();
3939
3940        assert_eq!(read.await.unwrap().unwrap(), original);
3941    }
3942
3943    #[cfg(unix)]
3944    #[tokio::test]
3945    async fn bounded_verified_get_refuses_a_symlink_leaf() {
3946        use std::os::unix::fs::symlink;
3947
3948        let dir = tempfile::tempdir().unwrap();
3949        let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
3950        let outside_bytes = b"outside but digest matching".to_vec();
3951        let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
3952        let outside = dir.path().join("outside");
3953        fs::write(&outside, &outside_bytes).unwrap();
3954        let leaf = shard_path(store.root(), &content_ref);
3955        fs::create_dir_all(leaf.parent().unwrap()).unwrap();
3956        symlink(&outside, &leaf).unwrap();
3957
3958        let err = store
3959            .get_bounded_verified(&content_ref, outside_bytes.len() as u64)
3960            .await
3961            .unwrap_err();
3962        assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
3963    }
3964
3965    #[cfg(unix)]
3966    #[tokio::test]
3967    async fn bounded_verified_get_refuses_a_symlinked_shard_component() {
3968        use std::os::unix::fs::symlink;
3969
3970        let dir = tempfile::tempdir().unwrap();
3971        let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
3972        let outside_bytes = b"outside through shard link".to_vec();
3973        let content_ref = ContentRef::from_digest_bytes(blake3::hash(&outside_bytes).as_bytes());
3974        let hex = content_ref.as_str();
3975        let outside_shard1 = dir.path().join("outside-shard1");
3976        let outside_shard2 = outside_shard1.join(&hex[2..4]);
3977        fs::create_dir_all(&outside_shard2).unwrap();
3978        fs::write(outside_shard2.join(hex), &outside_bytes).unwrap();
3979        symlink(&outside_shard1, store.root().join(&hex[0..2])).unwrap();
3980
3981        let err = store
3982            .get_bounded_verified(&content_ref, outside_bytes.len() as u64)
3983            .await
3984            .unwrap_err();
3985        assert!(matches!(err, StorageError::Driver { .. }), "got {err:?}");
3986    }
3987
3988    #[tokio::test]
3989    async fn put_content_ref_matches_blake3_digest() {
3990        let (_dir, store) = store(0);
3991        let bytes = b"digest check".to_vec();
3992        let content_ref = store.put(bytes.clone()).await.unwrap();
3993        let expected = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
3994        assert_eq!(content_ref, expected);
3995    }
3996
3997    #[tokio::test]
3998    async fn put_dedups_identical_content() {
3999        let (_dir, store) = store(0);
4000        let bytes = b"same bytes twice".to_vec();
4001        let first = store.put(bytes.clone()).await.unwrap();
4002        let second = store.put(bytes.clone()).await.unwrap();
4003        assert_eq!(first, second);
4004        assert_eq!(
4005            store
4006                .get_bounded_verified(&first, bytes.len() as u64)
4007                .await
4008                .unwrap(),
4009            bytes
4010        );
4011    }
4012
4013    #[tokio::test]
4014    async fn exists_reflects_put_and_delete() {
4015        let (_dir, store) = store(0);
4016        let bytes = b"exists check".to_vec();
4017        let content_ref = store.put(bytes).await.unwrap();
4018        assert!(store.exists(&content_ref).await.unwrap());
4019
4020        assert!(store.delete(&content_ref).await.unwrap());
4021        assert!(!store.exists(&content_ref).await.unwrap());
4022    }
4023
4024    #[tokio::test]
4025    async fn delete_missing_content_ref_returns_false() {
4026        let (_dir, store) = store(0);
4027        let missing = ContentRef::from_hex("f".repeat(64)).unwrap();
4028        assert!(!store.delete(&missing).await.unwrap());
4029    }
4030
4031    #[tokio::test]
4032    async fn size_reports_byte_length_for_a_present_object() {
4033        let (_dir, store) = store(0);
4034        let bytes = b"size check".to_vec();
4035        let content_ref = store.put(bytes.clone()).await.unwrap();
4036        assert_eq!(
4037            store.size(&content_ref).await.unwrap(),
4038            Some(bytes.len() as u64)
4039        );
4040    }
4041
4042    #[tokio::test]
4043    async fn size_returns_none_for_an_absent_object() {
4044        let (_dir, store) = store(0);
4045        let missing = ContentRef::from_hex("9".repeat(64)).unwrap();
4046        assert_eq!(store.size(&missing).await.unwrap(), None);
4047    }
4048
4049    #[tokio::test]
4050    async fn bounded_get_missing_content_ref_returns_not_found() {
4051        let (_dir, store) = store(0);
4052        let missing = ContentRef::from_hex("e".repeat(64)).unwrap();
4053        let err = store
4054            .get_bounded_verified(&missing, MAX_BLOB_WHOLE_BYTES)
4055            .await
4056            .unwrap_err();
4057        assert!(matches!(err, StorageError::NotFound { .. }));
4058    }
4059
4060    #[tokio::test]
4061    async fn put_refuses_below_free_space_floor() {
4062        // A floor no real disk clears -> put must fail closed, not silently
4063        // degrade or spill elsewhere (khive#292 SPEC-gate ruling).
4064        let (_dir, store) = store(u64::MAX);
4065        let err = store.put(b"too big a floor".to_vec()).await.unwrap_err();
4066        match err {
4067            StorageError::CapacityFloor {
4068                floor_bytes,
4069                available_bytes,
4070                ..
4071            } => {
4072                assert_eq!(floor_bytes, u64::MAX);
4073                assert!(available_bytes < u64::MAX);
4074            }
4075            other => panic!("expected CapacityFloor, got {other:?}"),
4076        }
4077    }
4078
4079    #[tokio::test]
4080    async fn capacity_floor_error_names_the_floor_and_volume() {
4081        let (_dir, store) = store(u64::MAX);
4082        let err = store.put(b"x".to_vec()).await.unwrap_err();
4083        let msg = err.to_string();
4084        assert!(
4085            msg.contains(&u64::MAX.to_string()),
4086            "must name the floor: {msg}"
4087        );
4088        assert!(msg.contains("Blob"), "must name the capability: {msg}");
4089    }
4090
4091    #[test]
4092    fn crosses_floor_is_write_size_aware_at_the_exact_boundary() {
4093        // Exact-boundary case, verbatim from the report: `available ==
4094        // floor_bytes + 1` must still refuse a 2-byte write. A floor-only
4095        // check (`available < floor_bytes`) would NOT catch this — 101 is
4096        // not below 100 — but the write's own size must be subtracted first.
4097        assert!(crosses_floor(101, 2, 100));
4098        assert!(!crosses_floor(101, 1, 100));
4099    }
4100
4101    #[test]
4102    fn crosses_floor_accepts_a_write_that_lands_exactly_on_the_floor() {
4103        assert!(!crosses_floor(100, 0, 100));
4104    }
4105
4106    #[test]
4107    fn crosses_floor_rejects_a_write_that_lands_one_byte_under_the_floor() {
4108        assert!(crosses_floor(100, 1, 100));
4109    }
4110
4111    #[test]
4112    fn crosses_floor_saturates_instead_of_underflowing_when_write_exceeds_available() {
4113        assert!(crosses_floor(10, 100, 50));
4114        // floor_bytes == 0 means "no floor enforced" (the convention every
4115        // other test in this file uses via `store(0)`) — even a write far
4116        // exceeding available space is not refused by the floor check itself
4117        // in that case; `saturating_sub` floors the subtraction at 0, and
4118        // `0 < 0` is false.
4119        assert!(!crosses_floor(10, 100, 0));
4120    }
4121
4122    #[test]
4123    fn put_refuses_a_write_that_would_cross_the_floor_even_though_available_alone_clears_it() {
4124        let dir = tempfile::tempdir().unwrap();
4125        let root = dir.path().join("blobs");
4126        fs::create_dir_all(&root).unwrap();
4127
4128        // Use the same blocking path as `FsBlobStore::put`, with a fixed
4129        // capacity snapshot: 101 bytes clears a 100-byte floor by itself,
4130        // but a pending two-byte write would leave only 99 bytes. Sampling
4131        // the host-wide APFS free-space gauge here made the old test flaky:
4132        // unrelated cleanup could legitimately replenish more than its
4133        // 64 MiB cushion between the test's sample and the put's sample.
4134        let err = put_blocking_with_space_probe(&root, 100, vec![7u8; 2], |_| Ok(101)).unwrap_err();
4135        assert!(
4136            matches!(err, StorageError::CapacityFloor { .. }),
4137            "a write-size-aware floor check must reject a write that pushes the volume \
4138             below the floor even though available space alone still clears it: {err:?}"
4139        );
4140    }
4141
4142    #[test]
4143    fn a_later_put_checks_a_fresh_capacity_snapshot() {
4144        let dir = tempfile::tempdir().unwrap();
4145        let root = dir.path().join("blobs");
4146        fs::create_dir_all(&root).unwrap();
4147
4148        // Model the two snapshots observed by serialized puts without tying
4149        // the assertion to a host-wide free-space gauge. The first two-byte
4150        // write may land exactly on the 100-byte floor from a 102-byte
4151        // snapshot. A later, different write sees 101 bytes and must refuse.
4152        // Mutual exclusion itself is covered deterministically below by the
4153        // shared-root lock test; together the tests prove the stale-snapshot
4154        // race is closed without relying on unrelated filesystem activity.
4155        let first = put_blocking_with_space_probe(&root, 100, vec![1u8; 2], |_| Ok(102));
4156        let second = put_blocking_with_space_probe(&root, 100, vec![2u8; 2], |_| Ok(101));
4157
4158        assert!(
4159            first.is_ok(),
4160            "the first put may land on the floor: {first:?}"
4161        );
4162        assert!(
4163            matches!(second, Err(StorageError::CapacityFloor { .. })),
4164            "the later put must use its lower capacity snapshot: {second:?}"
4165        );
4166    }
4167
4168    #[tokio::test]
4169    async fn concurrent_puts_from_two_independently_constructed_stores_share_the_root_lock() {
4170        // The actual gap in the prior fix: the
4171        // test above uses ONE `FsBlobStore` behind a shared `Arc`, so it
4172        // exercises only the per-instance mutex and cannot catch a missing
4173        // cross-instance guarantee. `StorageBackend::blob_store` constructs
4174        // a FRESH `FsBlobStore` on every call, even for the same root -- so
4175        // the real regression is two SEPARATELY CONSTRUCTED stores for the
4176        // same root. Before the shared canonical-root registry, each
4177        // store's `write_lock` was its own independent `Mutex`, and this
4178        // exact scenario would have let both puts pass the same free-space
4179        // snapshot.
4180        //
4181        // The earlier version of this test let
4182        // two real `tokio::spawn`ed puts race with no control over
4183        // interleaving -- it could PASS on the prior per-instance-mutex
4184        // bug purely because the blocking thread pool happened to run them
4185        // sequentially, which is not a deterministic regression guard.
4186        //
4187        // The first `sync_hook`-driven attempt kept proving exclusion
4188        // INDIRECTLY, through a free-space floor sized to admit exactly one
4189        // `payload_len` write -- but this dev box's real
4190        // `fs4::available_space` swings by many tens to hundreds of MB in
4191        // either direction over the several-second window the hook
4192        // orchestration takes (concurrent fleet `cargo clean`/build
4193        // activity), and no floor margin proved robust: it was observed to
4194        // both under-shoot (store_a's own write refused; available_bytes
4195        // 25521500160 vs floor_bytes 25517096960, a ~60 MiB drop) and
4196        // over-shoot (store_b's write unexpectedly SUCCEEDED after
4197        // store_a's landed) in back-to-back runs.
4198        //
4199        // Lock sharing is orthogonal to floor arithmetic -- the same
4200        // `crosses_floor`/`put_blocking` path runs regardless of which
4201        // `FsBlobStore` instance calls it, and that arithmetic is already
4202        // covered deterministically by `a_later_put_checks_a_fresh_capacity_snapshot`
4203        // and the pure `crosses_floor` unit tests above.
4204        //
4205        // The prior fix's negative proof (B
4206        // must not reach its own checkpoint) still leaned on a 200ms
4207        // `recv_timeout` as the CORRECTNESS decision -- under sufficiently
4208        // delayed scheduling, old per-instance-mutex code's B could simply
4209        // arrive after the window and every assertion would still pass,
4210        // silently defeating the regression guard. Fix: assert directly
4211        // and immediately (no timeout, no second hook, no second
4212        // `tokio::spawn` racing at all) that `store_b.write_lock` -- a
4213        // private field, reachable here because `tests` is a child module
4214        // of the module that declares it -- is ALREADY held the instant
4215        // store_a's put holds ITS guard. Under the fixed canonical-root
4216        // registry this is the exact same `Arc` store_a's own `write_lock`
4217        // resolves to, so `try_lock()` fails with zero timing dependence;
4218        // under the old per-instance-mutex code, `store_b.write_lock` is a
4219        // completely independent, unheld `Mutex`, so `try_lock()` would
4220        // succeed immediately, pinning the defect on the spot regardless
4221        // of scheduling.
4222        let dir = tempfile::tempdir().unwrap();
4223        let root = dir.path().join("blobs");
4224        fs::create_dir_all(&root).unwrap();
4225        let canonical_root = root.canonicalize().unwrap();
4226
4227        // Two INDEPENDENT `FsBlobStore::new` calls for the identical root --
4228        // exactly what `StorageBackend::blob_store` does on repeat calls.
4229        let store_a = std::sync::Arc::new(FsBlobStore::new(root.clone(), 0).unwrap());
4230        let store_b = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
4231
4232        let (a_reached, a_release, _a_done) = sync_hook::install(&canonical_root);
4233        let a = {
4234            let store_a = store_a.clone();
4235            tokio::spawn(async move { store_a.put(b"store_a payload".to_vec()).await })
4236        };
4237        assert!(
4238            recv_blocking(a_reached).await,
4239            "store_a's put must reach the sync_hook checkpoint"
4240        );
4241
4242        // The deterministic proof: store_b's OWN write_lock field must
4243        // already be unavailable while store_a holds its guard -- true
4244        // only if the two independently constructed stores share one
4245        // Arc<Mutex<()>>. No timeout, no scheduling dependence.
4246        assert!(
4247            store_b.write_lock.try_lock().is_err(),
4248            "store_b's write_lock was NOT held while store_a's put held its guard -- the two \
4249             independently constructed stores do NOT share one lock"
4250        );
4251
4252        // Release A and let it finish. Awaiting A's outer task
4253        // deterministically waits for the guard to be dropped too (see
4254        // `put`'s inner-block scoping).
4255        a_release.send(()).unwrap();
4256        let result_a = a.await.unwrap();
4257        assert!(result_a.is_ok(), "store_a's put must succeed: {result_a:?}");
4258
4259        // Liveness coverage: an ordinary put on store_b succeeds once
4260        // store_a has released the (shared) lock.
4261        let result_b = store_b.put(b"store_b payload".to_vec()).await;
4262        assert!(result_b.is_ok(), "store_b's put must succeed: {result_b:?}");
4263    }
4264
4265    #[tokio::test]
4266    async fn aborting_the_outer_put_future_does_not_release_the_guard_before_persist_completes() {
4267        // The prior fix held the write guard only
4268        // in `put`'s own async stack frame (`let _write_guard = ...
4269        // .lock().await`) while the `spawn_blocking` closure captured just
4270        // root/floor_bytes/bytes. Cancelling/dropping the outer `put`
4271        // future released that borrowed guard immediately, even though an
4272        // already-started blocking write kept running on its own thread --
4273        // a second put could then pass its floor check while the first
4274        // write was still landing.
4275        //
4276        // The earlier version of this test
4277        // proved the fix with a fixed 10ms sleep before abort and a fixed
4278        // 500x10ms poll loop waiting for the lock to free -- and the poll
4279        // loop actually FAILED once in a required-suite run (a
4280        // flaky gate, not a regression). This version uses the `sync_hook`
4281        // seam instead: `reached` fires only once execution is genuinely
4282        // inside the guarded closure (owned guard already moved in) and
4283        // blocks there until released; `done` fires only after the guard
4284        // has actually been dropped (see `put`'s inner-block scoping) --
4285        // both edges event-driven, no sleeps, no polling.
4286        let dir = tempfile::tempdir().unwrap();
4287        let root = dir.path().join("blobs");
4288        fs::create_dir_all(&root).unwrap();
4289        let canonical_root = root.canonicalize().unwrap();
4290
4291        let store = std::sync::Arc::new(FsBlobStore::new(root, 0).unwrap());
4292        let (reached, release, done) = sync_hook::install(&canonical_root);
4293        let handle = {
4294            let store = store.clone();
4295            tokio::spawn(async move { store.put(b"cancellation race payload".to_vec()).await })
4296        };
4297
4298        assert!(
4299            recv_blocking(reached).await,
4300            "put must reach the sync_hook checkpoint -- owned guard already moved into the \
4301             closure -- before this test can mean anything"
4302        );
4303
4304        handle.abort();
4305        let abort_result = handle.await;
4306        match &abort_result {
4307            Err(e) if e.is_cancelled() => {}
4308            other => panic!(
4309                "the outer task must actually have been cancelled for this test to be \
4310                 meaningful: {other:?}"
4311            ),
4312        }
4313
4314        let shared_lock = write_lock_for_root(&canonical_root).unwrap();
4315        assert!(
4316            shared_lock.try_lock().is_err(),
4317            "the guard must still be held by the detached blocking write immediately after \
4318             the outer future was cancelled -- if this is free, the guard was released with \
4319             the aborted frame instead of moving into the spawn_blocking closure"
4320        );
4321
4322        // Let the detached write proceed and finish, then wait for its
4323        // explicit completion signal -- no polling, no fixed durations.
4324        // `done` only fires after the guard is actually dropped (see
4325        // `put`), so the very next check is race-free.
4326        release.send(()).unwrap();
4327        assert!(
4328            recv_blocking(done).await,
4329            "the detached write must signal completion once it actually persists"
4330        );
4331        assert!(
4332            shared_lock.try_lock().is_ok(),
4333            "the guard must be free once the detached write's completion was observed"
4334        );
4335    }
4336
4337    /// `orphan_sweep` is disabled in this compatibility release: it has no
4338    /// `SqlAccess` capability with which to prove a completed V21 epoch, so
4339    /// every call refuses regardless of `live_refs` contents or `dry_run`,
4340    /// and nothing on disk is ever touched. This replaces the prior
4341    /// `orphan_sweep_race_demonstrates_the_documented_quiescence_requirement`
4342    /// regression, which pinned the now-eliminated caller-snapshot deletion
4343    /// hazard this disablement closes.
4344    #[tokio::test]
4345    async fn orphan_sweep_is_disabled_in_both_modes_regardless_of_live_refs() {
4346        let (_dir, store) = store(0);
4347        let blob = store
4348            .put(b"never swept by this API".to_vec())
4349            .await
4350            .unwrap();
4351        let mut live_refs = std::collections::HashSet::new();
4352        live_refs.insert(blob.clone());
4353
4354        for dry_run in [true, false] {
4355            let error = store
4356                .orphan_sweep(&BlobOrphanSweepConfig {
4357                    live_refs: live_refs.clone(),
4358                    dry_run,
4359                })
4360                .await
4361                .expect_err("caller-snapshot orphan_sweep must be disabled");
4362            assert!(
4363                matches!(error, StorageError::Unsupported { .. }),
4364                "expected typed Unsupported, got {error:?}"
4365            );
4366        }
4367        assert!(store.exists(&blob).await.unwrap());
4368    }
4369
4370    /// Rollout compatibility fence: the Phase-3 binary's V20 schema cannot
4371    /// represent a moodboard model's nested FANN network as SQL liveness.
4372    /// Both report-only and destructive transactional sweeps must therefore
4373    /// refuse before taking the root lock or mutating abandoned claims.
4374    #[tokio::test]
4375    async fn transactional_orphan_sweep_refuses_v20_before_root_or_claim_mutation() {
4376        let dir = tempfile::tempdir().unwrap();
4377        let db_path = dir.path().join("khive.db");
4378        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4379        {
4380            let mut writer = backend.pool().writer().unwrap();
4381            prepare_v20_gc_fixture(writer.conn_mut());
4382        }
4383
4384        let root = dir.path().join("blobs");
4385        let store = Arc::new(
4386            FsBlobStore::new(root, 0)
4387                .unwrap()
4388                .with_orphan_sweep_grace(Duration::ZERO),
4389        );
4390        let bundle = store.put(b"legacy model bundle".to_vec()).await.unwrap();
4391        let network = store.put(b"legacy FANN network".to_vec()).await.unwrap();
4392        let orphan = store.put(b"ordinary old orphan".to_vec()).await.unwrap();
4393        let abandoned_ref = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
4394        {
4395            let writer = backend.pool().writer().unwrap();
4396            writer
4397                .conn()
4398                .execute(
4399                    "INSERT INTO entities \
4400                     (id, namespace, kind, entity_type, name, tags, created_at, updated_at, \
4401                      content_ref) \
4402                     VALUES ('legacy-model', 'local', 'artifact', 'moodboard_model', \
4403                             'legacy model', '[]', 1, 1, ?1)",
4404                    [bundle.as_str()],
4405                )
4406                .unwrap();
4407            writer
4408                .conn()
4409                .execute(
4410                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4411                     VALUES ('abandoned-before-compat', ?1, 1)",
4412                    [abandoned_ref],
4413                )
4414                .unwrap();
4415        }
4416
4417        // If the compatibility gate is below the root lock, the call parks
4418        // here and the timeout fails. A correct V20 refusal never waits for it.
4419        let _root_guard = store.write_lock.clone().lock_owned().await;
4420        for dry_run in [true, false] {
4421            let outcome = tokio::time::timeout(
4422                Duration::from_secs(1),
4423                store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
4424            )
4425            .await
4426            .expect("V20 refusal must happen before waiting for the held root lock");
4427            let error = outcome.expect_err("V20 transactional sweep must be disabled");
4428            match error {
4429                StorageError::Unsupported {
4430                    capability: StorageCapability::Blob,
4431                    operation,
4432                    message,
4433                } => {
4434                    assert_eq!(operation, "transactional_orphan_sweep");
4435                    assert!(
4436                        message.contains("complete V21 attachment cutover"),
4437                        "unexpected compatibility diagnostic: {message}"
4438                    );
4439                }
4440                other => panic!("expected typed Unsupported refusal, got {other:?}"),
4441            }
4442        }
4443
4444        let reader = backend.pool().reader().unwrap();
4445        let abandoned: i64 = reader
4446            .conn()
4447            .query_row(
4448                "SELECT COUNT(*) FROM blob_gc_claims \
4449                 WHERE root_key = 'abandoned-before-compat' AND content_ref = ?1",
4450                [abandoned_ref],
4451                |row| row.get(0),
4452            )
4453            .unwrap();
4454        assert_eq!(abandoned, 1, "V20 refusal must not clean abandoned claims");
4455        drop(reader);
4456        assert!(store.exists(&bundle).await.unwrap());
4457        assert!(store.exists(&network).await.unwrap());
4458        assert!(store.exists(&orphan).await.unwrap());
4459    }
4460
4461    /// The durable marker is authoritative, not the mere presence of V21
4462    /// tables, triggers, or even a ledger row. An interrupted/inconsistent
4463    /// cutover remains non-sweepable for both modes and is rejected before
4464    /// the root wait or abandoned-claim recovery.
4465    #[tokio::test]
4466    async fn transactional_orphan_sweep_refuses_incomplete_v21_marker_without_mutation() {
4467        let dir = tempfile::tempdir().unwrap();
4468        let db_path = dir.path().join("khive.db");
4469        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4470        let abandoned_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
4471        {
4472            let mut writer = backend.pool().writer().unwrap();
4473            prepare_completed_v21_gc_fixture(writer.conn_mut());
4474            writer
4475                .conn_mut()
4476                .execute(
4477                    "UPDATE attachment_cutover_state \
4478                     SET state = 'incomplete', completed_at = NULL \
4479                     WHERE singleton = 1",
4480                    [],
4481                )
4482                .unwrap();
4483            writer
4484                .conn_mut()
4485                .execute(
4486                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4487                     VALUES ('abandoned-incomplete-v21', ?1, 1)",
4488                    [abandoned_ref],
4489                )
4490                .unwrap();
4491        }
4492
4493        let store = Arc::new(
4494            FsBlobStore::new(dir.path().join("blobs"), 0)
4495                .unwrap()
4496                .with_orphan_sweep_grace(Duration::ZERO),
4497        );
4498        let orphan = store.put(b"incomplete V21 orphan".to_vec()).await.unwrap();
4499        let _root_guard = store.write_lock.clone().lock_owned().await;
4500
4501        for dry_run in [true, false] {
4502            let outcome = tokio::time::timeout(
4503                Duration::from_secs(1),
4504                store.transactional_orphan_sweep(backend.sql().as_ref(), dry_run),
4505            )
4506            .await
4507            .expect("incomplete V21 must refuse before waiting for the root lock");
4508            assert!(
4509                matches!(outcome, Err(StorageError::Unsupported { .. })),
4510                "incomplete V21 must return typed Unsupported: {outcome:?}"
4511            );
4512        }
4513
4514        let remaining: i64 = backend
4515            .pool()
4516            .reader()
4517            .unwrap()
4518            .conn()
4519            .query_row(
4520                "SELECT COUNT(*) FROM blob_gc_claims \
4521                 WHERE root_key = 'abandoned-incomplete-v21' AND content_ref = ?1",
4522                [abandoned_ref],
4523                |row| row.get(0),
4524            )
4525            .unwrap();
4526        assert_eq!(remaining, 1, "refusal must not recover abandoned claims");
4527        assert!(store.exists(&orphan).await.unwrap());
4528    }
4529
4530    /// Both public sweep APIs must refuse under the same non-completed
4531    /// epochs, in both `dry_run` modes, without touching files or claims.
4532    /// `orphan_sweep` ignores the fixture's `SqlAccess` entirely (it has no
4533    /// epoch capability of its own -- that is exactly why it is disabled),
4534    /// so this proves the disablement holds even when a caller has a real,
4535    /// otherwise-plausible database sitting right next to it.
4536    #[tokio::test]
4537    async fn both_sweep_apis_refuse_v20_and_incomplete_v21_epochs_in_both_modes() {
4538        for incomplete_v21 in [false, true] {
4539            let dir = tempfile::tempdir().unwrap();
4540            let db_path = dir.path().join("khive.db");
4541            let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4542            {
4543                let mut writer = backend.pool().writer().unwrap();
4544                if incomplete_v21 {
4545                    prepare_completed_v21_gc_fixture(writer.conn_mut());
4546                    writer
4547                        .conn_mut()
4548                        .execute(
4549                            "UPDATE attachment_cutover_state \
4550                             SET state = 'incomplete', completed_at = NULL \
4551                             WHERE singleton = 1",
4552                            [],
4553                        )
4554                        .unwrap();
4555                } else {
4556                    prepare_v20_gc_fixture(writer.conn_mut());
4557                }
4558            }
4559
4560            let store = FsBlobStore::new(dir.path().join("blobs"), 0)
4561                .unwrap()
4562                .with_orphan_sweep_grace(Duration::ZERO);
4563            let orphan = store
4564                .put(format!("both-apis orphan (incomplete_v21={incomplete_v21})").into_bytes())
4565                .await
4566                .unwrap();
4567
4568            // A known claim, seeded once per epoch case, must survive every
4569            // refused arm below untouched -- refusal must never recover or
4570            // otherwise mutate an existing claim.
4571            let known_claim_ref = "f".repeat(64);
4572            {
4573                let writer = backend.pool().writer().unwrap();
4574                writer
4575                    .conn()
4576                    .execute(
4577                        "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4578                         VALUES ('both-api-known-claim', ?1, 1)",
4579                        [known_claim_ref.as_str()],
4580                    )
4581                    .unwrap();
4582            }
4583            let assert_known_claim_unchanged = |arm: &str| {
4584                let remaining: i64 = backend
4585                    .pool()
4586                    .reader()
4587                    .unwrap()
4588                    .conn()
4589                    .query_row(
4590                        "SELECT COUNT(*) FROM blob_gc_claims \
4591                         WHERE root_key = 'both-api-known-claim' AND content_ref = ?1",
4592                        [known_claim_ref.as_str()],
4593                        |row| row.get(0),
4594                    )
4595                    .unwrap();
4596                assert_eq!(
4597                    remaining, 1,
4598                    "incomplete_v21={incomplete_v21} arm={arm}: refusal must not mutate \
4599                     an existing claim"
4600                );
4601            };
4602
4603            for dry_run in [true, false] {
4604                let snapshot_error = store
4605                    .orphan_sweep(&BlobOrphanSweepConfig {
4606                        live_refs: std::collections::HashSet::new(),
4607                        dry_run,
4608                    })
4609                    .await
4610                    .expect_err("orphan_sweep must refuse regardless of epoch");
4611                assert!(
4612                    matches!(snapshot_error, StorageError::Unsupported { .. }),
4613                    "incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
4614                     from orphan_sweep, got {snapshot_error:?}"
4615                );
4616                assert_known_claim_unchanged(&format!("orphan_sweep dry_run={dry_run}"));
4617
4618                let transactional_error = store
4619                    .transactional_orphan_sweep(backend.sql().as_ref(), dry_run)
4620                    .await
4621                    .expect_err("transactional_orphan_sweep must refuse this epoch");
4622                assert!(
4623                    matches!(transactional_error, StorageError::Unsupported { .. }),
4624                    "incomplete_v21={incomplete_v21} dry_run={dry_run}: expected Unsupported \
4625                     from transactional_orphan_sweep, got {transactional_error:?}"
4626                );
4627                assert_known_claim_unchanged(&format!(
4628                    "transactional_orphan_sweep dry_run={dry_run}"
4629                ));
4630            }
4631
4632            assert!(
4633                store.exists(&orphan).await.unwrap(),
4634                "incomplete_v21={incomplete_v21}: a refused sweep must not delete anything"
4635            );
4636        }
4637    }
4638
4639    /// The epoch recheck taken under database ownership must run before the
4640    /// sweep ever waits on the root guard/lock, even when the epoch was
4641    /// still valid at the read-only preflight and only regressed while the
4642    /// call was waiting for database ownership. Without the fix this recheck
4643    /// used to run after the root guard, filesystem root lock, and directory
4644    /// walk -- this pins the corrected ordering with a real cross-task race
4645    /// instead of trusting the comment above it.
4646    ///
4647    /// The externally held blocking point is the OS-level root advisory lock
4648    /// (`acquire_root_write_lock`), not the process-local `write_lock`
4649    /// mutex. The old (pre-fix) ordering acquired the process-local mutex
4650    /// before ever reaching database ownership -- holding that mutex here
4651    /// would have starved the sweep task before it ever reached the hook
4652    /// below, hanging this test instead of failing it. The OS-level root
4653    /// lock is only ever acquired after database ownership under both the
4654    /// old and the new ordering, so both orderings reach the hook; only the
4655    /// old ordering then blocks trying to acquire the lock held here,
4656    /// because its (mis-placed) recheck runs after that acquisition instead
4657    /// of before it.
4658    #[tokio::test]
4659    async fn transactional_orphan_sweep_recheck_refuses_before_root_lock_when_epoch_regresses_after_db_ownership(
4660    ) {
4661        let dir = tempfile::tempdir().unwrap();
4662        let db_path = dir.path().join("khive.db");
4663        let backend = Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4664        {
4665            let mut writer = backend.pool().writer().unwrap();
4666            prepare_completed_v21_gc_fixture(writer.conn_mut());
4667        }
4668        assert!(blob_gc_fencing_complete(backend.sql().as_ref())
4669            .await
4670            .unwrap());
4671
4672        let blob_root = dir.path().join("blobs");
4673        let store = Arc::new(
4674            FsBlobStore::new(blob_root.clone(), 0)
4675                .unwrap()
4676                .with_orphan_sweep_grace(Duration::ZERO),
4677        );
4678        let orphan = store
4679            .put(b"epoch regressed after database ownership".to_vec())
4680            .await
4681            .unwrap();
4682
4683        // Hold the same OS-level root lock the real sweep acquires, on the
4684        // same canonicalized path it canonicalizes to.
4685        let canonical_root = blob_root.canonicalize().unwrap();
4686        let _root_write_guard = acquire_root_write_lock(&canonical_root).unwrap();
4687
4688        // Key the hook by the same canonicalized path `SqlAccess::database_path`
4689        // returns -- not the raw `db_path` above -- since that's what the
4690        // implementation looks the hook up by.
4691        let canonical_db_path = backend.sql().database_path();
4692        let (reached, release) = db_ownership_sync_hook::install(canonical_db_path.as_deref());
4693        let sweep_store = store.clone();
4694        let sweep_backend = backend.clone();
4695        let handle = tokio::spawn(async move {
4696            sweep_store
4697                .transactional_orphan_sweep(sweep_backend.sql().as_ref(), false)
4698                .await
4699        });
4700
4701        // The hook fires only once database ownership (process-local guard
4702        // plus the cross-process advisory lock) is held, strictly after the
4703        // read-only preflight already observed a valid epoch. Bounded so a
4704        // regression that drops the hook call (or blocks ahead of it) fails
4705        // this test instead of hanging it.
4706        let reached_signal = tokio::time::timeout(Duration::from_secs(1), recv_blocking(reached))
4707            .await
4708            .expect("the sweep must reach database ownership before this test's timeout");
4709        assert!(reached_signal, "hook sender was dropped before signaling");
4710        {
4711            let writer = backend.pool().writer().unwrap();
4712            writer
4713                .conn()
4714                .execute("DELETE FROM _schema_migrations WHERE version = 21", [])
4715                .unwrap();
4716        }
4717        release.send(()).unwrap();
4718
4719        let outcome = tokio::time::timeout(Duration::from_secs(1), handle)
4720            .await
4721            .expect(
4722                "the recheck must refuse before ever waiting on the externally held root lock -- \
4723                 under the old (pre-fix) ordering this join times out instead, because the \
4724                 sweep blocks acquiring the OS-level root lock held above",
4725            )
4726            .unwrap();
4727        assert!(
4728            matches!(outcome, Err(StorageError::Unsupported { .. })),
4729            "expected the regressed epoch to be caught immediately after database ownership: \
4730             {outcome:?}"
4731        );
4732        assert!(store.exists(&orphan).await.unwrap());
4733    }
4734
4735    /// A Phase-4a binary can remain in a mixed fleet after a newer binary has
4736    /// atomically completed V21. It must then use every attachment role as
4737    /// liveness, including a moodboard FANN network, and delete only the true
4738    /// orphan.
4739    #[tokio::test]
4740    async fn transactional_orphan_sweep_accepts_completed_v21_attachment_liveness() {
4741        let dir = tempfile::tempdir().unwrap();
4742        let db_path = dir.path().join("khive.db");
4743        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
4744        {
4745            let mut writer = backend.pool().writer().unwrap();
4746            prepare_completed_v21_gc_fixture(writer.conn_mut());
4747        }
4748
4749        let store = FsBlobStore::new(dir.path().join("blobs"), 0)
4750            .unwrap()
4751            .with_orphan_sweep_grace(Duration::ZERO);
4752        let bundle = store.put(b"V21 model bundle".to_vec()).await.unwrap();
4753        let network = store.put(b"V21 FANN network".to_vec()).await.unwrap();
4754        let orphan = store.put(b"V21 true orphan".to_vec()).await.unwrap();
4755        {
4756            let writer = backend.pool().writer().unwrap();
4757            writer
4758                .conn()
4759                .execute(
4760                    "INSERT INTO entities \
4761                     (id, namespace, kind, entity_type, name, tags, created_at, updated_at) \
4762                     VALUES ('model', 'local', 'artifact', 'moodboard_model', \
4763                             'model', '[]', 1, 1)",
4764                    [],
4765                )
4766                .unwrap();
4767            writer
4768                .conn()
4769                .execute(
4770                    "INSERT INTO attachments \
4771                     (record_uuid, substrate, role, content_ref, created_at) \
4772                     VALUES ('model', 'entity', 'content', ?1, 1), \
4773                            ('model', 'entity', 'fann-network', ?2, 1)",
4774                    rusqlite::params![bundle.as_str(), network.as_str()],
4775                )
4776                .unwrap();
4777        }
4778
4779        let dry_run = store
4780            .transactional_orphan_sweep(backend.sql().as_ref(), true)
4781            .await
4782            .expect("completed V21 dry run must be supported");
4783        assert_eq!(dry_run.would_delete, 1);
4784        assert_eq!(dry_run.deleted, 0);
4785
4786        let result = store
4787            .transactional_orphan_sweep(backend.sql().as_ref(), false)
4788            .await
4789            .expect("completed V21 destructive sweep must be supported");
4790        assert_eq!(result.deleted, 1);
4791        assert!(store.exists(&bundle).await.unwrap());
4792        assert!(store.exists(&network).await.unwrap());
4793        assert!(!store.exists(&orphan).await.unwrap());
4794    }
4795
4796    /// A `StorageBackend` constructed directly and never run through the
4797    /// versioned migration ledger (`run_migrations`/`prepare_core_schema`) —
4798    /// only the ad hoc, idempotent `entities` DDL a plain `entities()` call
4799    /// applies — has no completed V21 marker, attachment liveness table, or
4800    /// attachment fencing triggers. Without that set a reference committed between liveness
4801    /// selection and physical deletion would dangle, so the trait contract
4802    /// requires `StorageError::Unsupported` here rather than an unfenced
4803    /// sweep, and every candidate must survive.
4804    #[tokio::test]
4805    async fn transactional_orphan_sweep_refuses_without_the_blob_gc_claims_migration() {
4806        let dir = tempfile::tempdir().unwrap();
4807        let db_path = dir.path().join("khive.db");
4808        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4809        backend.entities().unwrap();
4810        {
4811            let reader = backend.pool().reader().unwrap();
4812            let present: bool = reader
4813                .conn()
4814                .query_row(
4815                    "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = 'table' \
4816                     AND name = 'blob_gc_claims'",
4817                    [],
4818                    |row| row.get(0),
4819                )
4820                .unwrap();
4821            assert!(
4822                !present,
4823                "this test's premise requires blob_gc_claims to be absent"
4824            );
4825        }
4826
4827        let root = dir.path().join("blobs");
4828        let store = std::sync::Arc::new(
4829            FsBlobStore::new(root.clone(), 0)
4830                .unwrap()
4831                .with_orphan_sweep_grace(Duration::ZERO),
4832        );
4833        let orphan = store.put(b"direct-backend orphan".to_vec()).await.unwrap();
4834
4835        let sql = backend.sql();
4836        let error = store
4837            .transactional_orphan_sweep(sql.as_ref(), false)
4838            .await
4839            .expect_err("sweep must refuse a backend without the blob_gc_claims fencing set");
4840        assert!(
4841            matches!(error, StorageError::Unsupported { .. }),
4842            "expected StorageError::Unsupported, got {error:?}"
4843        );
4844        assert!(
4845            store.exists(&orphan).await.unwrap(),
4846            "a refused sweep must not have deleted anything"
4847        );
4848    }
4849
4850    #[tokio::test]
4851    async fn transactional_orphan_sweep_refuses_an_incomplete_cutover_marker() {
4852        let dir = tempfile::tempdir().unwrap();
4853        let db_path = dir.path().join("khive.db");
4854        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4855        {
4856            let mut writer = backend.pool().writer().unwrap();
4857            prepare_completed_v21_gc_fixture(writer.conn_mut());
4858            writer
4859                .conn_mut()
4860                .execute_batch(
4861                    "UPDATE attachment_cutover_state \
4862                     SET state = 'incomplete', completed_at = NULL WHERE singleton = 1; \
4863                     DELETE FROM _schema_migrations WHERE version = 21;",
4864                )
4865                .unwrap();
4866        }
4867
4868        let store = FsBlobStore::new(dir.path().join("blobs"), 0)
4869            .unwrap()
4870            .with_orphan_sweep_grace(Duration::ZERO);
4871        let orphan = store
4872            .put(b"incomplete-cutover orphan".to_vec())
4873            .await
4874            .unwrap();
4875        let error = store
4876            .transactional_orphan_sweep(backend.sql().as_ref(), false)
4877            .await
4878            .expect_err("sweep must refuse every durable incomplete marker");
4879        assert!(matches!(error, StorageError::Unsupported { .. }));
4880        assert!(
4881            store.exists(&orphan).await.unwrap(),
4882            "refused incomplete-state sweep must preserve every blob"
4883        );
4884    }
4885
4886    /// The fencing gate must demand the complete V21 set, not just the
4887    /// claims table: with a fencing trigger dropped, a claim no longer
4888    /// blocks a concurrent attachment write from resurrecting the digest, so
4889    /// the sweep must refuse exactly as it does with no migration at all.
4890    #[tokio::test]
4891    async fn transactional_orphan_sweep_refuses_with_incomplete_fencing_triggers() {
4892        let dir = tempfile::tempdir().unwrap();
4893        let db_path = dir.path().join("khive.db");
4894        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4895        {
4896            let mut writer = backend.pool().writer().unwrap();
4897            prepare_completed_v21_gc_fixture(writer.conn_mut());
4898            writer
4899                .conn_mut()
4900                .execute_batch("DROP TRIGGER attachments_reject_claimed_blob_update")
4901                .unwrap();
4902            writer
4903                .conn_mut()
4904                .execute(
4905                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
4906                     VALUES ('abandoned-partial-fence', \
4907                             'dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd', \
4908                             1)",
4909                    [],
4910                )
4911                .unwrap();
4912        }
4913
4914        let root = dir.path().join("blobs");
4915        let store = std::sync::Arc::new(
4916            FsBlobStore::new(root.clone(), 0)
4917                .unwrap()
4918                .with_orphan_sweep_grace(Duration::ZERO),
4919        );
4920        let orphan = store.put(b"partial-fence orphan".to_vec()).await.unwrap();
4921
4922        let sql = backend.sql();
4923        let _root_guard = store.write_lock.clone().lock_owned().await;
4924        let error = tokio::time::timeout(
4925            Duration::from_secs(1),
4926            store.transactional_orphan_sweep(sql.as_ref(), false),
4927        )
4928        .await
4929        .expect("an incomplete V21 fence must refuse before the root wait")
4930        .expect_err("sweep must refuse when any V21 fencing trigger is missing");
4931        assert!(
4932            matches!(error, StorageError::Unsupported { .. }),
4933            "expected StorageError::Unsupported, got {error:?}"
4934        );
4935        assert!(
4936            store.exists(&orphan).await.unwrap(),
4937            "a refused sweep must not have deleted anything"
4938        );
4939        let remaining: i64 = backend
4940            .pool()
4941            .reader()
4942            .unwrap()
4943            .conn()
4944            .query_row(
4945                "SELECT COUNT(*) FROM blob_gc_claims \
4946                 WHERE root_key = 'abandoned-partial-fence'",
4947                [],
4948                |row| row.get(0),
4949            )
4950            .unwrap();
4951        assert_eq!(remaining, 1, "a refused sweep must not recover claims");
4952    }
4953
4954    /// The gate must verify the fence FUNCTIONS, not that three names exist
4955    /// in `sqlite_master`: triggers with the right names but no-op bodies
4956    /// pass any name census while letting a claimed `content_ref` become
4957    /// live during the released-writer deletion window. The fence probe must
4958    /// catch them and refuse, deleting nothing and leaving no probe residue.
4959    #[tokio::test]
4960    async fn transactional_orphan_sweep_refuses_same_named_noop_fencing_triggers() {
4961        let dir = tempfile::tempdir().unwrap();
4962        let db_path = dir.path().join("khive.db");
4963        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
4964        {
4965            let mut writer = backend.pool().writer().unwrap();
4966            prepare_completed_v21_gc_fixture(writer.conn_mut());
4967            writer
4968                .conn_mut()
4969                .execute_batch(
4970                    "DROP TRIGGER attachments_reject_claimed_blob_insert; \
4971                     DROP TRIGGER attachments_reject_claimed_blob_update; \
4972                     CREATE TRIGGER attachments_reject_claimed_blob_insert \
4973                     BEFORE INSERT ON attachments BEGIN SELECT 0; END; \
4974                     CREATE TRIGGER attachments_reject_claimed_blob_update \
4975                     BEFORE UPDATE OF content_ref ON attachments \
4976                     BEGIN SELECT 0; END;",
4977                )
4978                .unwrap();
4979        }
4980
4981        let root = dir.path().join("blobs");
4982        let store = std::sync::Arc::new(
4983            FsBlobStore::new(root.clone(), 0)
4984                .unwrap()
4985                .with_orphan_sweep_grace(Duration::ZERO),
4986        );
4987        let orphan = store.put(b"noop-trigger orphan".to_vec()).await.unwrap();
4988
4989        let sql = backend.sql();
4990        let error = store
4991            .transactional_orphan_sweep(sql.as_ref(), false)
4992            .await
4993            .expect_err("sweep must refuse when the fencing triggers are same-named no-ops");
4994        assert!(
4995            matches!(error, StorageError::Unsupported { .. }),
4996            "expected StorageError::Unsupported, got {error:?}"
4997        );
4998        assert!(
4999            store.exists(&orphan).await.unwrap(),
5000            "a refused sweep must not have deleted anything"
5001        );
5002
5003        // The probe must not leave residue behind either.
5004        let reader = backend.pool().reader().unwrap();
5005        let leftovers: i64 = reader
5006            .conn()
5007            .query_row(
5008                "SELECT (SELECT COUNT(*) FROM blob_gc_claims \
5009                         WHERE root_key GLOB '__fence_probe-*') \
5010                      + (SELECT COUNT(*) FROM attachments \
5011                         WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
5012                [],
5013                |row| row.get(0),
5014            )
5015            .unwrap();
5016        assert_eq!(leftovers, 0, "fence probe rows must not survive the probe");
5017    }
5018
5019    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5020    async fn fence_probe_refuses_id_collision_and_preserves_the_colliding_attachment() {
5021        let dir = tempfile::tempdir().unwrap();
5022        let db_path = dir.path().join("khive.db");
5023        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5024        {
5025            let mut writer = backend.pool().writer().unwrap();
5026            prepare_completed_v21_gc_fixture(writer.conn_mut());
5027            writer
5028                .conn_mut()
5029                .execute(
5030                    "INSERT INTO attachments \
5031                     (record_uuid, substrate, role, content_ref, media_type, created_at) \
5032                     VALUES ('victim-id', 'entity', 'content', \
5033                             '2222222222222222222222222222222222222222222222222222222222222222', \
5034                             'application/test', 7)",
5035                    [],
5036                )
5037                .unwrap();
5038        }
5039
5040        let sql = backend.sql();
5041        let error = super::blob_gc_fence_probe_with_ids(
5042            sql.as_ref(),
5043            "victim-id".to_string(),
5044            "victim-update-id".to_string(),
5045            "victim-insert2-id".to_string(),
5046            "victim-update2-id".to_string(),
5047            "victim-claim-key".to_string(),
5048        )
5049        .await
5050        .expect_err("the probe must refuse when an id it would delete already names a row");
5051        assert!(
5052            matches!(error, StorageError::Unsupported { .. }),
5053            "expected StorageError::Unsupported, got {error:?}"
5054        );
5055
5056        let reader = backend.pool().reader().unwrap();
5057        let (media_type, created_at): (String, i64) = reader
5058            .conn()
5059            .query_row(
5060                "SELECT media_type, created_at FROM attachments \
5061                 WHERE record_uuid = 'victim-id' AND role = 'content'",
5062                [],
5063                |row| Ok((row.get(0)?, row.get(1)?)),
5064            )
5065            .expect("the colliding attachment must survive the refused probe untouched");
5066        assert_eq!(media_type, "application/test");
5067        assert_eq!(created_at, 7);
5068    }
5069
5070    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5071    async fn fence_probe_does_not_touch_an_unrelated_retained_entity_sequence() {
5072        let dir = tempfile::tempdir().unwrap();
5073        let db_path = dir.path().join("khive.db");
5074        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5075        {
5076            let mut writer = backend.pool().writer().unwrap();
5077            prepare_completed_v21_gc_fixture(writer.conn_mut());
5078            // entities_seq rows intentionally survive entity hard deletion, so
5079            // an id can collide with the ledger alone — no entities row left
5080            // for the guard to trip on.
5081            writer
5082                .conn_mut()
5083                .execute(
5084                    "INSERT INTO entities \
5085                     (id, namespace, kind, name, tags, created_at, updated_at) \
5086                     VALUES ('retained-id', 'local', 'document', 'gone entity', '[]', 7, 7)",
5087                    [],
5088                )
5089                .unwrap();
5090            writer
5091                .conn_mut()
5092                .execute("DELETE FROM entities WHERE id = 'retained-id'", [])
5093                .unwrap();
5094            let retained: i64 = writer
5095                .conn_mut()
5096                .query_row(
5097                    "SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
5098                    [],
5099                    |row| row.get(0),
5100                )
5101                .unwrap();
5102            assert_eq!(retained, 1, "fixture requires a retained-only ledger row");
5103        }
5104
5105        let sql = backend.sql();
5106        super::blob_gc_fence_probe_with_ids(
5107            sql.as_ref(),
5108            "retained-id".to_string(),
5109            "retained-update-id".to_string(),
5110            "retained-insert2-id".to_string(),
5111            "retained-update2-id".to_string(),
5112            "retained-claim-key".to_string(),
5113        )
5114        .await
5115        .expect("attachment probe has no reason to mutate an entity sequence row");
5116
5117        let reader = backend.pool().reader().unwrap();
5118        let survivors: i64 = reader
5119            .conn()
5120            .query_row(
5121                "SELECT COUNT(*) FROM entities_seq WHERE entity_id = 'retained-id'",
5122                [],
5123                |row| row.get(0),
5124            )
5125            .unwrap();
5126        assert_eq!(
5127            survivors, 1,
5128            "the retained entity ledger row must survive the attachment probe"
5129        );
5130    }
5131
5132    fn nul_embedded_canonical_ref() -> String {
5133        let mut polluted = "a".repeat(64);
5134        polluted.push('\0');
5135        polluted.push_str("zz");
5136        polluted
5137    }
5138
5139    /// SQLite's `length()` counts characters before the first NUL and GLOB
5140    /// stops scanning there, so a 64-hex-then-NUL value passes both while the
5141    /// exact-equality liveness anti-join cannot match it. The byte-length arm
5142    /// must refuse the sweep on such a claim row.
5143    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5144    async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref() {
5145        let dir = tempfile::tempdir().unwrap();
5146        let db_path = dir.path().join("khive.db");
5147        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5148        {
5149            let mut writer = backend.pool().writer().unwrap();
5150            prepare_completed_v21_gc_fixture(writer.conn_mut());
5151            writer
5152                .conn_mut()
5153                .execute(
5154                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5155                     VALUES ('nul-claim-key', ?1, 0)",
5156                    rusqlite::params![nul_embedded_canonical_ref()],
5157                )
5158                .unwrap();
5159        }
5160
5161        let sql = backend.sql();
5162        let error = super::validate_blob_gc_evidence(sql.as_ref())
5163            .await
5164            .expect_err("a NUL-embedded claim ref must refuse the sweep");
5165        assert!(
5166            error.to_string().contains("blob_gc_claims"),
5167            "expected the claims-table refusal, got {error:?}"
5168        );
5169    }
5170
5171    /// The attachments schema CHECK uses the same NUL-blind `length()`/GLOB
5172    /// pair, so the polluted row INSERTS successfully — this test proves that
5173    /// on purpose — and the evidence validator must then be the backstop that
5174    /// refuses the sweep.
5175    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5176    async fn blob_gc_evidence_rejects_a_nul_embedded_attachment_ref() {
5177        let dir = tempfile::tempdir().unwrap();
5178        let db_path = dir.path().join("khive.db");
5179        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5180        {
5181            let mut writer = backend.pool().writer().unwrap();
5182            prepare_completed_v21_gc_fixture(writer.conn_mut());
5183            // The table CHECK now rejects a NUL-embedded ref at admission
5184            // time; bypass it to simulate a row that reached this state some
5185            // other way (e.g. a pre-fix legacy row) and prove the validator
5186            // is still defense-in-depth against it.
5187            writer
5188                .conn_mut()
5189                .execute_batch("PRAGMA ignore_check_constraints = ON")
5190                .unwrap();
5191            writer
5192                .conn_mut()
5193                .execute(
5194                    "INSERT INTO attachments \
5195                     (record_uuid, substrate, role, content_ref, created_at) \
5196                     VALUES ('nul-attachment-id', 'entity', 'content', ?1, 0)",
5197                    rusqlite::params![nul_embedded_canonical_ref()],
5198                )
5199                .expect("ignore_check_constraints must allow the corrupt row to insert");
5200            writer
5201                .conn_mut()
5202                .execute_batch("PRAGMA ignore_check_constraints = OFF")
5203                .unwrap();
5204        }
5205
5206        let sql = backend.sql();
5207        let error = super::validate_blob_gc_evidence(sql.as_ref())
5208            .await
5209            .expect_err("a NUL-embedded attachment ref must refuse the sweep");
5210        assert!(
5211            error.to_string().contains("attachments"),
5212            "expected the attachments-table refusal, got {error:?}"
5213        );
5214    }
5215
5216    /// A trigger rewrite that fences only the probe's fixed all-zero sentinel
5217    /// passes the sentinel arms; the second-digest arms must catch it. The
5218    /// healthy fixture is probed first as the positive control.
5219    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5220    async fn fence_probe_refuses_a_digest_restricted_trigger_rewrite() {
5221        let dir = tempfile::tempdir().unwrap();
5222        let db_path = dir.path().join("khive.db");
5223        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5224        {
5225            let mut writer = backend.pool().writer().unwrap();
5226            prepare_completed_v21_gc_fixture(writer.conn_mut());
5227        }
5228
5229        let sql = backend.sql();
5230        super::blob_gc_fence_probe(sql.as_ref())
5231            .await
5232            .expect("the healthy fence must pass all four probe arms");
5233
5234        {
5235            let mut writer = backend.pool().writer().unwrap();
5236            writer
5237                .conn_mut()
5238                .execute_batch(
5239                    "DROP TRIGGER attachments_reject_claimed_blob_insert; \
5240                     DROP TRIGGER attachments_reject_claimed_blob_update; \
5241                     CREATE TRIGGER attachments_reject_claimed_blob_insert \
5242                     BEFORE INSERT ON attachments \
5243                     WHEN NEW.content_ref = \
5244                         '0000000000000000000000000000000000000000000000000000000000000000' \
5245                       AND EXISTS (SELECT 1 FROM blob_gc_claims \
5246                                   WHERE content_ref = NEW.content_ref) \
5247                     BEGIN \
5248                         SELECT RAISE(ABORT, \
5249                             'content_ref is reserved by an active blob sweep'); \
5250                     END; \
5251                     CREATE TRIGGER attachments_reject_claimed_blob_update \
5252                     BEFORE UPDATE OF content_ref ON attachments \
5253                     WHEN NEW.content_ref = \
5254                         '0000000000000000000000000000000000000000000000000000000000000000' \
5255                       AND EXISTS (SELECT 1 FROM blob_gc_claims \
5256                                   WHERE content_ref = NEW.content_ref) \
5257                     BEGIN \
5258                         SELECT RAISE(ABORT, \
5259                             'content_ref is reserved by an active blob sweep'); \
5260                     END;",
5261                )
5262                .unwrap();
5263        }
5264
5265        let error = super::blob_gc_fence_probe(sql.as_ref())
5266            .await
5267            .expect_err("a sentinel-only fence must fail the second-digest arms");
5268        assert!(
5269            matches!(error, StorageError::Unsupported { .. }),
5270            "expected StorageError::Unsupported, got {error:?}"
5271        );
5272        assert!(
5273            error.to_string().contains("second-digest"),
5274            "the refusal must name a second-digest arm, got {error}"
5275        );
5276    }
5277
5278    /// A trigger rewrite restricted to the entity/content shape passes the
5279    /// sentinel arms; the note-shaped second-digest arms must catch it.
5280    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5281    async fn fence_probe_refuses_a_shape_restricted_trigger_rewrite() {
5282        let dir = tempfile::tempdir().unwrap();
5283        let db_path = dir.path().join("khive.db");
5284        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5285        {
5286            let mut writer = backend.pool().writer().unwrap();
5287            prepare_completed_v21_gc_fixture(writer.conn_mut());
5288            writer
5289                .conn_mut()
5290                .execute_batch(
5291                    "DROP TRIGGER attachments_reject_claimed_blob_insert; \
5292                     DROP TRIGGER attachments_reject_claimed_blob_update; \
5293                     CREATE TRIGGER attachments_reject_claimed_blob_insert \
5294                     BEFORE INSERT ON attachments \
5295                     WHEN NEW.substrate = 'entity' \
5296                       AND EXISTS (SELECT 1 FROM blob_gc_claims \
5297                                   WHERE content_ref = NEW.content_ref) \
5298                     BEGIN \
5299                         SELECT RAISE(ABORT, \
5300                             'content_ref is reserved by an active blob sweep'); \
5301                     END; \
5302                     CREATE TRIGGER attachments_reject_claimed_blob_update \
5303                     BEFORE UPDATE OF content_ref ON attachments \
5304                     WHEN NEW.substrate = 'entity' \
5305                       AND EXISTS (SELECT 1 FROM blob_gc_claims \
5306                                   WHERE content_ref = NEW.content_ref) \
5307                     BEGIN \
5308                         SELECT RAISE(ABORT, \
5309                             'content_ref is reserved by an active blob sweep'); \
5310                     END;",
5311                )
5312                .unwrap();
5313        }
5314
5315        let sql = backend.sql();
5316        let error = super::blob_gc_fence_probe(sql.as_ref())
5317            .await
5318            .expect_err("an entity-shape-only fence must fail the note-shaped arms");
5319        assert!(
5320            matches!(error, StorageError::Unsupported { .. }),
5321            "expected StorageError::Unsupported, got {error:?}"
5322        );
5323        assert!(
5324            error.to_string().contains("second-digest"),
5325            "the refusal must name a second-digest arm, got {error}"
5326        );
5327    }
5328
5329    /// SQLite fixes the text encoding when the database file is first
5330    /// initialized; creating (and dropping) a table under the pragma leaves
5331    /// an initialized UTF-16LE database that later writers inherit.
5332    fn initialize_utf16le_database(db_path: &std::path::Path) {
5333        let conn = rusqlite::Connection::open(db_path).unwrap();
5334        conn.execute_batch(
5335            "PRAGMA encoding = 'UTF-16le'; \
5336             CREATE TABLE __encoding_pin (x INTEGER); \
5337             DROP TABLE __encoding_pin;",
5338        )
5339        .unwrap();
5340    }
5341
5342    /// CAST(TEXT AS BLOB) yields the database encoding's bytes, so a fixed
5343    /// bytes=64 arm would reject every valid 64-char ref in a UTF-16LE
5344    /// database (128 bytes). The width-derived arm must pass them.
5345    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5346    async fn blob_gc_evidence_accepts_valid_refs_in_a_utf16le_database() {
5347        let dir = tempfile::tempdir().unwrap();
5348        let db_path = dir.path().join("khive.db");
5349        initialize_utf16le_database(&db_path);
5350        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5351        {
5352            let mut writer = backend.pool().writer().unwrap();
5353            let encoding: String = writer
5354                .conn_mut()
5355                .query_row("PRAGMA encoding", [], |row| row.get(0))
5356                .unwrap();
5357            assert_eq!(
5358                encoding, "UTF-16le",
5359                "the fixture database must actually be UTF-16le"
5360            );
5361            prepare_completed_v21_gc_fixture(writer.conn_mut());
5362            // The fixture copies attachments from the (empty) entities table,
5363            // so seed one valid canonical row into each scanned table — with
5364            // no rows the probes measure nothing and a broken byte arm would
5365            // still pass this test.
5366            writer
5367                .conn_mut()
5368                .execute(
5369                    "INSERT INTO attachments \
5370                     (record_uuid, substrate, role, content_ref, created_at) \
5371                     VALUES ('utf16-valid-attachment', 'entity', 'content', ?1, 0)",
5372                    rusqlite::params!["a".repeat(64)],
5373                )
5374                .unwrap();
5375            writer
5376                .conn_mut()
5377                .execute(
5378                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5379                     VALUES ('utf16-valid-claim-key', ?1, 0)",
5380                    rusqlite::params!["b".repeat(64)],
5381                )
5382                .unwrap();
5383        }
5384
5385        let sql = backend.sql();
5386        super::validate_blob_gc_evidence(sql.as_ref())
5387            .await
5388            .expect("valid canonical refs must pass in a UTF-16LE database");
5389    }
5390
5391    /// The NUL arm must stay red in UTF-16 as well: 64 chars + NUL + tail is
5392    /// 130+ bytes against the expected 128.
5393    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5394    async fn blob_gc_evidence_rejects_a_nul_embedded_claim_ref_in_a_utf16le_database() {
5395        let dir = tempfile::tempdir().unwrap();
5396        let db_path = dir.path().join("khive.db");
5397        initialize_utf16le_database(&db_path);
5398        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5399        {
5400            let mut writer = backend.pool().writer().unwrap();
5401            let encoding: String = writer
5402                .conn_mut()
5403                .query_row("PRAGMA encoding", [], |row| row.get(0))
5404                .unwrap();
5405            assert_eq!(
5406                encoding, "UTF-16le",
5407                "the fixture database must actually be UTF-16le"
5408            );
5409            prepare_completed_v21_gc_fixture(writer.conn_mut());
5410            writer
5411                .conn_mut()
5412                .execute(
5413                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5414                     VALUES ('nul-claim-key-utf16', ?1, 0)",
5415                    rusqlite::params![nul_embedded_canonical_ref()],
5416                )
5417                .unwrap();
5418        }
5419
5420        let sql = backend.sql();
5421        let error = super::validate_blob_gc_evidence(sql.as_ref())
5422            .await
5423            .expect_err("a NUL-embedded claim ref must refuse the sweep in UTF-16LE too");
5424        assert!(
5425            error.to_string().contains("blob_gc_claims"),
5426            "expected the claims-table refusal, got {error:?}"
5427        );
5428    }
5429
5430    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5431    async fn transactional_orphan_sweep_preserves_put_started_after_liveness_mark() {
5432        let dir = tempfile::tempdir().unwrap();
5433        let db_path = dir.path().join("khive.db");
5434        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5435        {
5436            let mut writer = backend.pool().writer().unwrap();
5437            prepare_completed_v21_gc_fixture(writer.conn_mut());
5438        }
5439        let root = dir.path().join("blobs");
5440        let store = std::sync::Arc::new(
5441            FsBlobStore::new(root.clone(), 0)
5442                .unwrap()
5443                .with_orphan_sweep_grace(Duration::ZERO),
5444        );
5445        let orphan = store.put(b"old orphan".to_vec()).await.unwrap();
5446        let canonical_root = root.canonicalize().unwrap();
5447        let (marked, release, _done) = sync_hook::install(&canonical_root);
5448
5449        let sweep = {
5450            let store = store.clone();
5451            let sql = backend.sql();
5452            tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5453        };
5454        assert!(
5455            recv_blocking(marked).await,
5456            "sweep must finish its liveness mark"
5457        );
5458
5459        assert!(
5460            store.write_lock.try_lock().is_err(),
5461            "the sweep must hold the same root lock used by blob writers"
5462        );
5463        let (started_tx, started_rx) = std::sync::mpsc::channel();
5464        let new_ref = {
5465            let root = root.clone();
5466            tokio::task::spawn_blocking(move || {
5467                let _ = started_tx.send(());
5468                put_blocking(&root, 0, b"new concurrent blob".to_vec())
5469            })
5470        };
5471        assert!(recv_blocking(started_rx).await, "blob put must start");
5472
5473        release.send(()).unwrap();
5474        let sweep_result = sweep.await.unwrap().unwrap();
5475        let new_ref = new_ref.await.unwrap().unwrap();
5476
5477        assert_eq!(sweep_result.deleted, 1);
5478        assert!(!store.exists(&orphan).await.unwrap());
5479        assert!(
5480            store.exists(&new_ref).await.unwrap(),
5481            "a blob put started between the liveness mark and physical sweep must survive"
5482        );
5483    }
5484
5485    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5486    async fn transactional_orphan_sweep_releases_sqlite_writer_before_physical_delete() {
5487        let dir = tempfile::tempdir().unwrap();
5488        let db_path = dir.path().join("khive.db");
5489        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5490        {
5491            let mut writer = backend.pool().writer().unwrap();
5492            prepare_completed_v21_gc_fixture(writer.conn_mut());
5493        }
5494        let root = dir.path().join("blobs");
5495        let store = std::sync::Arc::new(
5496            FsBlobStore::new(root.clone(), 0)
5497                .unwrap()
5498                .with_orphan_sweep_grace(Duration::ZERO),
5499        );
5500        let orphan = store
5501            .put(b"claim then delete outside sqlite".to_vec())
5502            .await
5503            .unwrap();
5504        let canonical_root = root.canonicalize().unwrap();
5505        let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
5506
5507        let sweep = {
5508            let store = store.clone();
5509            let sql = backend.sql();
5510            tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5511        };
5512        assert!(
5513            recv_blocking(claimed).await,
5514            "sweep must durably claim the orphan before physical deletion"
5515        );
5516        assert!(
5517            store.exists(&orphan).await.unwrap(),
5518            "the test seam must pause before the physical delete"
5519        );
5520        let external_database_lock = fs::OpenOptions::new()
5521            .read(true)
5522            .write(true)
5523            .open(database_gc_lock_path(&db_path))
5524            .unwrap();
5525        assert!(
5526            matches!(
5527                fs4::FileExt::try_lock(&external_database_lock),
5528                Err(fs4::TryLockError::WouldBlock)
5529            ),
5530            "the sweep must retain cross-process database ownership while SQLite's writer is free"
5531        );
5532
5533        // The destructive filesystem phase is deliberately paused. An
5534        // unrelated SQLite writer must nevertheless complete now: this is
5535        // the hold-time proof that external I/O is no longer inside the
5536        // sweep's BEGIN IMMEDIATE span.
5537        let unrelated = rusqlite::Connection::open(&db_path).unwrap();
5538        unrelated.busy_timeout(Duration::from_millis(100)).unwrap();
5539        unrelated
5540            .execute(
5541                "INSERT INTO entities \
5542                 (id, namespace, kind, name, tags, created_at, updated_at) \
5543                 VALUES ('unrelated-writer', 'local', 'concept', 'unrelated', '[]', 1, 1)",
5544                [],
5545            )
5546            .expect("external filesystem work must not retain SQLite's writer lock");
5547
5548        // The claim trigger is the cross-resource fence: while the file is
5549        // selected for deletion, a concurrent attachment writer cannot make it
5550        // newly live in the released-writer window.
5551        let claimed_err = unrelated
5552            .execute(
5553                "INSERT INTO attachments \
5554                 (record_uuid, substrate, role, content_ref, created_at) \
5555                 VALUES ('racing-reference', 'entity', 'content', ?1, 1)",
5556                [orphan.as_str()],
5557            )
5558            .expect_err("a claimed content_ref must fail closed before deletion");
5559        assert!(
5560            claimed_err.to_string().contains("active blob sweep"),
5561            "unexpected claim error: {claimed_err}"
5562        );
5563
5564        release_delete.send(()).unwrap();
5565        let result = sweep.await.unwrap().unwrap();
5566        assert_eq!(result.deleted, 1);
5567        assert!(!store.exists(&orphan).await.unwrap());
5568
5569        let remaining_claims: i64 = unrelated
5570            .query_row(
5571                "SELECT COUNT(*) FROM blob_gc_claims WHERE content_ref = ?1",
5572                [orphan.as_str()],
5573                |row| row.get(0),
5574            )
5575            .unwrap();
5576        assert_eq!(
5577            remaining_claims, 0,
5578            "successful deletion releases the claim"
5579        );
5580    }
5581
5582    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5583    async fn cancelling_sweep_during_delete_keeps_owner_locks_until_blocking_work_finishes() {
5584        let dir = tempfile::tempdir().unwrap();
5585        let db_path = dir.path().join("khive.db");
5586        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5587        {
5588            let mut writer = backend.pool().writer().unwrap();
5589            prepare_completed_v21_gc_fixture(writer.conn_mut());
5590        }
5591        let root = dir.path().join("blobs");
5592        let store = std::sync::Arc::new(
5593            FsBlobStore::new(root.clone(), 0)
5594                .unwrap()
5595                .with_orphan_sweep_grace(Duration::ZERO),
5596        );
5597        let orphan = store.put(b"cancelled sweep orphan".to_vec()).await.unwrap();
5598        let canonical_root = root.canonicalize().unwrap();
5599        let (claimed, release_delete, done) = sync_hook::install(&canonical_root);
5600
5601        let sweep = {
5602            let store = store.clone();
5603            let sql = backend.sql();
5604            tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5605        };
5606        assert!(recv_blocking(claimed).await);
5607        sweep.abort();
5608        assert!(sweep.await.unwrap_err().is_cancelled());
5609
5610        let external_root_lock = fs::OpenOptions::new()
5611            .read(true)
5612            .write(true)
5613            .open(root.join(ROOT_WRITE_LOCK_FILE))
5614            .unwrap();
5615        let external_database_lock = fs::OpenOptions::new()
5616            .read(true)
5617            .write(true)
5618            .open(database_gc_lock_path(&db_path))
5619            .unwrap();
5620        assert!(matches!(
5621            fs4::FileExt::try_lock(&external_root_lock),
5622            Err(fs4::TryLockError::WouldBlock)
5623        ));
5624        assert!(matches!(
5625            fs4::FileExt::try_lock(&external_database_lock),
5626            Err(fs4::TryLockError::WouldBlock)
5627        ));
5628
5629        release_delete.send(()).unwrap();
5630        let done_disconnected = tokio::task::spawn_blocking(move || done.recv().is_err())
5631            .await
5632            .unwrap();
5633        assert!(
5634            done_disconnected,
5635            "the cancelled outer task cannot send done"
5636        );
5637        assert!(fs4::FileExt::try_lock(&external_root_lock).is_ok());
5638        assert!(fs4::FileExt::try_lock(&external_database_lock).is_ok());
5639        drop(external_root_lock);
5640        drop(external_database_lock);
5641        assert!(!store.exists(&orphan).await.unwrap());
5642        let stranded_claims: i64 = rusqlite::Connection::open(&db_path)
5643            .unwrap()
5644            .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5645            .unwrap();
5646        assert_eq!(
5647            stranded_claims, 1,
5648            "cancellation leaves a fail-closed claim"
5649        );
5650
5651        let recovered = store
5652            .transactional_orphan_sweep(backend.sql().as_ref(), false)
5653            .await
5654            .unwrap();
5655        assert_eq!(recovered.deleted, 0);
5656        let remaining: i64 = rusqlite::Connection::open(&db_path)
5657            .unwrap()
5658            .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5659            .unwrap();
5660        assert_eq!(remaining, 0, "the next exclusive owner recovers the claim");
5661    }
5662
5663    #[tokio::test]
5664    async fn transactional_orphan_sweep_recovers_stale_claims_fail_closed() {
5665        let dir = tempfile::tempdir().unwrap();
5666        let db_path = dir.path().join("khive.db");
5667        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5668        {
5669            let mut writer = backend.pool().writer().unwrap();
5670            prepare_completed_v21_gc_fixture(writer.conn_mut());
5671        }
5672        let root = dir.path().join("blobs");
5673        let store = FsBlobStore::new(root.clone(), 0)
5674            .unwrap()
5675            .with_orphan_sweep_grace(Duration::from_secs(60));
5676        let bytes = b"republished after a crashed claim".to_vec();
5677        let content_ref = store.put(bytes.clone()).await.unwrap();
5678        let canonical_root = root.canonicalize().unwrap();
5679        let root_key = blob_root_key(&canonical_root);
5680        let absent_ref = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee";
5681        let former_probe_seed = "1111111111111111111111111111111111111111111111111111111111111111";
5682        {
5683            let writer = backend.pool().writer().unwrap();
5684            writer
5685                .conn()
5686                .execute(
5687                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5688                     VALUES (?1, ?2, 1), (?1, ?3, 1), (?1, ?4, 1)",
5689                    rusqlite::params![
5690                        root_key,
5691                        content_ref.as_str(),
5692                        absent_ref,
5693                        former_probe_seed
5694                    ],
5695                )
5696                .unwrap();
5697        }
5698
5699        // A publisher that recovered after the claiming process crashed
5700        // refreshes the digest's grace witness before its attachment write. The
5701        // next sweep must clear this protected claim, the claim whose file was
5702        // already removed, and the former fixed probe seed without letting a
5703        // healthy attachment fence block abandoned-claim recovery forever.
5704        assert_eq!(store.put(bytes).await.unwrap(), content_ref);
5705        let result = store
5706            .transactional_orphan_sweep(backend.sql().as_ref(), false)
5707            .await
5708            .unwrap();
5709        assert_eq!(result.deleted, 0);
5710        assert_eq!(result.grace_period_skipped, 1);
5711        assert!(store.exists(&content_ref).await.unwrap());
5712
5713        let remaining: i64 = backend
5714            .pool()
5715            .writer()
5716            .unwrap()
5717            .conn()
5718            .query_row(
5719                "SELECT COUNT(*) FROM blob_gc_claims WHERE root_key = ?1",
5720                [blob_root_key(&canonical_root)],
5721                |row| row.get(0),
5722            )
5723            .unwrap();
5724        assert_eq!(remaining, 0, "the next sweep recovers stale claims");
5725    }
5726
5727    #[tokio::test]
5728    async fn transactional_orphan_sweep_recovers_claims_after_root_relocation() {
5729        let dir = tempfile::tempdir().unwrap();
5730        let db_path = dir.path().join("khive.db");
5731        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5732        {
5733            let mut writer = backend.pool().writer().unwrap();
5734            prepare_completed_v21_gc_fixture(writer.conn_mut());
5735        }
5736
5737        let old_root = dir.path().join("old-blobs");
5738        let bytes = b"claim must follow a relocated blob root".to_vec();
5739        let content_ref = {
5740            let old_store = FsBlobStore::new(old_root.clone(), 0)
5741                .unwrap()
5742                .with_orphan_sweep_grace(Duration::from_secs(60));
5743            old_store.put(bytes).await.unwrap()
5744        };
5745        let old_root_key = blob_root_key(&old_root.canonicalize().unwrap());
5746        backend
5747            .pool()
5748            .writer()
5749            .unwrap()
5750            .conn()
5751            .execute(
5752                "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5753                 VALUES (?1, ?2, 1)",
5754                rusqlite::params![old_root_key, content_ref.as_str()],
5755            )
5756            .unwrap();
5757
5758        let new_root = dir.path().join("relocated-blobs");
5759        std::fs::rename(&old_root, &new_root).unwrap();
5760        let relocated_store = FsBlobStore::new(new_root, 0)
5761            .unwrap()
5762            .with_orphan_sweep_grace(Duration::from_secs(60));
5763        let result = relocated_store
5764            .transactional_orphan_sweep(backend.sql().as_ref(), false)
5765            .await
5766            .unwrap();
5767
5768        assert_eq!(
5769            result.deleted, 0,
5770            "a fresh relocated blob remains protected"
5771        );
5772        assert_eq!(result.grace_period_skipped, 1);
5773        assert!(relocated_store.exists(&content_ref).await.unwrap());
5774        let remaining: i64 = backend
5775            .pool()
5776            .writer()
5777            .unwrap()
5778            .conn()
5779            .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5780            .unwrap();
5781        assert_eq!(
5782            remaining, 0,
5783            "exclusive database sweep ownership makes every pre-existing claim abandoned, \
5784             even when its old path-derived root key no longer matches"
5785        );
5786    }
5787
5788    #[tokio::test]
5789    async fn transactional_orphan_sweep_recovers_claims_copied_by_database_restore() {
5790        let dir = tempfile::tempdir().unwrap();
5791        let source_path = dir.path().join("source.db");
5792        let restored_path = dir.path().join("restored.db");
5793        let bytes = b"claim copied in an online database backup".to_vec();
5794        let content_ref = ContentRef::from_digest_bytes(blake3::hash(&bytes).as_bytes());
5795        {
5796            let source = crate::StorageBackend::sqlite(&source_path).unwrap();
5797            let mut writer = source.pool().writer().unwrap();
5798            prepare_completed_v21_gc_fixture(writer.conn_mut());
5799            writer
5800                .conn()
5801                .execute(
5802                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5803                     VALUES ('source-root-before-backup', ?1, 1)",
5804                    [content_ref.as_str()],
5805                )
5806                .unwrap();
5807            writer
5808                .conn()
5809                .execute_batch("PRAGMA wal_checkpoint(TRUNCATE)")
5810                .unwrap();
5811        }
5812        std::fs::copy(&source_path, &restored_path).unwrap();
5813
5814        let restored = crate::StorageBackend::sqlite(&restored_path).unwrap();
5815        let restored_root = dir.path().join("restored-blobs");
5816        let store = FsBlobStore::new(restored_root, 0)
5817            .unwrap()
5818            .with_orphan_sweep_grace(Duration::from_secs(60));
5819        assert_eq!(store.put(bytes).await.unwrap(), content_ref);
5820        let result = store
5821            .transactional_orphan_sweep(restored.sql().as_ref(), false)
5822            .await
5823            .unwrap();
5824
5825        assert_eq!(result.deleted, 0);
5826        assert_eq!(result.grace_period_skipped, 1);
5827        let remaining: i64 = restored
5828            .pool()
5829            .writer()
5830            .unwrap()
5831            .conn()
5832            .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5833            .unwrap();
5834        assert_eq!(remaining, 0, "restored claims are abandoned ownership");
5835    }
5836
5837    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5838    async fn transactional_orphan_sweep_bounds_each_durable_claim_batch() {
5839        let dir = tempfile::tempdir().unwrap();
5840        let db_path = dir.path().join("khive.db");
5841        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
5842        {
5843            let mut writer = backend.pool().writer().unwrap();
5844            prepare_completed_v21_gc_fixture(writer.conn_mut());
5845        }
5846        let root = dir.path().join("blobs");
5847        let store = std::sync::Arc::new(
5848            FsBlobStore::new(root.clone(), 0)
5849                .unwrap()
5850                .with_orphan_sweep_grace(Duration::ZERO),
5851        );
5852        let candidate_count = BLOB_GC_CLAIM_BATCH_SIZE * 2 + 1;
5853        for index in 0..candidate_count {
5854            store
5855                .put(format!("bounded claim candidate {index}").into_bytes())
5856                .await
5857                .unwrap();
5858        }
5859        let canonical_root = root.canonicalize().unwrap();
5860        let (claimed, release_delete, _done) = sync_hook::install(&canonical_root);
5861
5862        let sweep = {
5863            let store = store.clone();
5864            let sql = backend.sql();
5865            tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
5866        };
5867        assert!(
5868            recv_blocking(claimed).await,
5869            "the first bounded claim batch must commit before deletion"
5870        );
5871        let active_claims: i64 = rusqlite::Connection::open(&db_path)
5872            .unwrap()
5873            .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5874            .unwrap();
5875        assert!(active_claims > 0);
5876        assert!(
5877            active_claims <= BLOB_GC_CLAIM_BATCH_SIZE as i64,
5878            "one transaction may expose at most {BLOB_GC_CLAIM_BATCH_SIZE} claim rows; \
5879             observed {active_claims}"
5880        );
5881
5882        release_delete.send(()).unwrap();
5883        let result = sweep.await.unwrap().unwrap();
5884        assert_eq!(result.deleted, candidate_count as u64);
5885        let remaining: i64 = rusqlite::Connection::open(&db_path)
5886            .unwrap()
5887            .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5888            .unwrap();
5889        assert_eq!(remaining, 0);
5890    }
5891
5892    #[tokio::test]
5893    async fn abandoned_claim_recovery_deletes_at_most_one_batch_per_writer_hold() {
5894        let dir = tempfile::tempdir().unwrap();
5895        let db_path = dir.path().join("khive.db");
5896        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5897        {
5898            let mut writer = backend.pool().writer().unwrap();
5899            crate::run_migrations(writer.conn_mut()).unwrap();
5900            let tx = writer.conn_mut().transaction().unwrap();
5901            for index in 0..(BLOB_GC_CLAIM_BATCH_SIZE + 1) {
5902                tx.execute(
5903                    "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5904                     VALUES ('abandoned-root', ?1, 1)",
5905                    [format!("{index:064x}")],
5906                )
5907                .unwrap();
5908            }
5909            tx.commit().unwrap();
5910        }
5911
5912        let released = release_abandoned_blob_gc_claim_batch(backend.sql().as_ref())
5913            .await
5914            .unwrap();
5915        assert_eq!(released, BLOB_GC_CLAIM_BATCH_SIZE as u64);
5916        let remaining: i64 = backend
5917            .pool()
5918            .writer()
5919            .unwrap()
5920            .conn()
5921            .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
5922            .unwrap();
5923        assert_eq!(remaining, 1);
5924    }
5925
5926    #[tokio::test]
5927    async fn transactional_orphan_sweep_refuses_corrupt_liveness_and_claim_evidence() {
5928        let dir = tempfile::tempdir().unwrap();
5929        let db_path = dir.path().join("khive.db");
5930        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
5931        {
5932            let mut writer = backend.pool().writer().unwrap();
5933            prepare_completed_v21_gc_fixture(writer.conn_mut());
5934        }
5935        let root = dir.path().join("blobs");
5936        let store = FsBlobStore::new(root.clone(), 0)
5937            .unwrap()
5938            .with_orphan_sweep_grace(Duration::ZERO);
5939        let orphan = store
5940            .put(b"must survive corrupt evidence".to_vec())
5941            .await
5942            .unwrap();
5943
5944        let conn = rusqlite::Connection::open(&db_path).unwrap();
5945        conn.execute_batch("PRAGMA ignore_check_constraints = ON")
5946            .unwrap();
5947        conn.execute(
5948            "INSERT INTO attachments \
5949             (record_uuid, substrate, role, content_ref, created_at) \
5950             VALUES ('corrupt-live', 'entity', 'content', 'not-a-content-ref', 1)",
5951            [],
5952        )
5953        .unwrap();
5954        conn.execute_batch("PRAGMA ignore_check_constraints = OFF")
5955            .unwrap();
5956        let live_error = store
5957            .transactional_orphan_sweep(backend.sql().as_ref(), false)
5958            .await
5959            .expect_err("corrupt live evidence must fail closed");
5960        assert!(matches!(live_error, StorageError::InvalidInput { .. }));
5961        assert!(
5962            store.exists(&orphan).await.unwrap(),
5963            "no file may be removed after corrupt live evidence"
5964        );
5965
5966        conn.execute(
5967            "DELETE FROM attachments WHERE record_uuid = 'corrupt-live'",
5968            [],
5969        )
5970        .unwrap();
5971        let root_key = blob_root_key(&root.canonicalize().unwrap());
5972        conn.execute(
5973            "INSERT INTO blob_gc_claims (root_key, content_ref, claimed_at) \
5974             VALUES (?1, 'also-not-a-content-ref', 1)",
5975            [root_key.as_str()],
5976        )
5977        .unwrap();
5978        let claim_error = store
5979            .transactional_orphan_sweep(backend.sql().as_ref(), false)
5980            .await
5981            .expect_err("corrupt durable claim evidence must fail closed");
5982        assert!(matches!(claim_error, StorageError::InvalidInput { .. }));
5983        assert!(
5984            store.exists(&orphan).await.unwrap(),
5985            "no file may be removed after corrupt claim evidence"
5986        );
5987        let remaining: i64 = conn
5988            .query_row(
5989                "SELECT COUNT(*) FROM blob_gc_claims \
5990                 WHERE root_key = ?1 AND content_ref = 'also-not-a-content-ref'",
5991                [root_key.as_str()],
5992                |row| row.get(0),
5993            )
5994            .unwrap();
5995        assert_eq!(
5996            remaining, 1,
5997            "corrupt claim evidence is not silently erased"
5998        );
5999        let probe_residue: i64 = conn
6000            .query_row(
6001                "SELECT (SELECT COUNT(*) FROM blob_gc_claims \
6002                         WHERE root_key GLOB '__fence_probe-*') \
6003                      + (SELECT COUNT(*) FROM attachments \
6004                         WHERE record_uuid GLOB '__blob-gc-fence-probe-*')",
6005                [],
6006                |row| row.get(0),
6007            )
6008            .unwrap();
6009        assert_eq!(
6010            probe_residue, 0,
6011            "invalid evidence must abort before the functional fence probe"
6012        );
6013    }
6014
6015    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6016    async fn transactional_orphan_sweep_republishes_deduplicated_external_put() {
6017        let dir = tempfile::tempdir().unwrap();
6018        let db_path = dir.path().join("khive.db");
6019        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
6020        {
6021            let mut writer = backend.pool().writer().unwrap();
6022            prepare_completed_v21_gc_fixture(writer.conn_mut());
6023        }
6024        let root = dir.path().join("blobs");
6025        let store = std::sync::Arc::new(
6026            FsBlobStore::new(root.clone(), 0)
6027                .unwrap()
6028                .with_orphan_sweep_grace(Duration::ZERO),
6029        );
6030        let payload = b"existing orphan republished during sweep".to_vec();
6031        let orphan = store.put(payload.clone()).await.unwrap();
6032        let canonical_root = root.canonicalize().unwrap();
6033        let (marked, release, _done) = sync_hook::install(&canonical_root);
6034
6035        let sweep = {
6036            let store = store.clone();
6037            let sql = backend.sql();
6038            tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), false).await })
6039        };
6040        assert!(
6041            recv_blocking(marked).await,
6042            "sweep must finish its liveness mark"
6043        );
6044
6045        let external_lock = fs::OpenOptions::new()
6046            .read(true)
6047            .write(true)
6048            .open(root.join(ROOT_WRITE_LOCK_FILE))
6049            .unwrap();
6050        assert!(
6051            matches!(
6052                fs4::FileExt::try_lock(&external_lock),
6053                Err(fs4::TryLockError::WouldBlock)
6054            ),
6055            "the sweep must exclude a publisher using an independently opened root lock"
6056        );
6057
6058        let (started_tx, started_rx) = std::sync::mpsc::channel();
6059        let republished = {
6060            let root = root.clone();
6061            tokio::task::spawn_blocking(move || {
6062                let _ = started_tx.send(());
6063                put_blocking(&root, 0, payload)
6064            })
6065        };
6066        assert!(recv_blocking(started_rx).await, "blob put must start");
6067
6068        release.send(()).unwrap();
6069        let sweep_result = sweep.await.unwrap().unwrap();
6070        let republished = republished.await.unwrap().unwrap();
6071
6072        assert_eq!(sweep_result.deleted, 1);
6073        assert_eq!(republished, orphan);
6074        assert!(
6075            store.exists(&republished).await.unwrap(),
6076            "a deduplicated put concurrent with the sweep must not return a deleted reference"
6077        );
6078    }
6079
6080    #[tokio::test]
6081    async fn transactional_orphan_sweep_uses_all_attachment_refs_as_live() {
6082        let dir = tempfile::tempdir().unwrap();
6083        let db_path = dir.path().join("khive.db");
6084        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6085        {
6086            let mut writer = backend.pool().writer().unwrap();
6087            prepare_completed_v21_gc_fixture(writer.conn_mut());
6088        }
6089        let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6090            .unwrap()
6091            .with_orphan_sweep_grace(Duration::ZERO);
6092        let live = store.put(b"live".to_vec()).await.unwrap();
6093        let soft_deleted = store.put(b"soft deleted".to_vec()).await.unwrap();
6094        let orphan = store.put(b"orphan".to_vec()).await.unwrap();
6095        {
6096            let writer = backend.pool().writer().unwrap();
6097            writer
6098                .conn()
6099                .execute_batch(
6100                    "INSERT INTO entities \
6101                     (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6102                     VALUES ('live', 'local', 'document', 'live', '[]', 1, 1, NULL), \
6103                            ('deleted', 'local', 'document', 'deleted', '[]', 1, 1, 2);",
6104                )
6105                .unwrap();
6106            writer
6107                .conn()
6108                .execute(
6109                    "INSERT INTO attachments \
6110                     (record_uuid, substrate, role, content_ref, created_at) \
6111                     VALUES ('live', 'entity', 'content', ?1, 1), \
6112                            ('deleted', 'entity', 'content', ?2, 1)",
6113                    rusqlite::params![live.as_str(), soft_deleted.as_str()],
6114                )
6115                .unwrap();
6116        }
6117
6118        let dry_run = store
6119            .transactional_orphan_sweep(backend.sql().as_ref(), true)
6120            .await
6121            .unwrap();
6122        assert_eq!(dry_run.would_delete, 1);
6123        assert_eq!(dry_run.deleted, 0);
6124        assert!(store.exists(&soft_deleted).await.unwrap());
6125        assert!(store.exists(&orphan).await.unwrap());
6126
6127        let result = store
6128            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6129            .await
6130            .unwrap();
6131
6132        assert_eq!(result.scanned, 3);
6133        assert_eq!(result.deleted, 1);
6134        assert!(store.exists(&live).await.unwrap());
6135        assert!(
6136            store.exists(&soft_deleted).await.unwrap(),
6137            "soft delete retains attachment rows and their blobs"
6138        );
6139        assert!(!store.exists(&orphan).await.unwrap());
6140    }
6141
6142    #[tokio::test]
6143    async fn transactional_orphan_sweep_protects_a_freshly_published_blob_before_its_reference_commits(
6144    ) {
6145        // The exact two-step client protocol hazard: `put` completes and
6146        // releases its write lock (step 1) while the attachment write that will
6147        // *later* commit a `content_ref` to this blob (step 2) has not
6148        // happened yet -- nothing in this store's locking serializes the
6149        // two, because they are separate calls the client makes with an
6150        // arbitrary gap in between. A sweep that lands in that gap must not
6151        // delete the blob: `attachments.content_ref` has no row for it yet
6152        // purely because the referencing write hasn't landed, not because
6153        // it is actually orphaned. Without the publish-grace window this
6154        // reproduces khive#1313's dangling-reference defect: the blob file
6155        // is deleted here, and the still-pending attachment write below would
6156        // commit a `content_ref` to nothing.
6157        let dir = tempfile::tempdir().unwrap();
6158        let db_path = dir.path().join("khive.db");
6159        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6160        {
6161            let mut writer = backend.pool().writer().unwrap();
6162            prepare_completed_v21_gc_fixture(writer.conn_mut());
6163        }
6164        // Default (non-zero) grace period -- this test exercises exactly
6165        // what it exists to protect.
6166        let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
6167
6168        // Step 1: put completes, lock released. No attachment anywhere
6169        // references this blob yet.
6170        let blob = store
6171            .put(b"published, reference not yet committed".to_vec())
6172            .await
6173            .unwrap();
6174
6175        // A sweep runs in the gap before step 2 (the attachment write) happens.
6176        let result = store
6177            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6178            .await
6179            .unwrap();
6180
6181        assert_eq!(result.deleted, 0, "the blob must survive: {result:?}");
6182        assert_eq!(
6183            result.would_delete, 0,
6184            "not treated as a deletable orphan: {result:?}"
6185        );
6186        assert_eq!(
6187            result.grace_period_skipped, 1,
6188            "must be reported as grace-protected rather than silently ignored: {result:?}"
6189        );
6190        assert!(
6191            store.exists(&blob).await.unwrap(),
6192            "a blob still inside its publish grace period must survive the sweep"
6193        );
6194
6195        // Step 2 now lands: the record-plus-attachment write commits content_ref to the
6196        // still-present blob.
6197        {
6198            let writer = backend.pool().writer().unwrap();
6199            writer
6200                .conn()
6201                .execute(
6202                    "INSERT INTO entities \
6203                     (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6204                     VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
6205                    [],
6206                )
6207                .unwrap();
6208            writer
6209                .conn()
6210                .execute(
6211                    "INSERT INTO attachments \
6212                     (record_uuid, substrate, role, content_ref, created_at) \
6213                     VALUES ('e1', 'entity', 'content', ?1, 1)",
6214                    [blob.as_str()],
6215                )
6216                .unwrap();
6217        }
6218
6219        // A later sweep now finds it live and keeps it for the ordinary
6220        // reason, independent of the grace window.
6221        let result = store
6222            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6223            .await
6224            .unwrap();
6225        assert_eq!(result.deleted, 0);
6226        assert!(store.exists(&blob).await.unwrap());
6227    }
6228
6229    #[tokio::test]
6230    async fn put_republishing_an_aged_orphan_restarts_its_grace_clock_before_the_reference_commits()
6231    {
6232        // The dedup fast path (`target.exists()`) used to return without
6233        // touching the file at all -- so a stale, already-orphaned blob
6234        // re-published by an identical `put` kept its OLD mtime, bypassed
6235        // the publish-grace check, and a transactional sweep landing in the
6236        // gap before the caller's follow-up attachment write could delete it
6237        // out from under that write (khive#1313). This reproduces the
6238        // race end to end and proves the mtime refresh closes it.
6239        let dir = tempfile::tempdir().unwrap();
6240        let db_path = dir.path().join("khive.db");
6241        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6242        {
6243            let mut writer = backend.pool().writer().unwrap();
6244            prepare_completed_v21_gc_fixture(writer.conn_mut());
6245        }
6246        let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6247            .unwrap()
6248            .with_orphan_sweep_grace(Duration::from_secs(60));
6249
6250        let bytes = b"old orphan re-published".to_vec();
6251        let first = store.put(bytes.clone()).await.unwrap();
6252
6253        // Age the blob well past the 60s grace floor -- no sleeps, same
6254        // backdating pattern as the existing older-than-grace test.
6255        let path = shard_path(store.root(), &first);
6256        let old_mtime = SystemTime::now() - Duration::from_secs(3600);
6257        fs::OpenOptions::new()
6258            .write(true)
6259            .open(&path)
6260            .unwrap()
6261            .set_modified(old_mtime)
6262            .unwrap();
6263
6264        // A deduplicating put republishes the identical bytes. No attachment
6265        // anywhere references this blob yet.
6266        let second = store.put(bytes).await.unwrap();
6267        assert_eq!(first, second);
6268
6269        // The sweep lands in the gap before the follow-up attachment write --
6270        // the refreshed mtime must keep it inside the grace window.
6271        let result = store
6272            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6273            .await
6274            .unwrap();
6275        assert_eq!(
6276            result.deleted, 0,
6277            "a dedup-republished blob must survive a sweep landing before its reference \
6278             commits: {result:?}"
6279        );
6280        assert_eq!(
6281            result.grace_period_skipped, 1,
6282            "must be reported as grace-protected, not silently ignored: {result:?}"
6283        );
6284        assert!(store.exists(&first).await.unwrap());
6285
6286        // The caller's follow-up record-plus-attachment write now lands.
6287        {
6288            let writer = backend.pool().writer().unwrap();
6289            writer
6290                .conn()
6291                .execute(
6292                    "INSERT INTO entities \
6293                     (id, namespace, kind, name, tags, created_at, updated_at, deleted_at) \
6294                     VALUES ('e1', 'local', 'document', 'e1', '[]', 1, 1, NULL)",
6295                    [],
6296                )
6297                .unwrap();
6298            writer
6299                .conn()
6300                .execute(
6301                    "INSERT INTO attachments \
6302                     (record_uuid, substrate, role, content_ref, created_at) \
6303                     VALUES ('e1', 'entity', 'content', ?1, 1)",
6304                    [first.as_str()],
6305                )
6306                .unwrap();
6307        }
6308
6309        let result = store
6310            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6311            .await
6312            .unwrap();
6313        assert_eq!(result.deleted, 0);
6314        assert!(
6315            store.exists(&first).await.unwrap(),
6316            "the blob must stay live once its reference has committed"
6317        );
6318    }
6319
6320    #[tokio::test]
6321    async fn put_dedup_mtime_refresh_has_no_observable_effect_under_zero_grace_period() {
6322        // The assumption the fix relies on for every zero-grace test in this
6323        // file: `within_publish_grace` with `Duration::ZERO` never protects a
6324        // candidate regardless of its mtime (`age < Duration::ZERO` is always
6325        // false), so refreshing the mtime on a deduplicated republish must
6326        // not change zero-grace sweep behavior. Verified directly rather
6327        // than assumed.
6328        let dir = tempfile::tempdir().unwrap();
6329        let db_path = dir.path().join("khive.db");
6330        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6331        {
6332            let mut writer = backend.pool().writer().unwrap();
6333            prepare_completed_v21_gc_fixture(writer.conn_mut());
6334        }
6335        let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6336            .unwrap()
6337            .with_orphan_sweep_grace(Duration::ZERO);
6338        let bytes = b"zero grace dedup refresh".to_vec();
6339        let first = store.put(bytes.clone()).await.unwrap();
6340        let second = store.put(bytes.clone()).await.unwrap();
6341        assert_eq!(first, second);
6342
6343        let result = store
6344            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6345            .await
6346            .unwrap();
6347        assert_eq!(
6348            result.deleted, 1,
6349            "a zero grace period must still delete an unreferenced blob even after a dedup \
6350             put refreshed its mtime: {result:?}"
6351        );
6352        assert!(!store.exists(&first).await.unwrap());
6353    }
6354
6355    #[tokio::test]
6356    async fn transactional_orphan_sweep_still_removes_orphans_older_than_the_grace_period() {
6357        // The grace window narrows the publish-vs-sweep race, it does not
6358        // disable sweeping outright: an object whose age already exceeds a
6359        // (short, for this test) grace period is removed exactly as before,
6360        // proving the fix bounds the exposure rather than papering over it.
6361        let dir = tempfile::tempdir().unwrap();
6362        let db_path = dir.path().join("khive.db");
6363        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6364        {
6365            let mut writer = backend.pool().writer().unwrap();
6366            prepare_completed_v21_gc_fixture(writer.conn_mut());
6367        }
6368        let store = FsBlobStore::new(dir.path().join("blobs"), 0)
6369            .unwrap()
6370            .with_orphan_sweep_grace(Duration::from_secs(60));
6371
6372        let orphan = store
6373            .put(b"actually orphaned, published long ago".to_vec())
6374            .await
6375            .unwrap();
6376        // Back-date the file's mtime well past the 60s grace period instead
6377        // of sleeping in the test.
6378        let path = shard_path(store.root(), &orphan);
6379        let old_mtime = SystemTime::now() - Duration::from_secs(3600);
6380        fs::OpenOptions::new()
6381            .write(true)
6382            .open(&path)
6383            .unwrap()
6384            .set_modified(old_mtime)
6385            .unwrap();
6386
6387        let result = store
6388            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6389            .await
6390            .unwrap();
6391
6392        assert_eq!(
6393            result.deleted, 1,
6394            "an orphan older than the grace period must still be swept: {result:?}"
6395        );
6396        assert_eq!(result.grace_period_skipped, 0);
6397        assert!(!store.exists(&orphan).await.unwrap());
6398    }
6399
6400    #[test]
6401    fn resolve_blob_root_prefers_env_var() {
6402        let _guard = ENV_LOCK.lock().unwrap();
6403        std::env::set_var("KHIVE_BLOB_ROOT", "/tmp/env-override-root");
6404        let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
6405        std::env::remove_var("KHIVE_BLOB_ROOT");
6406        assert_eq!(resolved.unwrap(), PathBuf::from("/tmp/env-override-root"));
6407    }
6408
6409    #[test]
6410    fn resolve_blob_root_prefers_config_over_default() {
6411        let _guard = ENV_LOCK.lock().unwrap();
6412        std::env::remove_var("KHIVE_BLOB_ROOT");
6413        let resolved = resolve_blob_root(Some(Path::new("/db/dir")), Some(Path::new("/cfg/root")));
6414        assert_eq!(resolved.unwrap(), PathBuf::from("/cfg/root"));
6415    }
6416
6417    #[test]
6418    fn resolve_blob_root_defaults_beside_db_dir() {
6419        let _guard = ENV_LOCK.lock().unwrap();
6420        std::env::remove_var("KHIVE_BLOB_ROOT");
6421        let resolved = resolve_blob_root(Some(Path::new("/db/dir")), None);
6422        assert_eq!(resolved.unwrap(), PathBuf::from("/db/dir/blobs"));
6423    }
6424
6425    #[test]
6426    fn resolve_blob_root_errors_with_no_env_config_or_db_dir() {
6427        let _guard = ENV_LOCK.lock().unwrap();
6428        std::env::remove_var("KHIVE_BLOB_ROOT");
6429        let resolved = resolve_blob_root(None, None);
6430        assert!(resolved.is_err());
6431    }
6432
6433    // `std::env::set_var`/`remove_var` mutate real process-global state, so the
6434    // four `resolve_blob_root` env-precedence tests must not interleave under
6435    // the crate's default parallel test runner.
6436    static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
6437
6438    #[cfg(unix)]
6439    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6440    async fn transactional_orphan_sweep_walk_ignores_a_leaf_swapped_for_an_outside_symlink_mid_scan(
6441    ) {
6442        // Regression for the PR #2201 review finding: the sweep's candidate
6443        // walk and grace-period mtime read used to be two separate,
6444        // path-based passes (`walk_blob_files` then `within_publish_grace`
6445        // via `fs::metadata(path)`), neither pinned to the retained
6446        // `root_handle`. A concurrent replacement of a leaf entry in the gap
6447        // between those passes could make the mtime read observe a file
6448        // OUTSIDE the retained root, letting a stale outside mtime evict the
6449        // grace period for a freshly published in-root blob. The fix folds
6450        // discovery and classification into one handle-relative
6451        // openat(..., O_NOFOLLOW) + fstat in
6452        // `walk_blob_files_from_root_handle`, so there is no later path
6453        // re-resolution left to race. This test forces exactly that swap —
6454        // via `walk_leaf_sync_hook`, between the leaf's hex-name discovery
6455        // and its classifying open — and proves the outside decoy's stale
6456        // mtime never reaches the sweep's counts.
6457        let dir = tempfile::tempdir().unwrap();
6458        let db_path = dir.path().join("khive.db");
6459        let backend = std::sync::Arc::new(crate::StorageBackend::sqlite(&db_path).unwrap());
6460        {
6461            let mut writer = backend.pool().writer().unwrap();
6462            prepare_completed_v21_gc_fixture(writer.conn_mut());
6463        }
6464
6465        let root = dir.path().join("blobs");
6466        let store = std::sync::Arc::new(
6467            FsBlobStore::new(root.clone(), 0)
6468                .unwrap()
6469                .with_orphan_sweep_grace(Duration::from_secs(3600)),
6470        );
6471
6472        // The one real, unreferenced candidate the sweep will discover. Its
6473        // grace period is wide (1h) so correct handle-relative
6474        // classification always protects it.
6475        let real = store
6476            .put(b"real freshly-published blob".to_vec())
6477            .await
6478            .unwrap();
6479        let real_path = shard_path(&root, &real);
6480
6481        // Outside decoy sharing the SAME leaf name (content-ref hex), aged
6482        // far past the grace period. If classification ever reads through
6483        // a swapped symlink, the candidate looks like a stale, deletable
6484        // orphan instead of a protected fresh publish.
6485        let outside_dir = dir.path().join("outside");
6486        fs::create_dir_all(&outside_dir).unwrap();
6487        let decoy_path = outside_dir.join(real.as_str());
6488        fs::write(&decoy_path, b"outside decoy, must never be observed").unwrap();
6489        let ancient = SystemTime::now() - Duration::from_secs(7200);
6490        fs::OpenOptions::new()
6491            .write(true)
6492            .open(&decoy_path)
6493            .unwrap()
6494            .set_modified(ancient)
6495            .unwrap();
6496
6497        let (reached, release) = walk_leaf_sync_hook::install(&root);
6498        let sweep = {
6499            let store = store.clone();
6500            let sql = backend.sql();
6501            tokio::spawn(async move { store.transactional_orphan_sweep(sql.as_ref(), true).await })
6502        };
6503        assert!(
6504            recv_blocking(reached).await,
6505            "sweep walk must reach the leaf classification pause"
6506        );
6507
6508        // Swap the real leaf out from under the paused walk: same name, now
6509        // a symlink resolving outside the retained root.
6510        fs::remove_file(&real_path).unwrap();
6511        std::os::unix::fs::symlink(&decoy_path, &real_path).unwrap();
6512
6513        release.send(()).unwrap();
6514        let result = sweep.await.unwrap().unwrap();
6515
6516        assert_eq!(
6517            result.would_delete, 0,
6518            "an outside decoy's stale mtime must never make an in-root candidate \
6519             eligible for deletion: {result:?}"
6520        );
6521        assert_eq!(
6522            result.grace_period_skipped, 0,
6523            "the swapped leaf is a symlink; `openat(..., O_NOFOLLOW)` refuses it, so it \
6524             must be dropped from candidates entirely rather than counted (real or \
6525             outside) at all: {result:?}"
6526        );
6527        assert_eq!(
6528            result.scanned, 0,
6529            "the symlinked leaf must never be scanned as a candidate: {result:?}"
6530        );
6531
6532        // Restore a real leaf and confirm the walk is not permanently wedged
6533        // by the hook: an un-swapped root still finds and reports a real,
6534        // grace-protected candidate normally.
6535        fs::remove_file(&real_path).unwrap();
6536        fs::write(&real_path, b"real freshly-published blob").unwrap();
6537        let control = store
6538            .transactional_orphan_sweep(backend.sql().as_ref(), true)
6539            .await
6540            .unwrap();
6541        assert_eq!(
6542            control.grace_period_skipped, 1,
6543            "control: the un-replaced root must still classify the real candidate as \
6544             grace-protected: {control:?}"
6545        );
6546        assert_eq!(control.would_delete, 0);
6547    }
6548
6549    #[cfg(unix)]
6550    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6551    async fn transactional_orphan_sweep_walk_still_finds_a_real_orphan_past_its_grace_period() {
6552        // Control for the swap test above, using an independent root (no
6553        // hook installed): the handle-relative walk must still classify and
6554        // report a genuine past-grace orphan as deletable, proving the fix
6555        // narrows the walk to descriptor-relative reads without disabling
6556        // orphan detection itself.
6557        let dir = tempfile::tempdir().unwrap();
6558        let db_path = dir.path().join("khive.db");
6559        let backend = crate::StorageBackend::sqlite(&db_path).unwrap();
6560        {
6561            let mut writer = backend.pool().writer().unwrap();
6562            prepare_completed_v21_gc_fixture(writer.conn_mut());
6563        }
6564        let root = dir.path().join("blobs");
6565        let store = FsBlobStore::new(root.clone(), 0)
6566            .unwrap()
6567            .with_orphan_sweep_grace(Duration::from_secs(60));
6568
6569        let orphan = store.put(b"aged real orphan".to_vec()).await.unwrap();
6570        let path = shard_path(&root, &orphan);
6571        let ancient = SystemTime::now() - Duration::from_secs(3600);
6572        fs::OpenOptions::new()
6573            .write(true)
6574            .open(&path)
6575            .unwrap()
6576            .set_modified(ancient)
6577            .unwrap();
6578
6579        let result = store
6580            .transactional_orphan_sweep(backend.sql().as_ref(), false)
6581            .await
6582            .unwrap();
6583        assert_eq!(
6584            result.deleted, 1,
6585            "a real orphan older than the grace period must still be swept: {result:?}"
6586        );
6587        assert!(!store.exists(&orphan).await.unwrap());
6588    }
6589}