use super::config::GcConfig;
use super::live_set::collect_live_set;
use super::run::{gc_namespace, gc_namespace_with_reverify_chunk};
use crate::checkpoint::advance_retention_floor;
use crate::checkpoint::record::release_checkpoint_record;
use crate::checkpoint::tests::{create_checkpoint, mutation_context, write_test_file};
use crate::commit_engine::{CommitCandidate, NamespaceCommitEngine};
use crate::context::MutationContext;
use crate::error::CoreError;
use crate::limits::{
CONTENT_RECLAMATION_GRACE_MS, FORK_CHECKPOINT_LEASE_MS, GC_MIN_GRACE_WINDOW_MS,
UPLOAD_SESSION_LEASE_MS,
};
use crate::path::write::{CommitRequest, FilesystemOperation};
use loonfs_api::v0::GcResponse;
use loonfs_api::wire::control::{
decode_control_object, CheckpointOwner, CheckpointRecordLifecycle, CheckpointRecordState,
ControlObjectKind, UploadSessionLifecycle, UploadSessionState,
};
use loonfs_api::{ContentRef, ContentStoreId, NamespaceId, UploadId};
use loonfs_objectstore::keys::{
checkpoint_prefix, metadata_manifest_object, metadata_manifest_prefix, metadata_table,
metadata_table_prefix, wal_segment, wal_segment_prefix,
};
use loonfs_objectstore::ObjectStore;
use std::collections::BTreeSet;
use std::num::NonZeroUsize;
use crate::commit_engine::delete_namespace;
use crate::namespace::bootstrap::bootstrap_namespace;
use crate::namespace::fork::fork_namespace;
use crate::options::DeleteNamespaceOptions;
use crate::path::read::{load_metadata_view, ReadLoadContext};
use bytes::Bytes;
use futures::stream::BoxStream;
use loonfs_objectstore::local_fs_store::LocalFsStore;
use loonfs_objectstore::{ByteRange, ObjectBody, ObjectMetadata, ObjectStoreError, PutMode};
use loonfs_test_support::stores::{
BlockingStore, CountingStore, KeyPredicate, MetadataMapStore, OperationContext, OperationKind,
};
use std::sync::atomic::{AtomicUsize, Ordering};
use tempfile::tempdir;
const GRACE_MS: u64 = 60 * 60 * 1000;
fn config() -> GcConfig {
GcConfig {
grace_window_ms: GRACE_MS,
max_objects: None,
cursor: None,
}
}
fn context(now_ms: u64) -> MutationContext {
mutation_context("gc-test", now_ms)
}
async fn checkpoint_lifecycle<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
checkpoint_id: &loonfs_api::CheckpointId,
) -> CheckpointRecordLifecycle {
crate::checkpoint::read_checkpoint_record(store, namespace_id, checkpoint_id)
.await
.expect("read checkpoint record")
.expect("checkpoint record exists")
.state
.state
}
async fn now_after_newest_object(
store: &LocalFsStore,
namespace_id: &NamespaceId,
offset_ms: u64,
) -> u64 {
let prefix = loonfs_objectstore::keys::namespace_prefix(namespace_id);
let mut newest = 0;
for key in store.list_prefix(&prefix).await.expect("list namespace") {
let modified = store
.head(&key)
.await
.expect("head object")
.expect("object exists")
.last_modified_ms
.expect("local fs provides timestamps");
newest = newest.max(modified);
}
assert!(newest > 0, "namespace tree must not be empty");
newest + offset_ms
}
async fn stat_root<S: ObjectStore>(store: &S, namespace_id: &NamespaceId) {
load_metadata_view(store, namespace_id, ReadLoadContext::latest())
.await
.expect("load latest view")
.resolve_path("/")
.await
.expect("resolve root");
}
#[derive(Debug)]
struct IncompleteGcAccountingStore {
inner: LocalFsStore,
deletes: AtomicUsize,
lists: AtomicUsize,
}
#[derive(Debug, Clone, Copy)]
enum BlockingControlCasTarget {
CheckpointReleased,
UploadCompleted,
UploadAborted,
}
impl BlockingControlCasTarget {
fn matches(self, bytes: &[u8]) -> bool {
match self {
BlockingControlCasTarget::CheckpointReleased => {
let Ok(envelope) = decode_control_object::<CheckpointRecordState>(
bytes,
ControlObjectKind::CheckpointRecord,
) else {
return false;
};
matches!(
envelope.state.state,
CheckpointRecordLifecycle::Released { .. }
)
}
BlockingControlCasTarget::UploadCompleted | BlockingControlCasTarget::UploadAborted => {
let Ok(envelope) = decode_control_object::<UploadSessionState>(
bytes,
ControlObjectKind::UploadSession,
) else {
return false;
};
match self {
BlockingControlCasTarget::UploadCompleted => matches!(
envelope.state.state,
UploadSessionLifecycle::Completed { .. }
),
BlockingControlCasTarget::UploadAborted => {
matches!(envelope.state.state, UploadSessionLifecycle::Aborted { .. })
}
_ => false,
}
}
}
}
}
fn blocking_control_cas_store(
inner: LocalFsStore,
target: BlockingControlCasTarget,
) -> BlockingStore<LocalFsStore> {
let store = BlockingStore::matching(inner, move |operation: &OperationContext<'_>| {
let bytes = match operation.kind() {
OperationKind::CompareAndSwap { bytes, .. }
| OperationKind::Put {
bytes,
mode: PutMode::CompareAndSwap { .. },
} => bytes,
_ => return false,
};
target.matches(bytes)
});
store.block_next();
store
}
#[async_trait::async_trait]
impl ObjectStore for IncompleteGcAccountingStore {
async fn head(&self, key: &str) -> Result<Option<ObjectMetadata>, ObjectStoreError> {
self.inner.head(key).await
}
async fn get_with_metadata(&self, key: &str) -> Result<Option<ObjectBody>, ObjectStoreError> {
self.inner.get_with_metadata(key).await
}
async fn get(
&self,
key: &str,
range: Option<ByteRange>,
) -> Result<Option<Bytes>, ObjectStoreError> {
self.inner.get(key, range).await
}
async fn put(
&self,
key: &str,
bytes: Bytes,
mode: PutMode,
) -> Result<ObjectMetadata, ObjectStoreError> {
self.inner.put(key, bytes, mode).await
}
async fn delete(&self, key: &str) -> Result<(), ObjectStoreError> {
self.deletes.fetch_add(1, Ordering::SeqCst);
self.inner.delete(key).await
}
fn list_prefix_stream(
&self,
prefix: &str,
) -> BoxStream<'static, Result<String, ObjectStoreError>> {
self.lists.fetch_add(1, Ordering::SeqCst);
self.inner.list_prefix_stream(prefix)
}
}
#[tokio::test]
async fn gc_rejects_grace_windows_below_the_derived_minimum() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let too_small = GcConfig {
grace_window_ms: GC_MIN_GRACE_WINDOW_MS - 1,
..GcConfig::default()
};
let error = gc_namespace(&store, &namespace_id, &too_small, &context(1_000))
.await
.expect_err("sub-minimum grace window must be rejected");
assert!(
matches!(&error, CoreError::InvalidGcConfig(message)
if message.contains("below the derived safety minimum")),
"expected invalid gc config, got {error:?}"
);
assert_eq!(
error.code(),
crate::error::ErrorCode::InvalidRequest,
"the rejection surfaces as invalid_request"
);
let zero_budget = GcConfig {
max_objects: Some(0),
..config()
};
let error = gc_namespace(&store, &namespace_id, &zero_budget, &context(1_000))
.await
.expect_err("zero budget must be rejected");
assert!(matches!(error, CoreError::InvalidGcConfig(_)));
}
#[tokio::test]
async fn gc_reaps_below_floor_segments_after_the_grace_window() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("checkpoint");
advance_retention_floor(&store, &namespace_id, &setup)
.await
.expect("advance floor");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert_eq!(report.deleted_wal_segments, 1);
assert!(!report.degraded_retention);
stat_root(&store, &namespace_id).await;
}
async fn write_upload_session(store: &LocalFsStore, namespace_id: &NamespaceId) -> String {
let upload_id = loonfs_api::UploadId::parse("upl_0123456789abcdef0123456789abcdef")
.expect("valid upload id");
let state = loonfs_api::wire::control::UploadSessionState {
namespace_id: namespace_id.clone(),
upload_id: upload_id.clone(),
content_id: loonfs_api::ContentId::generate(),
created_at_ms: 1_000,
transport: loonfs_api::wire::control::UploadSessionTransport::ServiceProxied {},
state: loonfs_api::wire::control::UploadSessionLifecycle::Open {
expires_at_ms: 1_000 + UPLOAD_SESSION_LEASE_MS,
staged_content: None,
},
};
let envelope = loonfs_api::wire::control::UploadSessionEnvelope::from_state(
loonfs_api::wire::control::ControlObjectKind::UploadSession,
state,
)
.expect("session envelope");
let bytes =
loonfs_api::wire::control::encode_control_object(&envelope).expect("encode session");
let key = loonfs_objectstore::keys::upload_session(namespace_id.as_str(), upload_id.as_str());
store
.put_if_absent(&key, bytes::Bytes::from(bytes))
.await
.expect("write session");
key
}
async fn stage_upload<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
context: &MutationContext,
) -> (UploadId, ContentRef, ContentStoreId) {
let begin = crate::protocol::begin_upload(
store,
namespace_id,
loonfs_api::v0::BeginUploadRequest::ServiceProxied {},
context,
)
.await
.expect("begin upload");
let staged =
crate::protocol::upload_content(store, namespace_id, &begin.upload_id, b"racing upload\n")
.await
.expect("stage upload");
let content_store_id =
crate::namespace::catalog::load_namespace_content_store_id(store, namespace_id)
.await
.expect("content store id");
(begin.upload_id, staged.content_ref, content_store_id)
}
async fn read_upload_session<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
upload_id: &UploadId,
) -> Option<UploadSessionState> {
let key = loonfs_objectstore::keys::upload_session(namespace_id.as_str(), upload_id.as_str());
let body = store.get(&key, None).await.expect("read upload session")?;
Some(
decode_control_object::<UploadSessionState>(&body, ControlObjectKind::UploadSession)
.expect("decode upload session")
.state,
)
}
#[tokio::test]
async fn active_record_with_a_missing_basis_is_released_not_degrading() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let pinned = create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("first checkpoint");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
let moved_on = crate::checkpoint::create_checkpoint(
&store,
&namespace_id,
CheckpointOwner::User {
name: "other-pin".to_owned(),
},
None,
&setup,
)
.await
.expect("second checkpoint");
assert_ne!(moved_on.manifest_id, pinned.manifest_id);
let record = crate::checkpoint::record::read_checkpoint_record(
&store,
&namespace_id,
&pinned.checkpoint_id,
)
.await
.expect("read record")
.expect("record exists")
.state;
let basis_key = metadata_manifest_object(namespace_id.as_str(), &record.manifest_object_id);
store.delete(&basis_key).await.expect("drop basis manifest");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert_eq!(report.released_missing_basis_checkpoints, 1);
assert!(
!report.degraded_retention,
"a verifiably absent basis is not ambiguity"
);
let released = crate::checkpoint::record::read_checkpoint_record(
&store,
&namespace_id,
&pinned.checkpoint_id,
)
.await
.expect("read record")
.expect("record still present")
.state;
assert_eq!(
released.state,
loonfs_api::wire::control::CheckpointRecordLifecycle::Released {
released_at_ms: aged.now_ms
}
);
let again = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &again)
.await
.expect("second gc pass");
assert_eq!(report.released_missing_basis_checkpoints, 0);
assert!(!report.degraded_retention);
stat_root(&store, &namespace_id).await;
}
#[tokio::test]
async fn deleted_namespace_reclaims_down_to_its_tombstone() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("user pin");
delete_namespace(
&store,
&namespace_id,
DeleteNamespaceOptions::default(),
&setup,
)
.await
.expect("delete namespace");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert!(report.deleted_wal_segments >= 1);
assert!(report.deleted_metadata_tables >= 1);
assert!(report.deleted_manifests >= 1);
assert_eq!(report.released_expired_checkpoints, 1);
assert!(!report.degraded_retention);
let reaped = context(aged.now_ms + GRACE_MS);
let report = gc_namespace(&store, &namespace_id, &config(), &reaped)
.await
.expect("gc pass past the release grace window");
assert!(report.deleted_checkpoint_records >= 1);
assert!(!report.degraded_retention);
for prefix in [
wal_segment_prefix(namespace_id.as_str()),
metadata_table_prefix(namespace_id.as_str()),
metadata_manifest_prefix(namespace_id.as_str()),
checkpoint_prefix(namespace_id.as_str()),
] {
assert!(
store.list_prefix(&prefix).await.expect("list").is_empty(),
"prefix `{prefix}` must be empty after reclamation"
);
}
for key in [
loonfs_objectstore::keys::wal_head(namespace_id.as_str()),
loonfs_objectstore::keys::metadata_root(namespace_id.as_str()),
] {
assert!(
store.head(&key).await.expect("head").is_some(),
"tombstone object `{key}` must survive"
);
}
let again = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &again)
.await
.expect("second gc pass");
assert_eq!(report.deleted_wal_segments, 0);
assert_eq!(report.deleted_manifests, 0);
assert!(!report.degraded_retention);
}
#[tokio::test]
async fn fork_protected_bases_survive_source_deletion_until_the_target_dies() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let source = NamespaceId::parse("source").expect("namespace id");
let clone = NamespaceId::parse("clone").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &source, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &source, "/docs/shared.txt", "gc-shared", &setup).await;
fork_namespace(&store, &source, &clone, &setup)
.await
.expect("fork");
delete_namespace(&store, &source, DeleteNamespaceOptions::default(), &setup)
.await
.expect("delete source");
let fork_record = read_fork_record(&store, &source).await;
let basis_key = metadata_manifest_object(source.as_str(), &fork_record.manifest_object_id);
let aged = context(now_after_newest_object(&store, &source, GRACE_MS + 1).await);
let report = gc_namespace(&store, &source, &config(), &aged)
.await
.expect("gc pass with live clone");
assert_eq!(report.released_fork_checkpoints, 0);
assert!(!report.degraded_retention);
assert!(
store.head(&basis_key).await.expect("head basis").is_some(),
"fork basis must survive while the clone lives"
);
let clone_view = load_metadata_view(&store, &clone, ReadLoadContext::latest())
.await
.expect("load clone view");
clone_view
.resolve_path("/docs/shared.txt")
.await
.expect("clone reads through the deleted source");
delete_namespace(&store, &clone, DeleteNamespaceOptions::default(), &setup)
.await
.expect("delete clone");
let aged = context(now_after_newest_object(&store, &source, GRACE_MS + 1).await);
let report = gc_namespace(&store, &source, &config(), &aged)
.await
.expect("gc pass after clone delete");
assert_eq!(report.released_fork_checkpoints, 1);
assert!(report.deleted_manifests >= 1);
assert!(
store.head(&basis_key).await.expect("head basis").is_none(),
"the basis ages out once no living target needs it"
);
let again = context(now_after_newest_object(&store, &source, GRACE_MS + 1).await);
let report = gc_namespace(&store, &source, &config(), &again)
.await
.expect("idempotent pass");
assert_eq!(report.released_fork_checkpoints, 0);
assert_eq!(report.deleted_manifests, 0);
assert!(!report.degraded_retention);
}
#[tokio::test]
async fn upload_gc_aborts_an_expired_session_then_reaps_it() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let (upload_id, content_ref, content_store_id) =
stage_upload(&store, &namespace_id, &setup).await;
let session_key =
loonfs_objectstore::keys::upload_session(namespace_id.as_str(), upload_id.as_str());
let content_key =
loonfs_objectstore::keys::content_blob(content_store_id.as_str(), &content_ref.content_id);
let inside = context(setup.now_ms + UPLOAD_SESSION_LEASE_MS - 1);
let report = gc_namespace(&store, &namespace_id, &config(), &inside)
.await
.expect("gc pass inside the lease");
assert_eq!(report.deleted_upload_sessions, 0);
assert!(store.head(&content_key).await.expect("head").is_some());
let expired = context(setup.now_ms + UPLOAD_SESSION_LEASE_MS + GRACE_MS + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &expired)
.await
.expect("gc pass past the lease");
assert_eq!(
report.deleted_upload_sessions, 0,
"the record outlives its abort"
);
let session = read_upload_session(&store, &namespace_id, &upload_id)
.await
.expect("aborted session retained");
assert!(matches!(
session.state,
UploadSessionLifecycle::Aborted { .. }
));
assert!(
store.head(&content_key).await.expect("head").is_none(),
"aborting deletes the object the session owned"
);
let reaped = context(expired.now_ms + GRACE_MS + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &reaped)
.await
.expect("gc pass past the abort grace");
assert_eq!(report.deleted_upload_sessions, 1);
assert_eq!(
report.deleted_content_objects, 0,
"the abort half's unconditional cleanup is not a reclamation it can count"
);
assert!(store.head(&session_key).await.expect("head").is_none());
let again = context(reaped.now_ms + GRACE_MS);
let report = gc_namespace(&store, &namespace_id, &config(), &again)
.await
.expect("gc pass after the sweep");
assert_eq!(report.deleted_upload_sessions, 0);
}
#[tokio::test]
async fn a_pass_reports_the_soonest_deadline_it_retained() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let (upload_id, ..) = stage_upload(&store, &namespace_id, &setup).await;
let expires_at_ms = setup.now_ms + UPLOAD_SESSION_LEASE_MS;
let inside = context(setup.now_ms + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &inside)
.await
.expect("gc pass inside the lease");
assert_eq!(report.retained_candidates, 1);
assert_eq!(
report.next_reclamation_at_ms,
Some(expires_at_ms + GRACE_MS),
"an open session's reclamation waits for its lease and then the grace window"
);
let completed_at = context(setup.now_ms + 2);
complete_staged_upload(&store, &namespace_id, &upload_id, &completed_at).await;
let report = gc_namespace(&store, &namespace_id, &config(), &completed_at)
.await
.expect("gc pass over the completed session");
assert_eq!(
report.next_reclamation_at_ms,
Some(completed_at.now_ms + CONTENT_RECLAMATION_GRACE_MS),
"a completed session's content is protected by the derived grace, not the configured one"
);
}
#[tokio::test]
async fn an_aborted_session_is_reclaimed_from_the_deadline_the_pass_reported() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let (upload_id, ..) = stage_upload(&store, &namespace_id, &setup).await;
let session_key =
loonfs_objectstore::keys::upload_session(namespace_id.as_str(), upload_id.as_str());
let expired = context(setup.now_ms + UPLOAD_SESSION_LEASE_MS + GRACE_MS + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &expired)
.await
.expect("gc pass past the lease");
assert_eq!(report.deleted_upload_sessions, 0);
let reclaim_at_ms = report
.next_reclamation_at_ms
.expect("the abort this pass performed is a deadline it created");
assert_eq!(reclaim_at_ms, expired.now_ms + GRACE_MS);
let reclaiming = context(reclaim_at_ms + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &reclaiming)
.await
.expect("gc pass at the reported deadline");
assert_eq!(report.deleted_upload_sessions, 1);
assert!(store.head(&session_key).await.expect("head").is_none());
assert_eq!(
report.next_reclamation_at_ms, None,
"a pass that reclaimed everything it found owes no later visit"
);
}
async fn complete_staged_upload<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
upload_id: &UploadId,
context: &MutationContext,
) {
let content_store_id =
crate::namespace::catalog::load_namespace_content_store_id(store, namespace_id)
.await
.expect("content store id");
let session = read_upload_session(store, namespace_id, upload_id)
.await
.expect("open session");
let content_ref = match session.state {
UploadSessionLifecycle::Open { staged_content, .. } => staged_content,
UploadSessionLifecycle::Completed { .. } | UploadSessionLifecycle::Aborted { .. } => None,
}
.expect("a staged session is open and carries the reference it wrote");
crate::protocol::complete_upload(
store,
namespace_id,
&content_store_id,
upload_id,
&loonfs_api::v0::CompleteUploadRequest::for_content_ref(content_ref),
context,
)
.await
.expect("complete upload");
}
#[tokio::test]
async fn upload_gc_reaps_a_session_that_never_staged_anything() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let session_key = write_upload_session(&store, &namespace_id).await;
let expired = context(1_000 + UPLOAD_SESSION_LEASE_MS + GRACE_MS + 1);
gc_namespace(&store, &namespace_id, &config(), &expired)
.await
.expect("gc pass past the lease");
let reaped = context(expired.now_ms + GRACE_MS + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &reaped)
.await
.expect("gc pass past the abort grace");
assert_eq!(report.deleted_upload_sessions, 1);
assert!(store.head(&session_key).await.expect("head").is_none());
}
#[tokio::test]
async fn upload_completion_wins_before_gc_abort_and_the_session_is_retained() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let (upload_id, content_ref, content_store_id) =
stage_upload(&store, &namespace_id, &setup).await;
let aged = context(setup.now_ms + UPLOAD_SESSION_LEASE_MS + GRACE_MS + 1);
let content_key =
loonfs_objectstore::keys::content_blob(content_store_id.as_str(), &content_ref.content_id);
let store = blocking_control_cas_store(store, BlockingControlCasTarget::UploadAborted);
let gc_config = config();
let gc = gc_namespace(&store, &namespace_id, &gc_config, &aged);
let complete = async {
store.wait_until_blocked().await;
let result = crate::protocol::complete_upload(
&store,
&namespace_id,
&content_store_id,
&upload_id,
&loonfs_api::v0::CompleteUploadRequest::for_content_ref(content_ref.clone()),
&aged,
)
.await;
store.release();
result
};
let (report, completion) = tokio::join!(gc, complete);
completion.expect("completion wins the blocked abort CAS");
let report = report.expect("gc pass");
assert_eq!(report.deleted_upload_sessions, 0);
let session = read_upload_session(&store, &namespace_id, &upload_id)
.await
.expect("completed session retained");
assert!(matches!(
session.state,
UploadSessionLifecycle::Completed { .. }
));
assert!(
store.head(&content_key).await.expect("head").is_some(),
"the losing abort must not clean up the winner's content"
);
}
#[tokio::test]
async fn gc_abort_wins_before_completion_and_completion_reports_not_found() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let (upload_id, content_ref, content_store_id) =
stage_upload(&store, &namespace_id, &setup).await;
let aged = context(setup.now_ms + UPLOAD_SESSION_LEASE_MS + GRACE_MS + 1);
let content_key =
loonfs_objectstore::keys::content_blob(content_store_id.as_str(), &content_ref.content_id);
let store = blocking_control_cas_store(store, BlockingControlCasTarget::UploadCompleted);
let request = loonfs_api::v0::CompleteUploadRequest::for_content_ref(content_ref.clone());
let completion = crate::protocol::complete_upload(
&store,
&namespace_id,
&content_store_id,
&upload_id,
&request,
&aged,
);
let abort = async {
store.wait_until_blocked().await;
let report = gc_namespace(&store, &namespace_id, &config(), &aged).await;
store.release();
report
};
let (completion, report) = tokio::join!(completion, abort);
let error = completion.expect_err("an aborted session is logically absent");
assert!(matches!(&error, CoreError::UploadNotFound { .. }));
assert_eq!(error.code(), crate::error::ErrorCode::UploadNotFound);
report.expect("gc pass");
let session = read_upload_session(&store, &namespace_id, &upload_id)
.await
.expect("aborted session retained for a grace window");
assert!(matches!(
session.state,
UploadSessionLifecycle::Aborted { .. }
));
assert!(
store.head(&content_key).await.expect("head").is_none(),
"the winning abort cleans up, and the losing completion does not resurrect"
);
}
async fn complete_upload_for_gc<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
bytes: &[u8],
context: &MutationContext,
) -> (
UploadId,
ContentRef,
ContentStoreId,
crate::publish::PreparedContent,
) {
let begin = crate::protocol::begin_upload(
store,
namespace_id,
loonfs_api::v0::BeginUploadRequest::ServiceProxied {},
context,
)
.await
.expect("begin upload");
let staged = crate::protocol::upload_content(store, namespace_id, &begin.upload_id, bytes)
.await
.expect("stage upload");
let content_store_id =
crate::namespace::catalog::load_namespace_content_store_id(store, namespace_id)
.await
.expect("content store id");
let completed = crate::protocol::complete_upload(
store,
namespace_id,
&content_store_id,
&begin.upload_id,
&loonfs_api::v0::CompleteUploadRequest::for_content_ref(staged.content_ref.clone()),
context,
)
.await
.expect("complete upload");
(
begin.upload_id,
staged.content_ref,
content_store_id,
completed.prepared,
)
}
async fn publish_completed_content<S: ObjectStore>(
store: &S,
namespace_id: &NamespaceId,
path: &str,
content_ref: ContentRef,
prepared: crate::publish::PreparedContent,
context: &MutationContext,
) {
NamespaceCommitEngine::new(namespace_id.clone())
.publish_batch(
store,
vec![CommitCandidate::prepared(
CommitRequest::single(
loonfs_api::CommitId::parse("publish-completed-content").expect("commit id"),
None,
FilesystemOperation::PutFile {
path: loonfs_api::AbsolutePath::parse(path).expect("path"),
content_ref,
behavior: loonfs_api::DestinationBehavior::NoReplace,
expected_revision_no: None,
},
),
vec![prepared],
)],
context,
&crate::protocol::PublishTailOptions::default(),
)
.await
.results
.pop()
.expect("one result")
.expect("published");
}
#[tokio::test]
async fn content_gc_retains_completed_content_inside_its_grace() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let (upload_id, content_ref, content_store_id, _prepared) =
complete_upload_for_gc(&store, &namespace_id, b"unpublished\n", &setup).await;
let content_key =
loonfs_objectstore::keys::content_blob(content_store_id.as_str(), &content_ref.content_id);
let inside = context(setup.now_ms + CONTENT_RECLAMATION_GRACE_MS - 1);
let report = gc_namespace(&store, &namespace_id, &config(), &inside)
.await
.expect("gc pass inside the content grace");
assert_eq!(report.deleted_upload_sessions, 0);
assert_eq!(report.deleted_content_objects, 0);
assert!(store.head(&content_key).await.expect("head").is_some());
assert!(read_upload_session(&store, &namespace_id, &upload_id)
.await
.is_some());
}
#[tokio::test]
async fn content_gc_reclaims_completed_content_nothing_references() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/other.txt", "gc-other", &setup).await;
let (upload_id, content_ref, content_store_id, _prepared) =
complete_upload_for_gc(&store, &namespace_id, b"unpublished\n", &setup).await;
let content_key =
loonfs_objectstore::keys::content_blob(content_store_id.as_str(), &content_ref.content_id);
let past = context(setup.now_ms + CONTENT_RECLAMATION_GRACE_MS + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &past)
.await
.expect("gc pass past the content grace");
assert_eq!(report.deleted_upload_sessions, 1);
assert_eq!(report.deleted_content_objects, 1);
assert!(
store.head(&content_key).await.expect("head").is_none(),
"completed content nothing published is reclaimable"
);
assert!(read_upload_session(&store, &namespace_id, &upload_id)
.await
.is_none());
}
#[tokio::test]
async fn content_gc_never_reclaims_published_content() {
for materialize in [false, true] {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let (upload_id, content_ref, content_store_id, prepared) =
complete_upload_for_gc(&store, &namespace_id, b"published\n", &setup).await;
publish_completed_content(
&store,
&namespace_id,
"/docs/published.txt",
content_ref.clone(),
prepared,
&setup,
)
.await;
if materialize {
crate::checkpoint::flush_wal(&store, &namespace_id, &setup)
.await
.expect("flush wal");
advance_retention_floor(&store, &namespace_id, &setup)
.await
.expect("advance floor");
}
let content_key = loonfs_objectstore::keys::content_blob(
content_store_id.as_str(),
&content_ref.content_id,
);
let past = context(setup.now_ms + CONTENT_RECLAMATION_GRACE_MS + 1);
let report = gc_namespace(&store, &namespace_id, &config(), &past)
.await
.expect("gc pass past the content grace");
assert_eq!(
report.deleted_upload_sessions, 1,
"materialize={materialize}"
);
assert_eq!(
report.deleted_content_objects, 0,
"materialize={materialize}"
);
assert!(
store.head(&content_key).await.expect("head").is_some(),
"published content survives its session (materialize={materialize})"
);
assert!(read_upload_session(&store, &namespace_id, &upload_id)
.await
.is_none());
assert!(!report.degraded_retention);
}
}
async fn namespace_with_a_scan_worth_bounding(
store: &LocalFsStore,
namespace_id: &NamespaceId,
setup: &MutationContext,
) {
bootstrap_namespace(store, namespace_id, setup, false)
.await
.expect("bootstrap");
for index in 0..3 {
write_test_file(
store,
namespace_id,
&format!("/docs/materialized-{index}.txt"),
&format!("scan-fixture-{index}"),
setup,
)
.await;
}
crate::checkpoint::flush_wal(store, namespace_id, setup)
.await
.expect("flush wal");
for index in 0..3 {
write_test_file(
store,
namespace_id,
&format!("/docs/tail-{index}.txt"),
&format!("scan-fixture-tail-{index}"),
setup,
)
.await;
}
}
#[tokio::test]
async fn a_budget_that_dies_inside_the_reference_scan_defers_and_walks_on() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
namespace_with_a_scan_worth_bounding(&store, &namespace_id, &setup).await;
let (upload_id, content_ref, content_store_id, _prepared) =
complete_upload_for_gc(&store, &namespace_id, b"unpublished\n", &setup).await;
let content_key =
loonfs_objectstore::keys::content_blob(content_store_id.as_str(), &content_ref.content_id);
let live = collect_live_set(&store, &namespace_id, &setup)
.await
.expect("collect live set");
assert!(
!live.manifests.is_empty() && !live.wal_segments.is_empty(),
"the fixture must give the scan more than one object to read"
);
let past = context(setup.now_ms + CONTENT_RECLAMATION_GRACE_MS + 1);
let mut tiny = config();
tiny.max_objects = Some(1);
let mut cursor: Option<String> = None;
let mut deferred = false;
let mut passes = 0;
loop {
passes += 1;
assert!(
passes <= 64,
"a one-object budget must still walk the namespace to the end"
);
tiny.cursor.clone_from(&cursor);
let pass = gc_namespace(&store, &namespace_id, &tiny, &past)
.await
.expect("one-object pass");
assert_eq!(pass.deleted_upload_sessions, 0);
assert_eq!(pass.deleted_content_objects, 0);
deferred |= pass.content_reclamation_deferred;
let Some(next) = pass.next_cursor else {
break;
};
assert_ne!(
Some(next.as_str()),
cursor.as_deref(),
"pass {passes} handed back the cursor it came in with"
);
cursor = Some(next);
}
assert!(
deferred,
"a one-object budget cannot afford the scan, and the pass must say so"
);
assert!(
store.head(&content_key).await.expect("head").is_some(),
"a deferred pass reclaims nothing"
);
assert!(
read_upload_session(&store, &namespace_id, &upload_id)
.await
.is_some(),
"the session that triggered the scan is retained, not reclaimed"
);
let mut enough = config();
enough.max_objects = Some(1_024);
let resumed = gc_namespace(&store, &namespace_id, &enough, &past)
.await
.expect("pass with room for the scan");
assert!(!resumed.content_reclamation_deferred);
assert_eq!(resumed.deleted_upload_sessions, 1);
assert_eq!(resumed.deleted_content_objects, 1);
assert!(store.head(&content_key).await.expect("head").is_none());
assert!(read_upload_session(&store, &namespace_id, &upload_id)
.await
.is_none());
}
#[tokio::test]
async fn no_budget_lets_a_partial_reference_set_decide_a_deletion() {
let temp_dir = tempdir().expect("tempdir");
let seed_root = temp_dir.path().join("seed");
let seed = LocalFsStore::new(&seed_root).expect("seed store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
namespace_with_a_scan_worth_bounding(&seed, &namespace_id, &setup).await;
let (upload_id, content_ref, content_store_id, prepared) =
complete_upload_for_gc(&seed, &namespace_id, b"published-last\n", &setup).await;
publish_completed_content(
&seed,
&namespace_id,
"/docs/published.txt",
content_ref.clone(),
prepared,
&setup,
)
.await;
let content_key =
loonfs_objectstore::keys::content_blob(content_store_id.as_str(), &content_ref.content_id);
let past = context(setup.now_ms + CONTENT_RECLAMATION_GRACE_MS + 1);
let mut some_budget_reached_the_verdict = false;
let mut some_budget_deferred_instead = false;
for max_objects in 1..=16 {
let trial_root = temp_dir.path().join(format!("trial-{max_objects}"));
copy_tree(&seed_root, &trial_root);
let store = LocalFsStore::new(&trial_root).expect("trial store");
let mut bounded = config();
bounded.max_objects = Some(max_objects);
let mut cursor: Option<String> = None;
let mut deferred = false;
let mut passes = 0;
loop {
passes += 1;
assert!(
passes <= 256,
"max_objects={max_objects}: the walk must reach its end"
);
bounded.cursor.clone_from(&cursor);
let pass = gc_namespace(&store, &namespace_id, &bounded, &past)
.await
.expect("bounded pass");
assert_eq!(pass.deleted_content_objects, 0, "max_objects={max_objects}");
assert!(
store.head(&content_key).await.expect("head").is_some(),
"max_objects={max_objects}: referenced content survives every budget"
);
deferred |= pass.content_reclamation_deferred;
let Some(next) = pass.next_cursor else {
break;
};
assert_ne!(
Some(next.as_str()),
cursor.as_deref(),
"max_objects={max_objects}: pass {passes} handed back its own cursor"
);
cursor = Some(next);
}
assert!(
store.head(&content_key).await.expect("head").is_some(),
"max_objects={max_objects}"
);
let decided = read_upload_session(&store, &namespace_id, &upload_id)
.await
.is_none();
assert_eq!(
deferred, !decided,
"max_objects={max_objects}: the session survives exactly when the scan was deferred"
);
some_budget_reached_the_verdict |= decided;
some_budget_deferred_instead |= deferred;
}
assert!(
some_budget_reached_the_verdict,
"a budget large enough to finish the scan must still decide the session"
);
assert!(
some_budget_deferred_instead,
"the sweep must actually run out mid-scan somewhere in this range, or \
the test proves nothing about partial reference sets"
);
}
#[tokio::test]
async fn gc_retains_everything_inside_the_grace_window() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("checkpoint");
advance_retention_floor(&store, &namespace_id, &setup)
.await
.expect("advance floor");
let young = context(now_after_newest_object(&store, &namespace_id, 0).await);
let report = gc_namespace(&store, &namespace_id, &config(), &young)
.await
.expect("gc pass");
assert_eq!(report.deleted_wal_segments, 0);
assert_eq!(report.deleted_metadata_tables, 0);
assert_eq!(report.deleted_manifests, 0);
assert!(report.retained_candidates > 0);
assert_eq!(reason_total(&report), report.retained_candidates);
assert!(report.retained.grace_window > 0);
assert_eq!(report.retained.no_provider_timestamp, 0);
stat_root(&store, &namespace_id).await;
}
fn reason_total(report: &GcResponse) -> u64 {
report
.retained
.by_reason()
.into_iter()
.map(|(_, count)| count)
.sum()
}
#[tokio::test]
async fn a_pass_names_a_checkpoint_record_it_could_not_advance() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let pinned = create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("checkpoint");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS * 2).await);
crate::checkpoint::release_checkpoint(&store, &namespace_id, &pinned.checkpoint_id, &aged)
.await
.expect("release checkpoint");
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert_eq!(report.deleted_checkpoint_records, 0);
assert_eq!(report.retained.checkpoint_not_releasable, 1);
assert_eq!(reason_total(&report), report.retained_candidates);
}
#[tokio::test]
async fn gc_never_deletes_the_live_replay_chain() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("checkpoint");
advance_retention_floor(&store, &namespace_id, &setup)
.await
.expect("advance floor");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert_eq!(report.deleted_wal_segments, 1);
let view = load_metadata_view(&store, &namespace_id, ReadLoadContext::latest())
.await
.expect("load view");
view.resolve_path("/docs/two.txt")
.await
.expect("tail commit stays readable");
}
#[tokio::test]
async fn gc_reaps_dead_checkpoints_before_their_basis_across_passes() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let first = create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("first checkpoint");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("second checkpoint");
let first_record =
crate::checkpoint::read_checkpoint_record(&store, &namespace_id, &first.checkpoint_id)
.await
.expect("read first record")
.expect("first record exists")
.state;
release_checkpoint_record(&store, &namespace_id, &first.checkpoint_id, setup.now_ms)
.await
.expect("mark first dead");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let first_pass = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("first gc pass");
assert_eq!(first_pass.deleted_checkpoint_records, 1);
assert!(!first_pass.degraded_retention);
assert!(
crate::checkpoint::read_checkpoint_record(&store, &namespace_id, &first.checkpoint_id)
.await
.expect("read record")
.is_none()
);
let basis = crate::checkpoint::load_namespace_manifest_envelope(
&store,
&namespace_id,
&first_record.manifest_object_id,
)
.await
.expect("dead basis manifest survives its record");
for file in &basis.payload.metadata_files {
assert!(
store
.head(&file.object_key)
.await
.expect("head table")
.is_some(),
"dead basis table survives its record"
);
}
let second_pass = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("second gc pass");
assert!(
second_pass.deleted_manifests >= 1,
"dead basis manifest reaped once its record is gone"
);
assert!(crate::checkpoint::load_namespace_manifest_envelope(
&store,
&namespace_id,
&first_record.manifest_object_id,
)
.await
.is_err());
stat_root(&store, &namespace_id).await;
}
#[tokio::test]
async fn gc_retains_unrecognized_manifest_keys() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let manifest_prefix = metadata_manifest_prefix(namespace_id.as_str());
let foreign_objects = [
(
format!("{manifest_prefix}notes.txt"),
b"foreign key".as_slice(),
),
(
format!("{manifest_prefix}invalid.manifest.json"),
b"invalid manifest id".as_slice(),
),
];
for (key, bytes) in &foreign_objects {
store
.put_if_absent(key, Bytes::copy_from_slice(bytes))
.await
.expect("write foreign manifest-prefix object");
}
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert_eq!(report.deleted_manifests, 0);
for (key, expected) in foreign_objects {
let actual = store
.get(&key, None)
.await
.expect("get foreign manifest-prefix object")
.expect("unrecognized object is retained");
assert_eq!(actual.as_ref(), expected);
}
}
#[tokio::test]
async fn gc_reclaims_manifests_superseded_by_wal_flushes() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
for round in 0..3 {
write_test_file(
&store,
&namespace_id,
&format!("/docs/file-{round}.txt"),
&format!("gc-adv-{round}"),
&setup,
)
.await;
crate::checkpoint::flush_wal(&store, &namespace_id, &setup)
.await
.expect("flush wal");
}
assert!(
store
.list_prefix(&checkpoint_prefix(namespace_id.as_str()))
.await
.expect("list checkpoint records")
.is_empty(),
"a wal flush must not create checkpoint records"
);
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert_eq!(report.deleted_manifests, 2);
assert!(!report.degraded_retention);
let manifests_left = store
.list_prefix(&metadata_manifest_prefix(namespace_id.as_str()))
.await
.expect("list manifests");
assert_eq!(manifests_left.len(), 1, "only the live root manifest stays");
let fold_policy = crate::checkpoint::MetadataLsmPolicy {
max_l0_runs: NonZeroUsize::MIN,
..Default::default()
};
for _ in 0..16 {
let report =
crate::checkpoint::reorganize_metadata_step(&store, &namespace_id, &setup, fold_policy)
.await
.expect("reorganize step");
if matches!(
report.outcome,
crate::checkpoint::MetadataReorganizeOutcome::NotNeeded { .. }
) {
break;
}
}
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let after_fold = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass after reorganization");
assert!(
after_fold.deleted_metadata_tables > 0,
"folded-away run tables become collectable"
);
assert!(!after_fold.degraded_retention);
stat_root(&store, &namespace_id).await;
let view = load_metadata_view(&store, &namespace_id, ReadLoadContext::latest())
.await
.expect("load view");
for round in 0..3 {
view.resolve_path(&format!("/docs/file-{round}.txt"))
.await
.expect("file readable after sweep");
}
}
#[tokio::test]
async fn gc_reaps_released_checkpoints_before_their_basis_across_passes() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let pinned = create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("pin checkpoint");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("advance past the pinned basis");
let first_release =
crate::checkpoint::release_checkpoint(&store, &namespace_id, &pinned.checkpoint_id, &setup)
.await
.expect("release");
assert!(first_release.was_active);
let repeat_release =
crate::checkpoint::release_checkpoint(&store, &namespace_id, &pinned.checkpoint_id, &setup)
.await
.expect("repeat release");
assert!(!repeat_release.was_active);
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let first_pass = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("first gc pass");
assert_eq!(first_pass.deleted_checkpoint_records, 1);
let second_pass = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("second gc pass");
assert!(second_pass.deleted_manifests >= 1);
let after_reap =
crate::checkpoint::release_checkpoint(&store, &namespace_id, &pinned.checkpoint_id, &setup)
.await
.expect("release after reap");
assert!(!after_reap.was_active);
stat_root(&store, &namespace_id).await;
}
#[tokio::test]
async fn caller_release_and_expiry_release_converge_on_the_winners_stamp() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let pin = |name: &'static str| {
crate::checkpoint::create_checkpoint(
&store,
&namespace_id,
CheckpointOwner::User {
name: name.to_owned(),
},
Some(setup.now_ms + GRACE_MS),
&setup,
)
};
let pass_first = pin("pass-first").await.expect("expiring checkpoint");
let caller_first = pin("caller-first").await.expect("expiring checkpoint");
let expired = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let caller_stamp = expired.now_ms + 1;
let released = crate::checkpoint::release_checkpoint(
&store,
&namespace_id,
&caller_first.checkpoint_id,
&context(caller_stamp),
)
.await
.expect("caller release");
assert!(released.was_active);
let report = gc_namespace(&store, &namespace_id, &config(), &expired)
.await
.expect("gc pass");
assert_eq!(
report.released_expired_checkpoints, 1,
"only the record the caller left alone is released here"
);
assert_eq!(
checkpoint_lifecycle(&store, &namespace_id, &caller_first.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: caller_stamp
},
"the winner's stamp stands"
);
assert_eq!(
checkpoint_lifecycle(&store, &namespace_id, &pass_first.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: expired.now_ms
}
);
let late = crate::checkpoint::release_checkpoint(
&store,
&namespace_id,
&pass_first.checkpoint_id,
&context(caller_stamp),
)
.await
.expect("a release that lost is still success");
assert!(!late.was_active);
assert_eq!(
checkpoint_lifecycle(&store, &namespace_id, &pass_first.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: expired.now_ms
},
"the loser rewrites nothing"
);
}
#[tokio::test]
async fn a_release_that_loses_its_etag_retains_without_erroring() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let pinned = crate::checkpoint::create_checkpoint(
&store,
&namespace_id,
CheckpointOwner::User {
name: "short-lived".to_owned(),
},
Some(setup.now_ms + GRACE_MS),
&setup,
)
.await
.expect("expiring checkpoint");
let expired = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let caller_stamp = expired.now_ms + 1;
let store = blocking_control_cas_store(store, BlockingControlCasTarget::CheckpointReleased);
let gc_config = config();
let pass = gc_namespace(&store, &namespace_id, &gc_config, &expired);
let caller = async {
store.wait_until_blocked().await;
let released = crate::checkpoint::release_checkpoint(
&store,
&namespace_id,
&pinned.checkpoint_id,
&context(caller_stamp),
)
.await;
store.release();
released
};
let (report, released) = tokio::join!(pass, caller);
assert!(released.expect("caller release").was_active);
let report = report.expect("the pass finishes");
assert_eq!(report.released_expired_checkpoints, 0);
assert_eq!(report.deleted_checkpoint_records, 0);
assert_eq!(
checkpoint_lifecycle(&store, &namespace_id, &pinned.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: caller_stamp
}
);
}
#[tokio::test]
async fn gc_deletes_a_released_record_only_after_its_release_ages() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let pinned = create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("pin checkpoint");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("advance past the pinned basis");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS * 4).await);
crate::checkpoint::release_checkpoint(&store, &namespace_id, &pinned.checkpoint_id, &aged)
.await
.expect("release");
assert_eq!(
checkpoint_lifecycle(&store, &namespace_id, &pinned.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: aged.now_ms
}
);
let inside_grace = context(aged.now_ms + GRACE_MS - 1);
let report = gc_namespace(&store, &namespace_id, &config(), &inside_grace)
.await
.expect("pass inside the release grace window");
assert_eq!(
report.deleted_checkpoint_records, 0,
"an old object with a young release is retained"
);
assert!(crate::checkpoint::read_checkpoint_record(
&store,
&namespace_id,
&pinned.checkpoint_id
)
.await
.expect("read record")
.is_some());
let past_grace = context(aged.now_ms + GRACE_MS);
let report = gc_namespace(&store, &namespace_id, &config(), &past_grace)
.await
.expect("pass past the release grace window");
assert_eq!(report.deleted_checkpoint_records, 1);
assert!(crate::checkpoint::read_checkpoint_record(
&store,
&namespace_id,
&pinned.checkpoint_id
)
.await
.expect("read record")
.is_none());
}
#[tokio::test]
async fn gc_reaps_expired_checkpoints_before_their_basis_across_passes() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let expiring = crate::checkpoint::create_checkpoint(
&store,
&namespace_id,
CheckpointOwner::User {
name: "short-lived".to_owned(),
},
Some(setup.now_ms + GRACE_MS),
&setup,
)
.await
.expect("expiring checkpoint");
let lasting = crate::checkpoint::create_checkpoint(
&store,
&namespace_id,
CheckpointOwner::User {
name: "long-lived".to_owned(),
},
Some(u64::MAX),
&setup,
)
.await
.expect("lasting checkpoint");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("advance past the expiring basis");
let expired = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
assert!(
expired.now_ms > 1_000 + GRACE_MS,
"provider clock sits past the expiry"
);
let first_pass = gc_namespace(&store, &namespace_id, &config(), &expired)
.await
.expect("post-expiry pass");
assert_eq!(first_pass.released_expired_checkpoints, 1);
assert_eq!(first_pass.deleted_checkpoint_records, 0);
assert_eq!(
checkpoint_lifecycle(&store, &namespace_id, &expiring.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: expired.now_ms
}
);
let aged_out = context(expired.now_ms + GRACE_MS);
let second_pass = gc_namespace(&store, &namespace_id, &config(), &aged_out)
.await
.expect("second post-expiry pass");
assert_eq!(second_pass.deleted_checkpoint_records, 1);
assert!(crate::checkpoint::read_checkpoint_record(
&store,
&namespace_id,
&expiring.checkpoint_id
)
.await
.expect("read record")
.is_none());
let survivor =
crate::checkpoint::read_checkpoint_record(&store, &namespace_id, &lasting.checkpoint_id)
.await
.expect("read lasting record")
.expect("lasting record survives")
.state;
assert!(crate::checkpoint::load_namespace_manifest_envelope(
&store,
&namespace_id,
&survivor.manifest_object_id,
)
.await
.is_ok());
stat_root(&store, &namespace_id).await;
}
#[tokio::test]
async fn gc_keeps_a_basis_pinned_by_another_owner_after_one_release() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let first = crate::checkpoint::create_checkpoint(
&store,
&namespace_id,
CheckpointOwner::User {
name: "keeper".to_owned(),
},
None,
&setup,
)
.await
.expect("first owner");
let second = crate::checkpoint::create_checkpoint(
&store,
&namespace_id,
CheckpointOwner::User {
name: "releaser".to_owned(),
},
None,
&setup,
)
.await
.expect("second owner");
assert_ne!(first.checkpoint_id, second.checkpoint_id);
assert_eq!(first.manifest_id, second.manifest_id);
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("advance past the shared basis");
crate::checkpoint::release_checkpoint(&store, &namespace_id, &second.checkpoint_id, &setup)
.await
.expect("release one owner");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let first_pass = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("first gc pass");
assert_eq!(first_pass.deleted_checkpoint_records, 1);
let second_pass = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("second gc pass");
let _ = second_pass;
assert!(
crate::checkpoint::read_checkpoint_record(&store, &namespace_id, &first.checkpoint_id)
.await
.expect("read keeper record")
.is_some(),
"the surviving owner's record stays"
);
let keeper =
crate::checkpoint::read_checkpoint_record(&store, &namespace_id, &first.checkpoint_id)
.await
.expect("read keeper record")
.expect("keeper record exists")
.state;
assert!(
crate::checkpoint::load_namespace_manifest_envelope(
&store,
&namespace_id,
&keeper.manifest_object_id,
)
.await
.is_ok(),
"shared basis survives while any owner remains"
);
}
#[tokio::test]
async fn fork_owned_checkpoints_reject_user_release() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let source = NamespaceId::parse("source").expect("namespace id");
let clone = NamespaceId::parse("clone").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &source, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &source, "/docs/one.txt", "gc-one", &setup).await;
fork_namespace(&store, &source, &clone, &setup)
.await
.expect("fork");
let fork_record = read_fork_record(&store, &source).await;
let error =
crate::checkpoint::release_checkpoint(&store, &source, &fork_record.checkpoint_id, &setup)
.await
.expect_err("fork-owned release must fail");
assert!(
matches!(
&error,
CoreError::InvalidCheckpointRequest(message)
if message.contains("owned by fork target")
),
"expected invalid checkpoint request, got {error:?}"
);
}
#[tokio::test]
async fn gc_retains_active_checkpoint_bases() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let first = create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("first checkpoint");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("second checkpoint");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert!(report.deleted_manifests <= 1);
assert_eq!(report.deleted_checkpoint_records, 0);
let first_record =
crate::checkpoint::read_checkpoint_record(&store, &namespace_id, &first.checkpoint_id)
.await
.expect("read first checkpoint")
.expect("first checkpoint exists")
.state;
assert!(crate::checkpoint::load_namespace_manifest_envelope(
&store,
&namespace_id,
&first_record.manifest_object_id,
)
.await
.is_ok());
}
async fn read_fork_record(store: &LocalFsStore, source: &NamespaceId) -> CheckpointRecordState {
for key in store
.list_prefix(&checkpoint_prefix(source.as_str()))
.await
.expect("list checkpoints")
{
let bytes = store
.get(&key, None)
.await
.expect("get record")
.expect("record exists");
let record = decode_control_object::<CheckpointRecordState>(
&bytes,
ControlObjectKind::CheckpointRecord,
)
.expect("decode record")
.state;
if matches!(record.owner, CheckpointOwner::Fork { .. }) {
return record;
}
}
unreachable!("fork leaves one fork-owned record");
}
#[tokio::test]
async fn gc_releases_fork_checkpoints_of_terminally_deleted_targets_across_passes() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let source = NamespaceId::parse("source").expect("namespace id");
let clone = NamespaceId::parse("clone").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &source, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &source, "/docs/one.txt", "gc-one", &setup).await;
fork_namespace(&store, &source, &clone, &setup)
.await
.expect("fork");
let fork_record = read_fork_record(&store, &source).await;
write_test_file(&store, &source, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &source, &setup)
.await
.expect("advance root past the fork basis");
let before = context(now_after_newest_object(&store, &source, GRACE_MS + 1).await);
let report = gc_namespace(&store, &source, &config(), &before)
.await
.expect("gc with live target");
assert_eq!(report.released_fork_checkpoints, 0);
delete_namespace(&store, &clone, DeleteNamespaceOptions::default(), &setup)
.await
.expect("terminal delete of the fork target");
let aged = context(now_after_newest_object(&store, &source, GRACE_MS + 1).await);
let first_pass = gc_namespace(&store, &source, &config(), &aged)
.await
.expect("first gc pass");
assert_eq!(first_pass.released_fork_checkpoints, 1);
assert_eq!(first_pass.deleted_checkpoint_records, 0);
assert_eq!(
checkpoint_lifecycle(&store, &source, &fork_record.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: aged.now_ms
}
);
let aged_out = context(aged.now_ms + GRACE_MS);
let second_pass = gc_namespace(&store, &source, &config(), &aged_out)
.await
.expect("second gc pass");
assert_eq!(second_pass.deleted_checkpoint_records, 1);
assert!(
crate::checkpoint::load_namespace_manifest_envelope(
&store,
&source,
&fork_record.manifest_object_id,
)
.await
.is_ok(),
"basis survives the pass that deletes its record"
);
let third_pass = gc_namespace(&store, &source, &config(), &aged_out)
.await
.expect("third gc pass");
assert!(third_pass.deleted_manifests >= 1);
assert!(crate::checkpoint::load_namespace_manifest_envelope(
&store,
&source,
&fork_record.manifest_object_id,
)
.await
.is_err());
stat_root(&store, &source).await;
}
#[tokio::test]
async fn gc_never_releases_a_fork_record_while_its_target_lives() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let source = NamespaceId::parse("source").expect("namespace id");
let clone = NamespaceId::parse("clone").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &source, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &source, "/docs/one.txt", "gc-one", &setup).await;
fork_namespace(&store, &source, &clone, &setup)
.await
.expect("fork");
let fork_record = read_fork_record(&store, &source).await;
assert!(
fork_record.expires_at_ms.is_some(),
"a fork record carries the attempt's lease"
);
write_test_file(&store, &source, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &source, &setup)
.await
.expect("advance root past the fork basis");
let lease = fork_record.expires_at_ms.expect("lease");
for now_ms in [
now_after_newest_object(&store, &source, GRACE_MS + 1).await,
lease,
lease + FORK_CHECKPOINT_LEASE_MS,
u64::MAX / 2,
] {
let report = gc_namespace(&store, &source, &config(), &context(now_ms))
.await
.expect("gc pass with a live target");
assert_eq!(report.released_fork_checkpoints, 0, "at {now_ms}");
assert_eq!(report.released_expired_checkpoints, 0, "at {now_ms}");
assert_eq!(
checkpoint_lifecycle(&store, &source, &fork_record.checkpoint_id).await,
CheckpointRecordLifecycle::Active {},
"a live target keeps its pin at {now_ms}"
);
}
assert!(crate::checkpoint::load_namespace_manifest_envelope(
&store,
&source,
&fork_record.manifest_object_id,
)
.await
.is_ok());
load_metadata_view(&store, &clone, ReadLoadContext::latest())
.await
.expect("target readable after every pass")
.resolve_path("/docs/one.txt")
.await
.expect("forked file readable");
}
#[tokio::test]
async fn gc_releases_abandoned_fork_checkpoints_once_the_lease_expires() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let source = NamespaceId::parse("source").expect("namespace id");
let clone = NamespaceId::parse("clone").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &source, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &source, "/docs/one.txt", "gc-one", &setup).await;
let tight = GcConfig {
grace_window_ms: GC_MIN_GRACE_WINDOW_MS,
max_objects: None,
cursor: None,
};
let attempt = context(now_after_newest_object(&store, &source, 0).await);
let lease = attempt.now_ms + FORK_CHECKPOINT_LEASE_MS;
let abandoned = crate::checkpoint::create_checkpoint(
&store,
&source,
CheckpointOwner::Fork {
target_namespace_id: clone.clone(),
},
Some(lease),
&attempt,
)
.await
.expect("leased fork record");
let fork_record = read_fork_record(&store, &source).await;
write_test_file(&store, &source, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &source, &setup)
.await
.expect("advance root past the abandoned basis");
assert!(
lease - 1 > attempt.now_ms + tight.grace_window_ms,
"the second clock below is past the grace window and still inside the lease"
);
for now_ms in [attempt.now_ms + tight.grace_window_ms + 1, lease - 1] {
let report = gc_namespace(&store, &source, &tight, &context(now_ms))
.await
.expect("gc inside the lease");
assert_eq!(report.released_fork_checkpoints, 0, "at {now_ms}");
assert_eq!(
checkpoint_lifecycle(&store, &source, &abandoned.checkpoint_id).await,
CheckpointRecordLifecycle::Active {}
);
assert!(crate::checkpoint::load_namespace_manifest_envelope(
&store,
&source,
&fork_record.manifest_object_id,
)
.await
.is_ok());
}
let expired = context(lease);
let report = gc_namespace(&store, &source, &tight, &expired)
.await
.expect("gc past the lease");
assert_eq!(report.released_fork_checkpoints, 1);
assert_eq!(
checkpoint_lifecycle(&store, &source, &abandoned.checkpoint_id).await,
CheckpointRecordLifecycle::Released {
released_at_ms: expired.now_ms
}
);
let aged_out = context(expired.now_ms + tight.grace_window_ms);
let reaping = gc_namespace(&store, &source, &tight, &aged_out)
.await
.expect("gc past the release grace window");
assert_eq!(reaping.deleted_checkpoint_records, 1);
stat_root(&store, &source).await;
}
#[tokio::test]
async fn a_fork_retry_after_abandonment_takes_a_record_of_its_own() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let source = NamespaceId::parse("source").expect("namespace id");
let clone = NamespaceId::parse("clone").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &source, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &source, "/docs/one.txt", "gc-one", &setup).await;
let abandoned = crate::checkpoint::create_checkpoint(
&store,
&source,
CheckpointOwner::Fork {
target_namespace_id: clone.clone(),
},
Some(setup.now_ms + FORK_CHECKPOINT_LEASE_MS),
&setup,
)
.await
.expect("leased fork record from the attempt that died");
fork_namespace(&store, &source, &clone, &setup)
.await
.expect("fork retry after abandonment");
let retry = store
.list_prefix(&checkpoint_prefix(source.as_str()))
.await
.expect("list checkpoints")
.len();
assert_eq!(retry, 2, "the retry pins for itself instead of reusing");
assert_eq!(
checkpoint_lifecycle(&store, &source, &abandoned.checkpoint_id).await,
CheckpointRecordLifecycle::Active {},
"the abandoned record is untouched; its lease ends it"
);
load_metadata_view(&store, &clone, ReadLoadContext::latest())
.await
.expect("target readable after retry")
.resolve_path("/docs/one.txt")
.await
.expect("forked file readable");
}
#[tokio::test]
async fn gc_retains_unreadable_checkpoint_records_and_degrades() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("checkpoint");
for key in store
.list_prefix(&checkpoint_prefix(namespace_id.as_str()))
.await
.expect("list checkpoints")
{
store
.put_overwrite(&key, bytes::Bytes::from_static(b"not json"))
.await
.expect("corrupt record");
}
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert!(report.degraded_retention);
assert_eq!(report.deleted_checkpoint_records, 0);
assert_eq!(report.deleted_manifests, 0);
assert_eq!(report.deleted_metadata_tables, 0);
assert!(
!store
.list_prefix(&checkpoint_prefix(namespace_id.as_str()))
.await
.expect("list checkpoints")
.is_empty(),
"unreadable record retained"
);
}
#[tokio::test]
async fn gc_retains_everything_without_provider_timestamps() {
let temp_dir = tempdir().expect("tempdir");
let store = MetadataMapStore::without_last_modified(
LocalFsStore::new(temp_dir.path()).expect("store"),
KeyPredicate::any(),
);
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("checkpoint");
advance_retention_floor(&store, &namespace_id, &setup)
.await
.expect("advance floor");
let aged = context(now_after_newest_object(store.inner(), &namespace_id, GRACE_MS + 1).await);
let report = gc_namespace(&store, &namespace_id, &config(), &aged)
.await
.expect("gc pass");
assert_eq!(report.deleted_wal_segments, 0);
assert_eq!(report.deleted_metadata_tables, 0);
assert_eq!(report.deleted_manifests, 0);
assert_eq!(report.deleted_checkpoint_records, 0);
assert_eq!(report.released_fork_checkpoints, 0);
assert!(report.retained_candidates > 0);
stat_root(&store, &namespace_id).await;
}
#[tokio::test]
async fn gc_sweep_reverification_chunks_preserve_outcomes() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &namespace_id, "/docs/one.txt", "gc-one", &setup).await;
let first = create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("first checkpoint");
write_test_file(&store, &namespace_id, "/docs/two.txt", "gc-two", &setup).await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("second checkpoint");
release_checkpoint_record(&store, &namespace_id, &first.checkpoint_id, setup.now_ms)
.await
.expect("mark first dead");
let aged = context(now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await);
let first_pass = gc_namespace_with_reverify_chunk(&store, &namespace_id, &config(), &aged, 1)
.await
.expect("first gc pass");
assert_eq!(first_pass.deleted_checkpoint_records, 1);
assert!(!first_pass.degraded_retention);
let second_pass = gc_namespace_with_reverify_chunk(&store, &namespace_id, &config(), &aged, 1)
.await
.expect("second gc pass");
assert!(second_pass.deleted_manifests >= 1);
stat_root(&store, &namespace_id).await;
}
#[tokio::test]
async fn gc_of_an_absent_namespace_lists_and_deletes_nothing() {
let temp_dir = tempdir().expect("tempdir");
let inner = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("orphan").expect("namespace id");
let store = IncompleteGcAccountingStore {
inner,
deletes: AtomicUsize::new(0),
lists: AtomicUsize::new(0),
};
let report = gc_namespace(&store, &namespace_id, &config(), &context(u64::MAX))
.await
.expect("gc absent namespace");
assert_eq!(report, GcResponse::empty(namespace_id.clone()));
assert_eq!(store.lists.load(Ordering::SeqCst), 0);
assert_eq!(store.deletes.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn gc_degrades_to_retention_when_a_pin_checkpoint_is_unreadable() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let source = NamespaceId::parse("source").expect("namespace id");
let clone = NamespaceId::parse("clone").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &source, &setup, false)
.await
.expect("bootstrap");
write_test_file(&store, &source, "/docs/one.txt", "gc-one", &setup).await;
fork_namespace(&store, &source, &clone, &setup)
.await
.expect("fork");
for key in store
.list_prefix(&loonfs_objectstore::keys::checkpoint_prefix(
source.as_str(),
))
.await
.expect("list checkpoints")
{
store
.put_overwrite(&key, bytes::Bytes::from_static(b"not json"))
.await
.expect("corrupt record");
}
let aged = context(now_after_newest_object(&store, &source, GRACE_MS + 1).await);
let report = gc_namespace(&store, &source, &config(), &aged)
.await
.expect("gc pass");
assert!(report.degraded_retention);
assert_eq!(report.deleted_manifests, 0);
assert_eq!(report.deleted_metadata_tables, 0);
}
async fn add_bounded_gc_fixture(
store: &LocalFsStore,
namespace_id: &NamespaceId,
setup: &MutationContext,
) {
bootstrap_namespace(store, namespace_id, setup, false)
.await
.expect("bootstrap");
let mut checkpoints = Vec::new();
for index in 0..6 {
write_test_file(
store,
namespace_id,
&format!("/docs/{index}.txt"),
&format!("bounded-gc-{index}"),
setup,
)
.await;
checkpoints.push(
create_checkpoint(store, namespace_id, setup)
.await
.expect("checkpoint"),
);
}
for checkpoint in &checkpoints[..checkpoints.len() - 1] {
release_checkpoint_record(store, namespace_id, &checkpoint.checkpoint_id, setup.now_ms)
.await
.expect("release checkpoint");
}
advance_retention_floor(store, namespace_id, setup)
.await
.expect("advance floor");
for index in 0..6 {
for key in [
wal_segment(
namespace_id.as_str(),
&format!("00000000000000000000-orphan-{index:02}"),
),
metadata_table(namespace_id.as_str(), &format!("000-orphan-{index:02}")),
format!(
"{}000-orphan-{index:02}.manifest.json",
metadata_manifest_prefix(namespace_id.as_str())
),
] {
store
.put_if_absent(&key, Bytes::from_static(b"orphan"))
.await
.expect("write orphan");
}
}
write_upload_session(store, namespace_id).await;
}
fn copy_tree(source: &std::path::Path, target: &std::path::Path) {
std::fs::create_dir_all(target).expect("create copied store directory");
for entry in std::fs::read_dir(source).expect("read source store") {
let entry = entry.expect("read source entry");
let source_path = entry.path();
let target_path = target.join(entry.file_name());
if entry.file_type().expect("read source file type").is_dir() {
copy_tree(&source_path, &target_path);
} else {
std::fs::copy(&source_path, &target_path).expect("copy store object");
}
}
}
async fn namespace_keys(store: &LocalFsStore, namespace_id: &NamespaceId) -> BTreeSet<String> {
store
.list_prefix(&loonfs_objectstore::keys::namespace_prefix(namespace_id))
.await
.expect("list namespace")
.into_iter()
.collect()
}
fn accumulate_report(total: &mut GcResponse, pass: &GcResponse) {
total.deleted_wal_segments += pass.deleted_wal_segments;
total.deleted_metadata_tables += pass.deleted_metadata_tables;
total.deleted_manifests += pass.deleted_manifests;
total.deleted_checkpoint_records += pass.deleted_checkpoint_records;
total.released_fork_checkpoints += pass.released_fork_checkpoints;
total.released_expired_checkpoints += pass.released_expired_checkpoints;
total.deleted_upload_sessions += pass.deleted_upload_sessions;
total.deleted_content_objects += pass.deleted_content_objects;
total.released_missing_basis_checkpoints += pass.released_missing_basis_checkpoints;
total.retained_candidates += pass.retained_candidates;
total.retained.add(&pass.retained);
total.degraded_retention |= pass.degraded_retention;
total.content_reclamation_deferred |= pass.content_reclamation_deferred;
total.next_reclamation_at_ms = match (total.next_reclamation_at_ms, pass.next_reclamation_at_ms)
{
(Some(a), Some(b)) => Some(a.min(b)),
(a, b) => a.or(b),
};
}
#[tokio::test]
async fn bounded_passes_delete_exactly_the_unbounded_pass_set() {
let temp_dir = tempdir().expect("tempdir");
let unbounded_root = temp_dir.path().join("unbounded");
let bounded_root = temp_dir.path().join("bounded");
let unbounded_store = LocalFsStore::new(&unbounded_root).expect("unbounded store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
add_bounded_gc_fixture(&unbounded_store, &namespace_id, &setup).await;
copy_tree(&unbounded_root, &bounded_root);
let bounded_store = LocalFsStore::new(&bounded_root).expect("bounded store");
let unbounded_now = now_after_newest_object(
&unbounded_store,
&namespace_id,
UPLOAD_SESSION_LEASE_MS + 2 * GRACE_MS + 1,
)
.await;
let unbounded_report = gc_namespace(
&unbounded_store,
&namespace_id,
&config(),
&context(unbounded_now),
)
.await
.expect("unbounded pass");
let bounded_now = now_after_newest_object(
&bounded_store,
&namespace_id,
UPLOAD_SESSION_LEASE_MS + 2 * GRACE_MS + 1,
)
.await;
let mut bounded_config = config();
bounded_config.max_objects = Some(3);
let mut bounded_report = GcResponse::empty(namespace_id.clone());
let mut passes = 0;
loop {
let pass = gc_namespace(
&bounded_store,
&namespace_id,
&bounded_config,
&context(bounded_now),
)
.await
.expect("bounded pass");
passes += 1;
accumulate_report(&mut bounded_report, &pass);
let Some(cursor) = pass.next_cursor else {
break;
};
bounded_config.cursor = Some(cursor);
}
assert!(passes > 5, "fixture should require substantial resumption");
assert_eq!(
namespace_keys(&bounded_store, &namespace_id).await,
namespace_keys(&unbounded_store, &namespace_id).await
);
assert_eq!(
(
bounded_report.deleted_wal_segments,
bounded_report.deleted_metadata_tables,
bounded_report.deleted_manifests,
bounded_report.deleted_checkpoint_records,
bounded_report.deleted_upload_sessions,
bounded_report.deleted_content_objects,
),
(
unbounded_report.deleted_wal_segments,
unbounded_report.deleted_metadata_tables,
unbounded_report.deleted_manifests,
unbounded_report.deleted_checkpoint_records,
unbounded_report.deleted_upload_sessions,
unbounded_report.deleted_content_objects,
)
);
}
#[tokio::test]
async fn budget_caps_candidate_operations_and_cursor_resumes_mid_family() {
let temp_dir = tempdir().expect("tempdir");
let inner = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&inner, &namespace_id, &setup, false)
.await
.expect("bootstrap");
let orphan_keys: Vec<String> = (0..5)
.map(|index| {
wal_segment(
namespace_id.as_str(),
&format!("00000000000000000000-orphan-{index:02}"),
)
})
.collect();
for key in &orphan_keys {
inner
.put_if_absent(key, Bytes::from_static(b"orphan"))
.await
.expect("write orphan");
}
let aged = context(now_after_newest_object(&inner, &namespace_id, GRACE_MS + 1).await);
let wal_prefix = wal_segment_prefix(namespace_id.as_str());
let store = CountingStore::new(inner, KeyPredicate::prefix(wal_prefix));
let mut bounded = config();
bounded.max_objects = Some(2);
let first = gc_namespace(&store, &namespace_id, &bounded, &aged)
.await
.expect("first bounded pass");
assert_eq!(first.deleted_wal_segments, 2);
assert!(first.next_cursor.is_some());
assert_eq!(store.snapshot().heads, 2);
assert_eq!(store.snapshot().deletes, 2);
for key in &orphan_keys[..2] {
assert!(store.head(key).await.expect("head orphan").is_none());
}
assert!(store
.head(&orphan_keys[2])
.await
.expect("head next orphan")
.is_some());
bounded.cursor = first.next_cursor;
store.reset();
let second = gc_namespace(&store, &namespace_id, &bounded, &aged)
.await
.expect("second bounded pass");
assert_eq!(second.deleted_wal_segments, 2);
assert!(second.next_cursor.is_some());
assert_eq!(store.snapshot().heads, 2);
assert_eq!(store.snapshot().deletes, 2);
for key in &orphan_keys[..4] {
assert!(store.head(key).await.expect("head orphan").is_none());
}
bounded.cursor = second.next_cursor;
loop {
store.reset();
let pass = gc_namespace(&store, &namespace_id, &bounded, &aged)
.await
.expect("remaining bounded pass");
assert!(store.snapshot().heads <= 2);
assert!(store.snapshot().deletes <= 2);
let Some(cursor) = pass.next_cursor else {
break;
};
bounded.cursor = Some(cursor);
}
for key in &orphan_keys {
assert!(store.head(key).await.expect("head orphan").is_none());
}
}
#[tokio::test]
async fn stale_cursor_rebuilds_roots_before_resuming() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("namespace id");
let setup = context(1_000);
bootstrap_namespace(&store, &namespace_id, &setup, false)
.await
.expect("bootstrap");
for index in 0..2 {
let key = wal_segment(
namespace_id.as_str(),
&format!("00000000000000000000-orphan-{index:02}"),
);
store
.put_if_absent(&key, Bytes::from_static(b"orphan"))
.await
.expect("write orphan");
}
let first_now = now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await;
let mut bounded = config();
bounded.max_objects = Some(1);
let first = gc_namespace(&store, &namespace_id, &bounded, &context(first_now))
.await
.expect("first bounded pass");
let cursor = first.next_cursor.expect("work remains");
write_test_file(
&store,
&namespace_id,
"/docs/new.txt",
"stale-cursor-new-wal",
&setup,
)
.await;
create_checkpoint(&store, &namespace_id, &setup)
.await
.expect("new checkpoint");
let resume_now = now_after_newest_object(&store, &namespace_id, GRACE_MS + 1).await;
let resume_context = context(resume_now);
let live = collect_live_set(&store, &namespace_id, &resume_context)
.await
.expect("collect advanced live set");
let mut resume = config();
resume.cursor = Some(cursor);
gc_namespace(&store, &namespace_id, &resume, &resume_context)
.await
.expect("resume stale cursor");
for key in live
.wal_segments
.iter()
.chain(live.tables.iter())
.chain(live.checkpoint_keys.iter())
{
assert!(
store.head(key).await.expect("head live object").is_some(),
"live object `{key}` must survive stale-cursor resumption"
);
}
for manifest_object_id in live.manifests {
let key = metadata_manifest_object(namespace_id.as_str(), &manifest_object_id);
assert!(
store
.head(&key)
.await
.expect("head live manifest")
.is_some(),
"live manifest `{key}` must survive stale-cursor resumption"
);
}
stat_root(&store, &namespace_id).await;
}