diff --git a/docs/en/antalya/cas/architecture/garbage-collection.md b/docs/en/antalya/cas/architecture/garbage-collection.md index 21953307f9d1..539750f87ea1 100644 --- a/docs/en/antalya/cas/architecture/garbage-collection.md +++ b/docs/en/antalya/cas/architecture/garbage-collection.md @@ -57,7 +57,7 @@ follower or a deferred round execution returns before that commit. | 13 | `round_commit` | fold | Retention-prune old generations, then publish the single `gc/state` `CAS` that adopts the whole round | | 14 | `handoff_reclaim` | post-`CAS` | Reclaim a generation a ref moved off during this round, which the ordinary retention prune already skipped and will not revisit | | 15 | `manifest_deletes` | post-`CAS` | Delete manifest bodies whose owner-removal minus-one edge the `CAS` in phase 13 just adopted | -| 16 | `namespace_cleanup` | leader; suppressed on `DEFER` | One bounded page of the perpetual namespace janitor, reclaiming dead-life debris | +| 16 | `namespace_cleanup` | leader; suppressed on `DEFER` | Pages of the perpetual namespace janitor under a 20 s soft budget, reclaiming dead-life debris | | 17 | `ref_object_cleanup` | post-`CAS` | Prune ref logs and snapshots once both fold coverage and a live snapshot make them safe to delete | | 18 | `orphan_sweep` | post-`CAS` | Exact-token deletion for the [orphan-manifest sweep](/antalya/cas/architecture/manifests-and-refs#orphan-sweep), after phase 13 adopted each candidate's blob-source retirements and the cursor | @@ -518,23 +518,36 @@ then picked up by the orphan-manifest sweep (phase 18). ## Phase 16 — namespace cleanup {#phase-16-namespace-cleanup} -One bounded page of the perpetual namespace janitor: deletes the physical objects of namespace lives -no longer in the catalog (dead-life debris). +Pages of the perpetual namespace janitor: deletes the physical objects of namespace lives +no longer in the catalog (dead-life debris). The phase takes the next page while the previous one +deleted something and published its cursor, until a 20 s soft budget: no page starts after it, the +page in progress finishes. After the last page of `cas/ns/` the next one starts from its beginning, so +a pass that began mid-stream also reaches the debris before its cursor; the first page that deletes +nothing ends the pass. A pool without debris costs one `LIST` per round. - **Runs on:** fold path here; also on the deferred path right after phase 4 with `suppress_destructive` forced on -- **Reads:** the durable `janitor_cursor`; one `LIST` page (≤ 1000 keys) of `cas/ns/`; a fresh - ref-catalog snapshot; `gc/state` per fence re-check -- **Writes / deletes:** exact-token `DELETE` per dead-life `_log` / `_snap` / `_ckpt` / `_files` - object; one `CAS` on the maintenance state when the page is decided -- **Safety:** each delete is under a GC fence re-check (`lease.owner` / `lease.seq`) before it and - once at the end; the incarnation segment in every key makes an old life's objects structurally +- **Reads:** per page: the durable `janitor_cursor`; one `LIST` page (≤ 1000 keys) of `cas/ns/`; a + fresh ref-catalog snapshot; `gc/state` for the fence re-check +- **Writes / deletes:** per page: batch `DELETE`s of the dead-life `_log` / `_snap` objects, which + are write-once and need no token; the page's keys are split evenly across the GC I/O pool + (`cas_gc_io_concurrency`), at most `cas_gc_bulk_delete_chunk_keys` per request, and a storage + without a batch delete gets one request per key from the same jobs; exact-token `DELETE` per + dead-life `_ckpt` / `_files` object; one `CAS` on the maintenance state when the page is decided +- **Safety:** each request is under a GC fence re-check (`lease.owner` / `lease.seq`), read once per + page and checked again before the cursor is published; the incarnation segment in every key makes an old life's objects structurally unreachable from a reborn same-name namespace, so a missed key can only leak storage, never expose it -- **Fails the round if:** nothing — the whole page is wrapped in a catch-all ("namespace janitor - skipped this round") +- **Fails the round if:** nothing — the whole phase is wrapped in a catch-all ("namespace janitor + stopped this round"). A failed batch `DELETE` leaks its keys and the cursor still advances. A failed + `LIST` resets the cursor to the stream start only on the first page of the phase; on a later page + the cursor the previous page published stays - **Observability:** phase row `namespace_cleanup`; metrics `janitor_pages`, `janitor_keys`, - `janitor_deleted`, `leaked` + `janitor_deleted` (keys of successful batch requests, absent ones included, plus exact removals), + `leaked` (failed keys of published pages; an unpublished page is listed again), `delete_requests` + (delete calls: one per exact-token delete and per batch, including a batch the storage refused + before its keys went one by one), `budget_exhausted`. Batch deletes run on the GC I/O pool, so the + row's `ProfileEvents` do not include their requests The cursor advances only when the whole page was decided under a held fence and an unambiguous catalog; under suppression it lists and classifies but deletes nothing and does not advance. @@ -930,13 +943,17 @@ backend without batch delete: the refused bulk call plus one `DELETE` per key). ### Phase 16 — namespace cleanup {#cost-phase-16} +Per page; the phase takes one page in a round without debris and more under its 20 s budget while +pages delete. + | Key | Operation | Requests | |---|---|---:| | `/gc/maintenance_state` | `GET` | 1 (durable `janitor_cursor`) | | `/cas/ns/` | `LIST` | one page | | `/cas/ref_catalog` | `GET` | 1 | -| `/gc/state` | `GET` | one per fence check | -| dead-life object | `DELETE` | one per object (plus one `HEAD` per object whose `LIST` entry carried no token) | +| `/gc/state` | `GET` | 1 (fence check) | +| dead-life `_log` / `_snap` objects | batch `DELETE` | up to `cas_gc_io_concurrency` requests for the page's keys; one per key on a storage without a batch delete | +| dead-life `_ckpt` / `_files` object | `DELETE` | one per object (plus one `HEAD` per object whose `LIST` entry carried no token) | | `/gc/maintenance_state` | `CAS` | 1 when the page is decided | ### Phase 17 — ref object cleanup {#cost-phase-17} diff --git a/docs/en/antalya/cas/architecture/namespaces.md b/docs/en/antalya/cas/architecture/namespaces.md index c38b65b0ce17..f3ba271ee4ad 100644 --- a/docs/en/antalya/cas/architecture/namespaces.md +++ b/docs/en/antalya/cas/architecture/namespaces.md @@ -151,7 +151,7 @@ proven its ref history is fully drained. | Catalog row (`cas/ref_catalog` entry) | `GC` phase 2, `pre_fold_ref_drain` | The round after the fold that sealed cleanup evidence for this life | | Part manifest bodies | Ordinary owner-removal ([phase 15](/antalya/cas/architecture/garbage-collection#the-round)) for anything that had a committed or precommit binding, the [orphan-manifest sweep](/antalya/cas/architecture/manifests-and-refs#orphan-sweep) for anything that never got that far | As each owning ref is dropped by the removal transaction itself, independent of the catalog row | | Blob bodies | The ordinary condemn/graduate/delete pipeline | Whenever the manifests that named them stop being live, same as any other blob | -| Ref stream/state objects (`_log`, `_snap`, `_ckpt`, `_files`) under the dead `life_id` | The perpetual namespace janitor ([phase 16](/antalya/cas/architecture/garbage-collection#the-round)) | Best-effort, one bounded `LIST` page at a time, whenever it next lists a key whose `life_id` a fresh catalog cut no longer names — independent of, and not gated on, catalog-row deletion | +| Ref stream/state objects (`_log`, `_snap`, `_ckpt`, `_files`) under the dead `life_id` | The perpetual namespace janitor ([phase 16](/antalya/cas/architecture/garbage-collection#the-round)) | Best-effort, in bounded `LIST` pages under a per-round time budget, whenever it next lists a key whose `life_id` a fresh catalog cut no longer names — independent of, and not gated on, catalog-row deletion | The janitor is leak-only: it never fails a round, never blocks progress on an unreadable key, and a crash mid-page simply leaves debris for its next page. diff --git a/docs/en/antalya/cas/configuration.md b/docs/en/antalya/cas/configuration.md index 17b1e3482768..bdc7a3f2f6e2 100644 --- a/docs/en/antalya/cas/configuration.md +++ b/docs/en/antalya/cas/configuration.md @@ -106,7 +106,7 @@ entirely before release. Treat this table as a snapshot of the current build, no | `cas_part_folder_cache_max_entry_bytes` | 16 MiB | Oversized part-folder views bypass retention above this size | | `cas_manifest_decode_cache_bytes` | 128 MiB | Manifest decode cache byte budget (`0` disables) | | `cas_gc_meta_pool_size` | `16` | Bounded pool size for GC per-hash freshness-meta writes | -| `cas_gc_io_concurrency` | `16` | Bounded pool size for GC object-storage requests that run in parallel: the fold's read-ahead (checkpoints, ref logs, manifests, zero-candidate HEADs), the orphan-manifest sweep planning reads, the `SYSTEM CAS GC REBUILD` read-ahead, and the `pending_deletes` blob `HEAD` + conditional `DELETE` fan-out. Not covered: meta writes (`cas_gc_meta_pool_size`) and all other GC requests, which run on the round thread. `1` runs the covered requests sequentially. `cas_gc_read_concurrency` is rejected without an alias; use `cas_gc_io_concurrency` instead | +| `cas_gc_io_concurrency` | `16` | Bounded pool size for GC object-storage requests that run in parallel: the fold's read-ahead (checkpoints, ref logs, manifests, zero-candidate HEADs), the orphan-manifest sweep planning reads, the `SYSTEM CAS GC REBUILD` read-ahead, the `pending_deletes` blob `HEAD` + conditional `DELETE` fan-out, and the namespace janitor's batch deletes of dead ref logs and snapshots. Not covered: meta writes (`cas_gc_meta_pool_size`) and all other GC requests, which run on the round thread. `1` runs the covered requests sequentially. `cas_gc_read_concurrency` is rejected without an alias; use `cas_gc_io_concurrency` instead | | `cas_attempt_timeout_ms` | `5000` | Budget for one HTTP attempt of a writable Native mount's control-plane requests (read, head, list, remove, conditional write), at least 1. Together with the connect cap it forms the attempt envelope (`cas_attempt_timeout_ms + 2 × cap`; the cap is `cas_attempt_timeout_ms` itself when the disk's `connect_timeout_ms` is `0`, else `min(connect_timeout_ms, cas_attempt_timeout_ms)`) that the lease arithmetic reserves: one TCP connect and one TLS handshake under the cap each, send/receive bounded per socket operation by `cas_attempt_timeout_ms`. With background renewal the cadence check requires `cas_mount_renew_period_ms + 2 × envelope + cas_lease_safety_margin_ms < cas_mount_lease_ttl_ms`, which puts an effective ceiling on the frozen connect cap: under the defaults (TTL 30000, period 10000, margin 2000) the envelope must stay under 9000, so a disk `connect_timeout_ms` of 2000 ms or more refuses to open writable — lower the connect timeout or raise the TTL if you hit this | | `cas_lease_safety_margin_ms` | `2000` | Startup-only margin validated against the mount lease TTL: the attempt envelope + `cas_lease_safety_margin_ms` must be strictly less than the mount lease TTL, and `cas_mount_renew_period_ms` + 2 × envelope + `cas_lease_safety_margin_ms` too, or the disk refuses to open writable | | `cas_unsafe_remount_no_delay` | `0` | Reclaim a mount slot that carries this server's own uuid at once after a hard restart, without observing the slot's token for the lease TTL. Unsafe whenever two processes can hold the same `server_uuid` (a copied uuid file, a stalled predecessor). After such a reclaim the predecessor can still start conditional writes until its own cutoff (`confirmed deadline − cas_lease_safety_margin_ms − 2 × envelope`) or until its next renewal meets the token guard, and a request it already sent may still materialize later. That is not a data hazard: ref-log keys carry `(writer_epoch, sequence)` and creates are conditional, so two writers can never commit different bodies to one key, and recovery's epoch seal settles any straggler (recovery fails closed after 64 successive seal-create attempts displaced by newly materializing old-epoch transactions). The exposure is availability, not data. Intended for test stands and deployments that guarantee one process per uuid | @@ -179,7 +179,7 @@ for the remaining caps, `0` means unbounded. |---|---|---|---| | `cas_manifest_sweep_list_budget_keys` | `1000` | `UInt64` | Orphan-manifest sweep `LIST` budget per round | | `cas_manifest_sweep_delete_budget_keys` | `100` | `UInt64` | Orphan-manifest sweep `DELETE` budget per round | -| `cas_gc_bulk_delete_chunk_keys` | `1000` | `1`–`1000` | Keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots) | +| `cas_gc_bulk_delete_chunk_keys` | `1000` | `1`–`1000` | Most keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots, the namespace janitor's dead ref logs and snapshots) | | `cas_gc_round_graduation_budget` | `5000` | `0` = unbounded | Blob-graduation (`condemned` → `delete_pending`) cohort cap per round | | `cas_gc_round_redelete_budget` | `5000` | `0` = unbounded | Exact-token re-delete cohort cap for prior `delete_pending` rows per round | | `cas_gc_round_sweep_namespace_budget` | `20` | `0` = unbounded | Distinct namespaces per orphan-manifest sweep page whose protection view may be built | diff --git a/docs/en/operations/system-tables/cas_gc_log.md b/docs/en/operations/system-tables/cas_gc_log.md index 46fcd559cfb0..53309724de68 100644 --- a/docs/en/operations/system-tables/cas_gc_log.md +++ b/docs/en/operations/system-tables/cas_gc_log.md @@ -86,7 +86,7 @@ The phases, in execution order: | `round_commit` | The generation-retention prune and the round's single `gc/state` compare-and-swap. | prune `LIST`s and deletes, one compare-and-swap | | `handoff_reclaim` | Wholesale-reclaim generations a moved run ref stranded below the retention cursor. | prefix `LIST`s and deletes | | `manifest_deletes` | Exact-token deletes of owner-removed manifest bodies, after their decrements were adopted. | one `DELETE` per body | -| `namespace_cleanup` | Run one bounded `cas/ns/` page across the stream and state subtrees for the perpetual dead-life janitor. This phase is physical reclamation, not a lifecycle gate. | one namespace-root page `LIST`, catalog cut, exact-token deletes | +| `namespace_cleanup` | Run `cas/ns/` pages across the stream and state subtrees for the perpetual dead-life janitor, under a 20 s soft budget; `budget_exhausted` is 1 when the budget stopped the phase, `delete_requests` counts the phase's delete calls, which run on the GC I/O pool and are missing from this row's `ProfileEvents`. This phase is physical reclamation, not a lifecycle gate. | per page: one namespace-root page `LIST`, catalog cut, batch deletes of dead ref logs and snapshots, exact-token deletes of the rest | | `ref_object_cleanup` | Delete ref logs covered by both the durable fold cursor and a durable snapshot, plus superseded snapshots. | one `HEAD` + one `DELETE` per deletable object | | `orphan_sweep` | The budgeted, cursor-paced orphan part-manifest backstop. | budgeted `LIST` and deletes | diff --git a/src/Common/ProfileEvents.cpp b/src/Common/ProfileEvents.cpp index a9986b5aa485..3ef8a1752aa0 100644 --- a/src/Common/ProfileEvents.cpp +++ b/src/Common/ProfileEvents.cpp @@ -927,7 +927,7 @@ The server successfully detected this situation and will download merged part fr M(CASGCRefWalkPlansBuilt, "Number of complete catalog-authoritative CAS ref walk plans constructed by ordinary GC and rebuild. A regular or rebuilding invocation that reaches the post-LIST catalog cut increments this exactly once, including a round that later defers.", ValueType::Number) \ M(CASGCUnmatchedAdoptedParentLives, "Number of adopted-parent CAS ref-life rows dropped because the post-LIST catalog cut has no matching physical life. Each occurrence is inert for planning and suppression and is logged with its exact physical life id; a persistent nonzero rate indicates old generation state is outliving catalog removal.", ValueType::Number) \ M(CASGCStuckRemovals, "Number of adopted CAS GC rounds that observed a Removing namespace at or beyond the diagnostic age threshold without terminal cleanup evidence. Incremented and warned every such round; diagnostic only, with no effect on folding, suppression, appends, or deletion.", ValueType::Number) \ - M(CASGCNamespaceCleanupLeaks, "Number of dead-life namespace objects whose reclamation the perpetual janitor could not confirm because HEAD or exact-delete failed. Each occurrence is logged with the exact key and remains leak-only: it neither suppresses destructive GC nor blocks catalog lifecycle progress.", ValueType::Number) \ + M(CASGCNamespaceCleanupLeaks, "Number of dead-life namespace objects whose reclamation the perpetual janitor could not confirm because HEAD, exact-delete or a batch delete failed. Each occurrence is logged with the key, or the first key of the batch, and remains leak-only: it neither suppresses destructive GC nor blocks catalog lifecycle progress.", ValueType::Number) \ M(CASDetachedWorkDrainTimeouts, "Counts CAS storage teardowns whose bounded wait for detached background work expired with work still in flight. The teardown proceeds, but for that teardown it could not be established that no tracked task still holds the pool. Expected to stay at zero.", ValueType::Number) \ M(CASEventDroppedContextExpired, "Number of CAS system-log events dropped because the storage's `Context` reference expired before delivery. A non-zero value indicates event production outlived the owning server context.", ValueType::Number) \ M(CASGCUnappliedFoldedTransactions, "Number of ref transactions a GC round folded and merged but whose blob deltas never reached a shard reducer. Always 0 on a healthy round; a nonzero value fails the round closed, because the round would otherwise advance its fold cursor past a transaction it never applied.", ValueType::Number) \ diff --git a/src/Common/setThreadName.h b/src/Common/setThreadName.h index 8a0222dc9639..b34b48cc5171 100644 --- a/src/Common/setThreadName.h +++ b/src/Common/setThreadName.h @@ -35,6 +35,7 @@ namespace DB M(CAS_ANOMALY_DIAG, "CasAnomalyDiag") \ M(CAS_GC_HEARTBEAT, "CasGcHeartbeat") \ M(CAS_GC_REDELETE, "CasGcRedelete") \ + M(CAS_GC_JANITOR, "CasGcJanitor") \ M(CAS_GC_SCHEDULER, "CasGcSched") \ M(CAS_LEASE_RENEWER, "CasLeaseRenewer") \ M(CAS_REF_SNAPSHOT_PUBLISH, "CasRefSnapPub") \ diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp index 7331d2bad29a..73f8380c668d 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp @@ -79,8 +79,8 @@ constexpr std::string_view CAS_KEY_PREFIX = "cas_"; DECLARE(UInt64, part_folder_cache_max_entry_bytes, 16ULL << 20, "Oversized part-folder views bypass retention above this size", 0) \ DECLARE(UInt64, manifest_decode_cache_bytes, 128ULL << 20, "Manifest DECODE cache byte budget (0 disables)", 0) \ DECLARE(UInt64, gc_meta_pool_size, 16, "Bounded pool size for GC per-hash freshness-meta writes", 0) \ - DECLARE(UInt64, gc_io_concurrency, 16, "Maximum number of threads in the GC I/O pool. Used for fold and rebuild read-ahead, orphan-manifest sweep planning reads, and pending_deletes HEAD plus conditional DELETE. Per-hash meta writes use gc_meta_pool_size; other GC requests run on the round thread. 1 disables parallel GC I/O", 0) \ - DECLARE(UInt64, gc_bulk_delete_chunk_keys, 1000, "Keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots); 1 to 1000", 0) \ + DECLARE(UInt64, gc_io_concurrency, 16, "Maximum number of threads in the GC I/O pool. Used for fold and rebuild read-ahead, orphan-manifest sweep planning reads, pending_deletes HEAD plus conditional DELETE, and namespace janitor batch deletes. Per-hash meta writes use gc_meta_pool_size; other GC requests run on the round thread. 1 disables parallel GC I/O", 0) \ + DECLARE(UInt64, gc_bulk_delete_chunk_keys, 1000, "Most keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots, the namespace janitor's dead ref logs and snapshots); 1 to 1000", 0) \ DECLARE(UInt64, attempt_timeout_ms, 5000, "Budget for one HTTP attempt of a writable Native mount's control-plane requests (read, head, list, remove, conditional write), at least 1. With the connect cap it forms the attempt envelope the lease arithmetic reserves", 0) \ DECLARE(UInt64, lease_safety_margin_ms, 2000, "Startup-only margin validated against the mount lease TTL: attempt envelope + this must be strictly less than the TTL, and renew period + 2 × envelope + this too", 0) \ DECLARE(String, staging_backend, "local", "Blob staging backend (local | s3); s3 is opt-in", 0) \ diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Formats/CasLayout.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Formats/CasLayout.h index 5ab57bb7854b..856b623b99b6 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Formats/CasLayout.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Formats/CasLayout.h @@ -158,7 +158,7 @@ class Layout /// not try an uncompressed variant. String refLogKey(const NamespaceLifeId & ns_id, const RefTxnId & id) const { - return namespaceStreamPrefix(ns_id) + "_log/" + renderRefTxnId(id) + String(storedSuffix(FormatId::RefLog)); + return refStreamObjectKey(ns_id.incarnation, RefObjectKind::Log, id); } /// Writer-published table snapshot at `.../_snap/.zst`. The snapshot @@ -166,7 +166,7 @@ class Layout /// `X` reuses the `RefTxnId` of the last log it covers. String refSnapshotKey(const NamespaceLifeId & ns_id, const RefTxnId & id) const { - return namespaceStreamPrefix(ns_id) + "_snap/" + renderRefTxnId(id) + String(storedSuffix(FormatId::RefSnapshot)); + return refStreamObjectKey(ns_id.incarnation, RefObjectKind::Snap, id); } /// The write-once forms of `refLogKey` and `refSnapshotKey`: the same strings, typed as keys a @@ -181,6 +181,17 @@ class Layout return WriteOnceKey(refSnapshotKey(ns_id, id)); } + /// The write-once form of a listed `_log`/`_snap` key that `parseRefObjectKey` accepted. `std::nullopt` + /// when `listed_key` is not the canonical key of the parsed identity: a precondition-free delete must + /// name only keys this layout writes. + std::optional writeOnceStreamKey(const ParsedRefObjectKey & parsed, std::string_view listed_key) const + { + String canonical = refStreamObjectKey(parsed.life_id, parsed.kind, parsed.txn_id); + if (canonical != listed_key) + return std::nullopt; + return WriteOnceKey(std::move(canonical)); + } + /// The life's checkpoint object (spec INV-4) at `/cas/ns/state//_ckpt`. Unlike /// immutable stream objects it is mutable (token-CAS), carries no transaction id, and therefore lives /// in the point/path-addressed state tree rather than a `_log`/`_snap` directory -- @@ -486,6 +497,13 @@ class Layout /// Parses the one physical-id segment after the rest of a life-owned key identified its family. NamespaceLifePhysicalId namespaceLifePhysicalIdOf(std::string_view key, std::string_view segment) const; + String refStreamObjectKey(NamespaceLifePhysicalId life_id, RefObjectKind kind, const RefTxnId & id) const + { + const bool log = kind == RefObjectKind::Log; + return namespaceStreamRootPrefix() + renderIncarnation(life_id) + (log ? "/_log/" : "/_snap/") + renderRefTxnId(id) + + String(storedSuffix(log ? FormatId::RefLog : FormatId::RefSnapshot)); + } + /// Build ///. /// Throws BAD_ARGUMENTS if id is shorter than 2 characters. String shardedKey(const String & ns, const String & id) const diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp index f871182c8edb..2debbc2787e7 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp @@ -85,6 +85,9 @@ namespace DB::Cas namespace { +/// Soft limit on the namespace janitor's phase: no page starts after it, the page in progress finishes. +constexpr uint64_t kJanitorBudgetMs = 20'000; + /// The `on_page_fetched` hook GC passes to every `forEachListedKey`/`recoverRefTable` /// call it owns (never passed by fsck/offline-repair callers of those shared helpers) -- one increment /// per physical LIST page, never per listed key. @@ -344,13 +347,14 @@ Gc::Gc(PoolPtr store_, UInt128 gc_id_, std::function now_ms_fn_, /// `store->poolConfig()` AFTER the null check above. meta_writer = std::make_unique( store, logger, static_cast(store->poolConfig().gc_meta_pool_size)); - /// The GC I/O pool (read-ahead and the `pending_deletes` fan-out), built here for the same reason. + /// The GC I/O pool (read-ahead, the `pending_deletes` fan-out and the namespace janitor's batch deletes), + /// built here for the same reason. /// The queue is UNBOUNDED because the hinting sites throttle themselves against /// `GcReadAhead::window`; a bounded queue would only move the throttle into /// `scheduleOrThrowOnError`, blocking the round thread instead of the hint loop that already knows /// how much it wants in flight. const size_t io_concurrency = std::max(1, store->poolConfig().gc_io_concurrency); - /// Keep `shutdown_on_exception` false: every task on this pool (read-ahead and re-delete) is submitted + /// Keep `shutdown_on_exception` false: every task on this pool (read-ahead, re-delete, janitor delete) is submitted /// through a callback runner that stores its exception in the task's future, and no caller invokes /// `wait` on this pool. io_pool = std::make_unique( @@ -360,34 +364,126 @@ Gc::Gc(PoolPtr store_, UInt128 gc_id_, std::function now_ms_fn_, io_pool_refuse_at_for_test = store->poolConfig().gc_io_pool_refuse_at_for_test; } -void Gc::runNamespaceJanitorPage( +void Gc::runNamespaceJanitor( const GcState & leased_state, bool suppress_destructive, uint64_t cleanup_evidence_rows) { GcPhaseTimer t(phase_sink, "namespace_cleanup"); t.metric("evidence_rows", cleanup_evidence_rows); NamespaceJanitorResult janitor_result; + bool budget_exhausted = false; try { CasRequests & requests = store->openRequests(); const Layout & layout = store->layout(); - NamespaceJanitor janitor(requests, layout, 1000); - /// ONE authority read per page, made here rather than from the predicate: the janitor's - /// operation samples its liveness before every request, and a page walks up to a thousand keys. - refreshAuthority(leased_state.lease.seq); - janitor_result = janitor.runOnePage(suppress_destructive, [this] { return authority_held; }); - for (const String & anomaly : janitor_result.anomalies) - LOG_WARNING(logger, "CAS namespace janitor: {}", anomaly); - if (janitor_result.leaked) - ProfileEvents::increment(ProfileEvents::CASGCNamespaceCleanupLeaks, janitor_result.leaked); + NamespaceJanitor janitor(requests, layout, 1000, + [this](CasOperation & op, const std::vector & keys, NamespaceJanitorResult & out) + { + removeWriteOnceOnIoPool(op, keys, out); + }); + const uint64_t deadline_ms = mono_ms_fn() + kJanitorBudgetMs; + while (true) + { + /// ONE authority read per page, made here rather than from the predicate: the janitor's + /// operation samples its liveness before every request, and a page walks up to a thousand keys. + refreshAuthority(leased_state.lease.seq); + const NamespaceJanitorResult page = janitor.runOnePage( + suppress_destructive, [this] { return authority_held; }, /*first_page_of_pass=*/janitor_result.pages == 0); + janitor_result.pages += page.pages; + janitor_result.keys += page.keys; + janitor_result.deleted += page.deleted; + janitor_result.leaked += page.leaked; + janitor_result.delete_requests += page.delete_requests; + for (const String & anomaly : page.anomalies) + LOG_WARNING(logger, "CAS namespace janitor: {}", anomaly); + if (!page.more) + break; + if (mono_ms_fn() >= deadline_ms) + { + budget_exhausted = true; + break; + } + } } catch (const std::exception & e) { - LOG_WARNING(logger, "CAS namespace janitor skipped this round: {}", e.what()); + LOG_WARNING(logger, "CAS namespace janitor stopped this round: {}", e.what()); } + if (janitor_result.leaked) + ProfileEvents::increment(ProfileEvents::CASGCNamespaceCleanupLeaks, janitor_result.leaked); t.metric("janitor_pages", janitor_result.pages); t.metric("janitor_keys", janitor_result.keys); t.metric("janitor_deleted", janitor_result.deleted); t.metric("leaked", janitor_result.leaked); + t.metric("delete_requests", janitor_result.delete_requests); + t.metric("budget_exhausted", budget_exhausted ? 1 : 0); +} + +void Gc::removeWriteOnceOnIoPool(CasOperation & op, const std::vector & keys, NamespaceJanitorResult & out) +{ + /// An even split, so a store without a batch delete sends its one-key requests from every pool thread. + const size_t concurrency = std::max(1, store->poolConfig().gc_io_concurrency); + const size_t chunk_keys = std::min( + std::clamp(store->poolConfig().gc_bulk_delete_chunk_keys, 1, kBulkDeleteMaxKeys), + (keys.size() + concurrency - 1) / concurrency); + std::vector> chunks; + for (size_t begin = 0; begin < keys.size(); begin += chunk_keys) + chunks.emplace_back(keys.begin() + begin, keys.begin() + std::min(keys.size(), begin + chunk_keys)); + + /// Written by the job of the same index, read after every job finished. + std::vector requests(chunks.size(), 0); + ThreadPoolCallbackRunnerLocal runner(*io_pool, ThreadName::CAS_GC_JANITOR); + std::vector::Task>> handles; + handles.reserve(chunks.size()); + SCOPE_EXIT_SAFE({ ThreadPoolCallbackRunnerLocal::waitForAllToFinish(handles); }); + try + { + for (size_t i = 0; i < chunks.size(); ++i) + { + if (io_pool_refuse_at_for_test == i) + { + io_pool_refuse_at_for_test.reset(); + throw Exception(ErrorCodes::CANNOT_SCHEDULE_TASK, "Injected CAS namespace janitor enqueue refusal at index {}", i); + } + /// The round thread waits for every job before it refreshes `authority_held` again. + handles.emplace_back(runner.enqueueAndGiveOwnership( + [this, chunk = &chunks[i], chunk_requests = &requests[i], admitted_generation = op.generation()] + { + CasOperation job_op = store->openRequests().resume(admitted_generation, [this] { return authority_held; }); + /// A job that fails made at least the one call. + *chunk_requests = 1; + *chunk_requests = removeChunkWriteOnceOrOneByOne(job_op, *chunk, Retry::standard()); + })); + } + } + catch (...) + { + /// Leak-only like a failed request: the keys stay for the cursor's next pass over this page. + uint64_t unscheduled = 0; + for (size_t i = handles.size(); i < chunks.size(); ++i) + unscheduled += chunks[i].size(); + out.leaked += unscheduled; + out.anomalies.push_back(fmt::format( + "leaked {} dead-life objects starting at '{}': delete job not scheduled: {}", + unscheduled, chunks[handles.size()].front().str(), getCurrentExceptionMessage(false))); + } + ThreadPoolCallbackRunnerLocal::waitForAllToFinish(handles); + + for (size_t i = 0; i < handles.size(); ++i) + { + out.delete_requests += requests[i]; + try + { + handles[i]->future.get(); + out.deleted += chunks[i].size(); + } + catch (...) + { + out.leaked += chunks[i].size(); + out.anomalies.push_back(fmt::format( + "leaked {} dead-life objects starting at '{}': batch delete failed: {}", + chunks[i].size(), chunks[i].front().str(), getCurrentExceptionMessage(false))); + } + } } uint64_t removeChunkWriteOnceOrOneByOne(CasOperation & op, const std::vector & chunk, const Retry & policy) @@ -820,7 +916,7 @@ RoundReport Gc::runRegularRound(std::function on_lease_acquired, bool al /// DEFER has no `FoldResult`, hence no complete global destructive verdict. The janitor still /// takes its bounded page and catalog cut, but suppression keeps both deletes and valid-page /// cursor progress at the same position for the bounded forced fold to retry. - runNamespaceJanitorPage(state, /*suppress_destructive=*/true, /*cleanup_evidence_rows=*/0); + runNamespaceJanitor(state, /*suppress_destructive=*/true, /*cleanup_evidence_rows=*/0); return report; /// no fold, no pre-CAS deletes, no gc/state CAS — sealed generation stays pinned } @@ -1331,7 +1427,7 @@ RoundReport Gc::runRegularRound(std::function on_lease_acquired, bool al uint64_t cleanup_evidence_rows = 0; for (const auto & [life_id, ref_life_state] : folded.fold_seal.ref_lives) cleanup_evidence_rows += ref_life_state.cleanup_evidence ? 1 : 0; - runNamespaceJanitorPage(state, suppress_destructive, cleanup_evidence_rows); + runNamespaceJanitor(state, suppress_destructive, cleanup_evidence_rows); /// PHASE 17/18 `ref_object_cleanup`. Emitted even when the whole pass is skipped (`trim_enabled` is /// a test seam, `suppressed` gates the deletes), because "this phase did nothing and why" is exactly /// what a reader of a round that reclaimed nothing needs to see. diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h index 33f2bdcf09a8..8e6c43a7c986 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h @@ -249,6 +249,7 @@ struct RefScanSummary class RefPlan; class RoundInput; +struct NamespaceJanitorResult; RefPlan buildRefWalkPlan(RoundInput && round_input); namespace tests @@ -548,13 +549,18 @@ class Gc /// It performs no physical LIST or delete. CatalogLifecycleReconcileResult drainCompletedRemoving(const GcState & leased_state); - /// Run exactly one independently paced physical namespace-maintenance page. The caller supplies + /// Run the independently paced physical namespace maintenance: pages from the persisted cursor while + /// each one deleted dead-life objects, under a soft time budget. The caller supplies /// the one round-wide destructive verdict when it exists; DEFER passes suppression because it has - /// no folded frontier verdict. This helper owns only janitor I/O and phase metrics, never lifecycle - /// transitions or the hot stream walk plan. - void runNamespaceJanitorPage( + /// no folded frontier verdict, which takes one page. This helper owns only janitor I/O and phase + /// metrics, never lifecycle transitions or the hot stream walk plan. + void runNamespaceJanitor( const GcState & leased_state, bool suppress_destructive, uint64_t cleanup_evidence_rows); + /// The janitor's `RemoveWriteOnce`: splits `keys` evenly across the GC I/O pool and waits for every job. + /// A failed job leaks its keys, and so does a job the pool refused to schedule. + void removeWriteOnceOnIoPool(CasOperation & op, const std::vector & keys, NamespaceJanitorResult & out); + void reportStuckRemovals(const RefPlan & plan, uint64_t current_round); /// What one fold produced. The blob deltas are sealed @@ -1024,7 +1030,7 @@ class Gc std::unique_ptr meta_writer; /// The GC I/O pool: the fold's and rebuild's read-ahead, the orphan-manifest sweep planning reads, - /// and the `pending_deletes` fan-out, sized by `gc_io_concurrency`. A `unique_ptr` for the same reason as `meta_writer`: the size comes from + /// the `pending_deletes` fan-out and the namespace janitor's batch deletes, sized by `gc_io_concurrency`. A `unique_ptr` for the same reason as `meta_writer`: the size comes from /// `store->poolConfig()`, which may only be read after the constructor body has validated `store`. std::unique_ptr io_pool; std::optional io_pool_refuse_at_for_test; diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.cpp index f8c153ed0f20..dba6b927c97c 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.cpp @@ -2,6 +2,8 @@ #include #include #include +#include +#include namespace DB::Cas { @@ -19,9 +21,40 @@ void throwOnRefusedOrGaveUp(WriteResult && result, std::string_view what) (void)orThrow(std::move(result), what); } +void removeOnCallerThread(CasOperation & op, const std::vector & keys, NamespaceJanitorResult & out) +{ + for (size_t begin = 0; begin < keys.size(); begin += kBulkDeleteMaxKeys) + { + const std::vector batch( + keys.begin() + begin, keys.begin() + std::min(keys.size(), begin + kBulkDeleteMaxKeys)); + ++out.delete_requests; + try + { + op.removeManyWriteOnce(batch, Retry::standard()); + out.deleted += batch.size(); + } + catch (const std::exception & e) + { + out.leaked += batch.size(); + out.anomalies.push_back(fmt::format( + "leaked {} dead-life objects starting at '{}': batch delete failed: {}", + batch.size(), batch.front().str(), e.what())); + } + } } -NamespaceJanitorResult NamespaceJanitor::runOnePage(bool suppress_deletes, Liveness liveness) +} + +NamespaceJanitor::NamespaceJanitor( + CasRequests & requests_, const Layout & layout_, size_t page_budget_, RemoveWriteOnce remove_write_once_) + : requests(requests_) + , layout(layout_) + , page_budget(page_budget_) + , remove_write_once(remove_write_once_ ? std::move(remove_write_once_) : removeOnCallerThread) +{ +} + +NamespaceJanitorResult NamespaceJanitor::runOnePage(bool suppress_deletes, Liveness liveness, bool first_page_of_pass) { NamespaceJanitorResult result; CasOperation op = requests.admit(std::move(liveness)); @@ -43,7 +76,8 @@ NamespaceJanitorResult NamespaceJanitor::runOnePage(bool suppress_deletes, Liven } catch (...) { - (void)casGcMaintenanceState(op, layout, progress.etag, GcMaintenanceState{}, Retry::once()); + if (first_page_of_pass) + (void)casGcMaintenanceState(op, layout, progress.etag, GcMaintenanceState{}, Retry::once()); throw; } result.pages = 1; @@ -70,16 +104,19 @@ NamespaceJanitorResult NamespaceJanitor::runOnePage(bool suppress_deletes, Liven /// absent objects and token mismatches are final per-key outcomes and therefore do not by /// themselves prevent progress. bool page_decided = !ambiguous && !suppress_deletes; + std::vector dead_stream_keys; for (const ListedKey & listed : page.keys) { std::optional life_id; + std::optional stream_key; try { if (listed.key.starts_with(layout.namespaceStreamRootPrefix())) { - if (const auto parsed = layout.parseRefObjectKey(listed.key)) - life_id = parsed->life_id; + stream_key = layout.parseRefObjectKey(listed.key); + if (stream_key) + life_id = stream_key->life_id; } else if (listed.key.starts_with(layout.namespaceStateRootPrefix())) { @@ -103,6 +140,15 @@ NamespaceJanitorResult NamespaceJanitor::runOnePage(bool suppress_deletes, Liven if (ambiguous || suppress_deletes || catalog_cut.life_index.resolve(*life_id)) continue; + if (stream_key) + { + if (std::optional write_once = layout.writeOnceStreamKey(*stream_key, listed.key)) + { + dead_stream_keys.push_back(std::move(*write_once)); + continue; + } + } + std::optional etag = listed.etag; if (!etag) { @@ -126,6 +172,7 @@ NamespaceJanitorResult NamespaceJanitor::runOnePage(bool suppress_deletes, Liven page_decided = false; break; } + ++result.delete_requests; try { if (op.remove(listed.key, *etag, Retry::standard()) == Removal::Removed) @@ -139,26 +186,38 @@ NamespaceJanitorResult NamespaceJanitor::runOnePage(bool suppress_deletes, Liven } } + if (page_decided && !dead_stream_keys.empty() && op.admitted()) + remove_write_once(op, dead_stream_keys, result); + /// Recheck even when the page had no dead candidate. A tenure that observes fence loss after LIST - /// or after the last exact delete must not publish progress. Loss after this check may still race + /// or after the last delete must not publish progress. Loss after this check may still race /// with the leak-only maintenance CAS; already completed exact deletes remain safe to repeat. if (page_decided && !op.admitted()) page_decided = false; + bool published = false; if (page_decided) { const GcMaintenanceState next{.janitor_cursor = page.next_cursor}; try { - const WriteResult published = casGcMaintenanceState(op, layout, progress.etag, next, Retry::standard()); - if (std::holds_alternative(published) || std::holds_alternative(published)) + const WriteResult outcome = casGcMaintenanceState(op, layout, progress.etag, next, Retry::standard()); + if (std::holds_alternative(outcome) || std::holds_alternative(outcome)) result.anomalies.push_back("cursor publication did not commit"); + else if (std::holds_alternative(outcome)) + { + published = true; + result.more = result.deleted > 0 && !(cursor.empty() && page.next_cursor.empty()); + } } catch (const std::exception & e) { result.anomalies.push_back("cursor publication failed: " + String(e.what())); } } + /// An unpublished page is listed again, so its failed keys are not leaked yet. + if (!published) + result.leaked = 0; return result; } diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.h index d23d21402c73..801d3967bf66 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasNamespaceJanitor.h @@ -13,15 +13,28 @@ struct NamespaceJanitorResult uint64_t keys = 0; uint64_t deleted = 0; uint64_t leaked = 0; + /// Delete calls the page made: one per exact-token delete and per batch, including a batch the storage + /// refused before its keys went one by one. Reissues of a call are not counted. + uint64_t delete_requests = 0; + /// The page deleted something, published its cursor and did not cover the whole stream by itself, so + /// the next page is worth taking now. After the last page of the stream the next one starts from its + /// beginning. + bool more = false; std::vector anomalies; }; -/// Runs one bounded, leak-only page over the physical namespace ownership tree. +/// Deletes one page's dead `_log`/`_snap` keys under the page's operation and adds to `out` the keys it +/// deleted, the keys it could not confirm deleted with an anomaly for each group of them, and its delete +/// calls. An exception leaves the page unpublished. +using RemoveWriteOnce + = std::function & keys, NamespaceJanitorResult & out)>; + +/// Runs one bounded, leak-only page over the physical namespace ownership tree; `Gc` repeats it under a time budget. class NamespaceJanitor { public: - NamespaceJanitor(CasRequests & requests_, const Layout & layout_, size_t page_budget_) - : requests(requests_), layout(layout_), page_budget(page_budget_) {} + /// An empty `remove_write_once_` deletes on the caller's thread, one request per `kBulkDeleteMaxKeys` keys. + NamespaceJanitor(CasRequests & requests_, const Layout & layout_, size_t page_budget_, RemoveWriteOnce remove_write_once_ = {}); /// `liveness` is admitted once for the whole page (one `CasOperation` covers the read, the list, /// every delete and the cursor publication): a fact the fence cannot see, such as "this tenure @@ -31,12 +44,22 @@ class NamespaceJanitor /// that returns false ends whichever request was about to be sent: a read verb (the maintenance /// read, the list, a HEAD) throws out of this call, and a write verb (a delete, the cursor /// publication) reports it as `GaveUp` rather than sending anything. - NamespaceJanitorResult runOnePage(bool suppress_deletes, Liveness liveness); + /// + /// A dead life's canonical `_log`/`_snap` keys are write-once, so the page deletes them together through + /// `remove_write_once` with no per-key `HEAD` or token; `_ckpt` and `_files` keep the exact-token delete. + /// + /// A failed `LIST` resets the persisted cursor only when `first_page_of_pass`: that cursor comes from an + /// earlier round and the store may reject it forever, while the one a page of this pass just published + /// is fresh, so the failure is transient and the next round resumes from it. + /// + /// `leaked` is zero for a page that did not publish its cursor: the page is listed again. + NamespaceJanitorResult runOnePage(bool suppress_deletes, Liveness liveness, bool first_page_of_pass = true); private: CasRequests & requests; const Layout & layout; size_t page_budget; + RemoveWriteOnce remove_write_once; }; } diff --git a/src/Disks/tests/gtest_cas_gc_frontier_gate.cpp b/src/Disks/tests/gtest_cas_gc_frontier_gate.cpp index d8c90fa7dc1c..a35f8c50e6a9 100644 --- a/src/Disks/tests/gtest_cas_gc_frontier_gate.cpp +++ b/src/Disks/tests/gtest_cas_gc_frontier_gate.cpp @@ -241,7 +241,7 @@ class DrainRaceBackend final : public CountingBackend std::function after_read_hook; }; -class PostFoldUnreadableTerminalBackend final : public CountingBackend +class DeadCheckpointHeadFailureBackend final : public CountingBackend { public: /// Unhide the names the primitive overrides below would otherwise shadow. @@ -260,7 +260,7 @@ class PostFoldUnreadableTerminalBackend final : public CountingBackend std::optional head(const String & key, TransportAccess & access) override { if (!bypass_fault && key == unreadable_key) - throw std::runtime_error("injected post-fold terminal read failure for " + key); + throw std::runtime_error("injected dead checkpoint HEAD failure for " + key); return CountingBackend::head(key, access); } @@ -2821,10 +2821,11 @@ TEST(CASGCFrontierGate, CleanupEvidenceLeavesRemovedNamespaceCheckpointForJanito /// Once a terminal has folded, a later physical read failure is janitor debt, not lifecycle evidence /// loss. Removing this per-key leak handling would either make the signal disappear or let one dead -/// object prevent the janitor from considering the rest of its page. -TEST(CASGCFrontierGate, PostFoldUnreadableTerminalIsCountedWithoutSuppressingProgress) +/// object prevent the janitor from considering the rest of its page. The unreadable object is the dead +/// checkpoint: dead `_log`/`_snap` keys are deleted token-free and never read. +TEST(CASGCFrontierGate, PostFoldUnreadableDeadCheckpointIsCountedWithoutSuppressingProgress) { - auto backend = std::make_shared(); + auto backend = std::make_shared(); auto store = openPoolForTest(backend, /*gc_fold_max_defer_rounds=*/0); CasRequests requests = openRequestsForTest(backend); CasOperation op = requests.admit(); @@ -2876,10 +2877,11 @@ TEST(CASGCFrontierGate, PostFoldUnreadableTerminalIsCountedWithoutSuppressingPro dropRefTransition(*backend, layout, progressing, "victim", manifest); const String terminal_key = layout.refLogKey(removed_life, RefTxnId{1, 2}); + const String checkpoint_key = layout.refCkptKey(removed_life); const String later_dead_residue = layout.refLogKey(removed_life, RefTxnId{1, 3}); ASSERT_TRUE(std::holds_alternative( op.create(later_dead_residue, "dead residue after the folded terminal", Retry::once()))); - backend->makeUnreadable(terminal_key); + backend->makeUnreadable(checkpoint_key); std::map namespace_cleanup; const uint64_t leaks_before @@ -2900,7 +2902,8 @@ TEST(CASGCFrontierGate, PostFoldUnreadableTerminalIsCountedWithoutSuppressingPro EXPECT_EQ(report.manifests_deleted, 1u) << "the janitor leak cannot promote itself into pool-wide destructive suppression"; EXPECT_FALSE(op.head(layout.manifestKey(manifest_id), Retry::once()).has_value()); - EXPECT_TRUE(backend->existsIgnoringFault(terminal_key)); + EXPECT_TRUE(backend->existsIgnoringFault(checkpoint_key)); + EXPECT_FALSE(backend->existsIgnoringFault(terminal_key)); EXPECT_FALSE(backend->existsIgnoringFault(later_dead_residue)) << "one unreadable key cannot stop the perpetual janitor from deciding the rest of its page"; ASSERT_FALSE(namespace_cleanup.empty()); @@ -2909,7 +2912,7 @@ TEST(CASGCFrontierGate, PostFoldUnreadableTerminalIsCountedWithoutSuppressingPro ProfileEvents::global_counters[ProfileEvents::CASGCNamespaceCleanupLeaks].load() - leaks_before, 1u); const String captured = log_capture.captured(); - EXPECT_NE(captured.find(terminal_key), String::npos); + EXPECT_NE(captured.find(checkpoint_key), String::npos); EXPECT_NE(captured.find("leak"), String::npos); } diff --git a/src/Disks/tests/gtest_cas_gc_round_defer.cpp b/src/Disks/tests/gtest_cas_gc_round_defer.cpp index 43a6a909e86a..4de654aa23f5 100644 --- a/src/Disks/tests/gtest_cas_gc_round_defer.cpp +++ b/src/Disks/tests/gtest_cas_gc_round_defer.cpp @@ -472,14 +472,14 @@ TEST(CASGCRoundDefer, DeferredRoundRetriesPartialJanitorPageAtForcedFoldWithoutP ASSERT_TRUE(folded.acquired_lease); ASSERT_FALSE(folded.deferred) << "gc_fold_max_defer_rounds=1 forces the round immediately following one DEFER to fold"; - EXPECT_EQ(backend->listCount(layout.namespaceRootPrefix()), 1u) - << "the authoritative fold must run the janitor exactly once, not once per call site"; + EXPECT_EQ(backend->listCount(layout.namespaceRootPrefix()), 2u) + << "one janitor run, not one per call site: the page that deletes, then one from the stream start"; const auto folded_cleanup = std::find_if(phases.begin(), phases.end(), [](const GcPhaseRecord & phase) { return phase.phase == "namespace_cleanup"; }); ASSERT_NE(folded_cleanup, phases.end()); - EXPECT_EQ(folded_cleanup->metrics.at("janitor_pages"), 1u); + EXPECT_EQ(folded_cleanup->metrics.at("janitor_pages"), 2u); EXPECT_GE(folded_cleanup->metrics.at("janitor_keys"), 1u); EXPECT_EQ(folded_cleanup->metrics.at("janitor_deleted"), 1u); EXPECT_EQ(static_cast(op.head(key_a, Retry::once()).has_value()) + static_cast(op.head(key_b, Retry::once()).has_value()), 0u) diff --git a/src/Disks/tests/gtest_cas_namespace_janitor_batches.cpp b/src/Disks/tests/gtest_cas_namespace_janitor_batches.cpp new file mode 100644 index 000000000000..ecb6accdd109 --- /dev/null +++ b/src/Disks/tests/gtest_cas_namespace_janitor_batches.cpp @@ -0,0 +1,422 @@ +#include "cas_test_helpers.h" +#include +#include +#include +#include + +#include +#include +#include + +/// Batch deletes of dead-life `_log`/`_snap` keys: one page through `NamespaceJanitor`, then whole rounds +/// through `Gc`, which adds the page loop, the time budget and the GC I/O pool. + +using namespace DB::Cas; +using namespace DB::Cas::tests; + +namespace DB::ErrorCodes +{ + extern const int NETWORK_ERROR; + extern const int NOT_IMPLEMENTED; +} + +namespace +{ + +void createObj(Backend & backend, const String & key, const String & bytes) +{ + OperationForTest op(backend); + ASSERT_TRUE(std::holds_alternative((*op).create(key, bytes, Retry::once()))); +} + +bool present(Backend & backend, const String & key) +{ + OperationForTest op(backend); + return (*op).head(key, Retry::standard()).has_value(); +} + +size_t countPresent(Backend & backend, const std::vector & keys) +{ + size_t n = 0; + for (const String & key : keys) + n += present(backend, key) ? 1 : 0; + return n; +} + +NamespaceLifeId life(const char * name, uint64_t id) +{ + return NamespaceLifeId::fromCatalogEntry(RootNamespace{name}, DB::UInt128{id}); +} + +std::vector seedLogs(Backend & backend, const Layout & layout, const NamespaceLifeId & of, uint64_t count) +{ + std::vector keys; + for (uint64_t i = 1; i <= count; ++i) + { + keys.push_back(layout.refLogKey(of, RefTxnId{.writer_epoch = 1, .ref_sequence = i})); + createObj(backend, keys.back(), "log"); + } + return keys; +} + +GcMaintenanceReadResult readState(CasRequests & requests, const Layout & layout) +{ + auto op = requests.admit(); + return readGcMaintenanceState(op, layout); +} + +/// Every batch delete fails; single-key behaviour is untouched. +class FailingBulkBackend : public CountingBackend +{ +public: + void removeManyWriteOnce(const std::vector &, TransportAccess &) override + { + attempted = true; + throw DB::Exception(DB::ErrorCodes::NETWORK_ERROR, "bulk delete failed"); + } + std::atomic attempted{false}; +}; + +/// Fails every `LIST` of the namespace root after the first `allowed_lists` ones. +class FailingNamespaceListBackend : public CountingBackend +{ +public: + RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override + { + if (prefix.ends_with("/cas/ns/") && lists++ >= allowed_lists) + throw DB::Exception(DB::ErrorCodes::NETWORK_ERROR, "namespace LIST failed"); + return CountingBackend::list(prefix, cursor, limit, access); + } + std::atomic lists{0}; + std::atomic allowed_lists{std::numeric_limits::max()}; +}; + +/// The first batch delete goes through and flips `delete_done`, which a liveness predicate reads. +class FlipAfterBulkBackend : public CountingBackend +{ +public: + void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override + { + CountingBackend::removeManyWriteOnce(keys, access); + delete_done = true; + } + std::atomic delete_done{false}; +}; + +/// A store without a batch delete: more than one key in a request is `NOT_IMPLEMENTED`. +class NoBatchDeleteBackend : public CountingBackend +{ +public: + void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override + { + if (keys.size() > 1) + { + ++refused_batches; + throw DB::Exception(DB::ErrorCodes::NOT_IMPLEMENTED, "no batch delete"); + } + ++single_key_requests; + CountingBackend::removeManyWriteOnce(keys, access); + } + std::atomic refused_batches{0}; + std::atomic single_key_requests{0}; +}; + +/// Each `LIST` of the namespace root moves the clock, so a test can spend the janitor's budget in one page. +class ClockOnNamespaceListBackend : public CountingBackend +{ +public: + RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override + { + if (prefix.ends_with("/cas/ns/")) + clock_ms += step_ms; + return CountingBackend::list(prefix, cursor, limit, access); + } + std::atomic clock_ms{1}; + uint64_t step_ms = 0; +}; + +template +PoolPtr openPool(std::shared_ptr backend, std::optional refuse_at = std::nullopt) +{ + PoolConfig config{.pool_prefix = "p", .server_root_id = "test"}; + config.gc_fold_max_defer_rounds = 0; + config.gc_io_concurrency = 4; + config.gc_io_pool_refuse_at_for_test = refuse_at; + return Pool::open(std::move(backend), std::move(config)); +} + +/// One reclaiming round; returns the `namespace_cleanup` phase metrics. +std::map runRoundAndGetJanitorMetrics(Gc & gc) +{ + std::map metrics; + gc.setPhaseSink([&](const GcPhaseRecord & record) + { + if (record.phase == "namespace_cleanup") + metrics = record.metrics; + }); + const RoundReport report = runRegularRoundReclaiming(gc); + gc.setPhaseSink({}); + EXPECT_TRUE(report.acquired_lease); + return metrics; +} + +} + +TEST(CASNamespaceJanitorBatches, DeadStreamKeysGoInOneBatchAndStateKeysStayExact) +{ + auto backend = std::make_shared(); + CasRequests requests(backend, Fence::open()); + const Layout layout("p"); + createObj(*backend, layout.refCatalogKey(), encodeRefCatalog({})); + const NamespaceLifeId dead = life("dead", 301); + std::vector stream = seedLogs(*backend, layout, dead, 3); + stream.push_back(layout.refSnapshotKey(dead, RefTxnId{.writer_epoch = 1, .ref_sequence = 2})); + createObj(*backend, stream.back(), "snap"); + const String ckpt = layout.refCkptKey(dead); + createObj(*backend, ckpt, "ckpt"); + backend->resetCounts(); + + const NamespaceJanitorResult result = NamespaceJanitor(requests, layout, 100).runOnePage(false, [] { return true; }); + + EXPECT_EQ(backend->bulkRemoveCalls(), 1u); + EXPECT_EQ(backend->headTotal(), 0u); + EXPECT_EQ(result.delete_requests, 2u) << "one batch and one exact-token delete"; + EXPECT_EQ(result.deleted, 5u); + EXPECT_EQ(result.leaked, 0u); + EXPECT_EQ(countPresent(*backend, stream), 0u); + EXPECT_FALSE(present(*backend, ckpt)); +} + +TEST(CASNamespaceJanitorBatches, RebornLifeOnTheSamePageKeepsItsKeys) +{ + auto backend = std::make_shared(); + CasRequests requests(backend, Fence::open()); + const Layout layout("p"); + const CatalogEntry reborn{.ns = RootNamespace{"t"}, .state = NsState::Live, .incarnation = DB::UInt128{302}}; + createObj(*backend, layout.refCatalogKey(), encodeRefCatalog(RefCatalog{.entries = {reborn}})); + const std::vector dead_stream = seedLogs(*backend, layout, life("t", 301), 3); + const std::vector live_stream + = seedLogs(*backend, layout, NamespaceLifeId::fromCatalogEntry(reborn.ns, reborn.incarnation), 3); + backend->resetCounts(); + + const NamespaceJanitorResult result = NamespaceJanitor(requests, layout, 100).runOnePage(false, [] { return true; }); + + EXPECT_EQ(result.deleted, 3u); + EXPECT_EQ(backend->deleteTotal(), 3u); + EXPECT_EQ(countPresent(*backend, dead_stream), 0u); + EXPECT_EQ(countPresent(*backend, live_stream), 3u); +} + +TEST(CASNamespaceJanitorBatches, FailedBatchLeaksAndCursorAdvances) +{ + auto backend = std::make_shared(); + CasRequests requests(backend, Fence::open()); + const Layout layout("p"); + createObj(*backend, layout.refCatalogKey(), encodeRefCatalog({})); + const std::vector stream = seedLogs(*backend, layout, life("dead", 303), 4); + + const NamespaceJanitorResult result = NamespaceJanitor(requests, layout, 2).runOnePage(false, [] { return true; }); + + EXPECT_EQ(result.deleted, 0u); + EXPECT_EQ(result.leaked, 2u); + EXPECT_EQ(countPresent(*backend, stream), 4u); + const auto state = readState(requests, layout); + ASSERT_EQ(state.status, GcMaintenanceReadStatus::Valid); + EXPECT_FALSE(state.state->janitor_cursor.empty()) << "a failed batch is leak-only: it must not pin the cursor"; +} + +TEST(CASNamespaceJanitorBatches, UnpublishedPageReportsNoLeaks) +{ + auto backend = std::make_shared(); + CasRequests requests(backend, Fence::open()); + const Layout layout("p"); + createObj(*backend, layout.refCatalogKey(), encodeRefCatalog({})); + seedLogs(*backend, layout, life("dead", 305), 4); + + const NamespaceJanitorResult result + = NamespaceJanitor(requests, layout, 2).runOnePage(false, [&] { return !backend->attempted.load(); }); + + EXPECT_EQ(readState(requests, layout).status, GcMaintenanceReadStatus::Absent); + EXPECT_EQ(result.leaked, 0u) << "the page is listed again, so nothing on it is leaked yet"; + EXPECT_FALSE(result.anomalies.empty()); +} + +TEST(CASNamespaceJanitorBatches, AuthorityLostAfterBatchHoldsCursor) +{ + auto backend = std::make_shared(); + CasRequests requests(backend, Fence::open()); + const Layout layout("p"); + createObj(*backend, layout.refCatalogKey(), encodeRefCatalog({})); + const std::vector stream = seedLogs(*backend, layout, life("dead", 304), 4); + + const NamespaceJanitorResult result + = NamespaceJanitor(requests, layout, 2).runOnePage(false, [&] { return !backend->delete_done.load(); }); + + EXPECT_EQ(result.deleted, 2u); + EXPECT_EQ(countPresent(*backend, stream), 2u); + EXPECT_EQ(readState(requests, layout).status, GcMaintenanceReadStatus::Absent) + << "a tenure that lost authority must not publish progress"; +} + +TEST(CASNamespaceJanitorBatches, RoundDrainsEveryPageOfDebris) +{ + auto backend = std::make_shared(); + auto store = openPool(backend); + const Layout & layout = store->layout(); + const std::vector stream = seedLogs(*backend, layout, life("dead", 311), 2500); + backend->resetCounts(); + + Gc gc(store, DB::UInt128{312}); + const auto metrics = runRoundAndGetJanitorMetrics(gc); + + EXPECT_EQ(countPresent(*backend, stream), 0u); + /// Three pages of debris, then one empty page from the stream start ends the pass. + EXPECT_EQ(metrics.at("janitor_pages"), 4u); + EXPECT_EQ(metrics.at("janitor_deleted"), 2500u); + EXPECT_EQ(metrics.at("leaked"), 0u); + EXPECT_EQ(backend->listCount(layout.namespaceRootPrefix()), 4u); + EXPECT_EQ(metrics.at("delete_requests"), 12u); + /// Every page is split across the 4 pool threads. + EXPECT_EQ(backend->bulkRemoveCalls(), 12u); +} + +TEST(CASNamespaceJanitorBatches, RoundDrainsDebrisOnBothSidesOfTheCursor) +{ + auto backend = std::make_shared(); + auto store = openPool(backend); + const Layout & layout = store->layout(); + const std::vector stream = seedLogs(*backend, layout, life("dead", 320), 2500); + { + CasOperation op = store->openRequests().admit(); + const GcMaintenanceReadResult progress = readGcMaintenanceState(op, layout); + ASSERT_TRUE(std::holds_alternative(casGcMaintenanceState( + op, layout, progress.etag, GcMaintenanceState{.janitor_cursor = stream[1499]}, Retry::standard()))); + } + + Gc gc(store, DB::UInt128{321}); + const auto metrics = runRoundAndGetJanitorMetrics(gc); + + EXPECT_EQ(metrics.at("janitor_deleted"), 2500u); + EXPECT_EQ(countPresent(*backend, stream), 0u) << "the pass continues from the stream start after its last page"; +} + +TEST(CASNamespaceJanitorBatches, ListFailureAfterTheFirstPageKeepsTheCursor) +{ + auto backend = std::make_shared(); + auto store = openPool(backend); + const Layout & layout = store->layout(); + const std::vector stream = seedLogs(*backend, layout, life("dead", 322), 2500); + backend->lists = 0; + backend->allowed_lists = 1; + + Gc gc(store, DB::UInt128{323}); + const auto failed = runRoundAndGetJanitorMetrics(gc); + EXPECT_EQ(failed.at("janitor_deleted"), 1000u); + EXPECT_EQ(countPresent(*backend, stream), 1500u); + const auto state = readState(store->openRequests(), layout); + ASSERT_EQ(state.status, GcMaintenanceReadStatus::Valid); + EXPECT_EQ(state.state->janitor_cursor, stream[999]) << "the cursor the first page published stays"; + + backend->allowed_lists = std::numeric_limits::max(); + const auto resumed = runRoundAndGetJanitorMetrics(gc); + EXPECT_EQ(resumed.at("janitor_deleted"), 1500u); + EXPECT_EQ(countPresent(*backend, stream), 0u); +} + +TEST(CASNamespaceJanitorBatches, BudgetStopsThePassAfterThePageInProgress) +{ + auto backend = std::make_shared(); + auto store = openPool(backend); + const Layout & layout = store->layout(); + const std::vector stream = seedLogs(*backend, layout, life("dead", 313), 2500); + backend->step_ms = 30'000; + + Gc gc(store, DB::UInt128{314}, {}, [&] { return backend->clock_ms.load(); }); + const auto first = runRoundAndGetJanitorMetrics(gc); + + EXPECT_EQ(first.at("janitor_pages"), 1u); + EXPECT_EQ(first.at("janitor_deleted"), 1000u); + ASSERT_TRUE(first.contains("budget_exhausted")); + EXPECT_EQ(first.at("budget_exhausted"), 1u); + EXPECT_EQ(countPresent(*backend, stream), 1500u); + + backend->step_ms = 0; + const auto second = runRoundAndGetJanitorMetrics(gc); + EXPECT_EQ(second.at("budget_exhausted"), 0u); + EXPECT_EQ(countPresent(*backend, stream), 0u) << "the next round resumes from the published cursor"; +} + +TEST(CASNamespaceJanitorBatches, QuietPoolCostsOneListPerRound) +{ + /// More than one page of live keys, and a clock that spends the budget on the first `LIST`: a pass that + /// went on after a page without debris would stop on the budget instead. + auto backend = std::make_shared(); + auto store = openPool(backend); + const Layout & layout = store->layout(); + const RootNamespace live_namespace{"00/live@cas@"}; + fixture::admitLive(*backend, layout, live_namespace); + const NamespaceLifeId live = fixture::fixtureLife(live_namespace); + createObj(*backend, layout.refCkptKey(live), + encodeRefCkpt(RefCkpt{.life_epoch = std::optional{1}, + .checkpoint_snapshot_id = std::nullopt, .last_epoch_seal = std::nullopt})); + std::vector files; + for (size_t i = 0; i < 1200; ++i) + { + files.push_back(layout.namespaceFilesPrefix(live) + fmt::format("file-{:04}", i)); + createObj(*backend, files.back(), "file"); + } + backend->step_ms = 30'000; + backend->resetCounts(); + + Gc gc(store, DB::UInt128{315}, {}, [&] { return backend->clock_ms.load(); }); + const auto metrics = runRoundAndGetJanitorMetrics(gc); + + EXPECT_EQ(metrics.at("janitor_pages"), 1u); + EXPECT_EQ(metrics.at("janitor_deleted"), 0u); + EXPECT_EQ(metrics.at("budget_exhausted"), 0u); + EXPECT_EQ(backend->listCount(layout.namespaceRootPrefix()), 1u); + EXPECT_EQ(countPresent(*backend, files), 1200u); + /// The page was decided, not suppressed: its cursor is published. + const auto state = readState(store->openRequests(), layout); + ASSERT_EQ(state.status, GcMaintenanceReadStatus::Valid); + EXPECT_FALSE(state.state->janitor_cursor.empty()); +} + +TEST(CASNamespaceJanitorBatches, StoreWithoutBatchDeleteDrainsKeyByKeyOnThePool) +{ + auto backend = std::make_shared(); + auto store = openPool(backend); + const Layout & layout = store->layout(); + const std::vector stream = seedLogs(*backend, layout, life("dead", 316), 1200); + + Gc gc(store, DB::UInt128{317}); + const auto metrics = runRoundAndGetJanitorMetrics(gc); + + EXPECT_EQ(countPresent(*backend, stream), 0u); + EXPECT_EQ(metrics.at("janitor_deleted"), 1200u); + EXPECT_EQ(metrics.at("leaked"), 0u); + EXPECT_EQ(backend->single_key_requests.load(), 1200u); + /// Two pages, each split into one job per pool thread. + EXPECT_EQ(backend->refused_batches.load(), 8u); + EXPECT_EQ(metrics.at("delete_requests"), 1208u); +} + +TEST(CASNamespaceJanitorBatches, RefusedEnqueueLeaksTheUnscheduledJobsAndTheNextRoundDrains) +{ + auto backend = std::make_shared(); + auto store = openPool(backend, /*refuse_at=*/1); + const Layout & layout = store->layout(); + const std::vector stream = seedLogs(*backend, layout, life("dead", 318), 100); + + /// Four jobs of 25 keys: the first runs, the second is refused and takes the rest with it. + Gc gc(store, DB::UInt128{319}); + const auto refused = runRoundAndGetJanitorMetrics(gc); + EXPECT_EQ(refused.at("janitor_deleted"), 25u); + EXPECT_EQ(refused.at("leaked"), 75u); + EXPECT_EQ(countPresent(*backend, stream), 75u); + + const auto drained = runRoundAndGetJanitorMetrics(gc); + EXPECT_EQ(drained.at("janitor_deleted"), 75u); + EXPECT_EQ(drained.at("leaked"), 0u); + EXPECT_EQ(countPresent(*backend, stream), 0u); +} diff --git a/src/Disks/tests/gtest_cas_write_once_key.cpp b/src/Disks/tests/gtest_cas_write_once_key.cpp index c05eb6b1e345..d9fd69c6756d 100644 --- a/src/Disks/tests/gtest_cas_write_once_key.cpp +++ b/src/Disks/tests/gtest_cas_write_once_key.cpp @@ -30,3 +30,25 @@ TEST(CASWriteOnceKey, FactoriesMintTheSameStringsAsThePlainKeyFunctions) EXPECT_TRUE(layout.parseManifestKey(layout.writeOnceManifestKey(manifest).str()).has_value()); EXPECT_TRUE(layout.parseRefObjectKey(layout.writeOnceRefLogKey(life, id).str()).has_value()); } + +TEST(CASWriteOnceKey, StreamKeyMintedFromAListedKeyMustBeThatKey) +{ + const Layout layout{"p"}; + const NamespaceLifeId life = NamespaceLifeId::fromCatalogEntry(RootNamespace{"test/aa@cas@"}, DB::UInt128(0x1234)); + const RefTxnId id{5, 7}; + + for (const String & listed : {layout.refLogKey(life, id), layout.refSnapshotKey(life, id)}) + { + const auto parsed = layout.parseRefObjectKey(listed); + ASSERT_TRUE(parsed.has_value()) << listed; + const auto minted = layout.writeOnceStreamKey(*parsed, listed); + ASSERT_TRUE(minted.has_value()) << listed; + EXPECT_EQ(minted->str(), listed); + } + + const auto parsed_log = layout.parseRefObjectKey(layout.refLogKey(life, id)); + ASSERT_TRUE(parsed_log.has_value()); + EXPECT_FALSE(layout.writeOnceStreamKey(*parsed_log, layout.refSnapshotKey(life, id)).has_value()); + EXPECT_FALSE(layout.writeOnceStreamKey(*parsed_log, layout.refLogKey(life, RefTxnId{5, 8})).has_value()); + EXPECT_FALSE(layout.writeOnceStreamKey(*parsed_log, "q" + layout.refLogKey(life, id).substr(1)).has_value()); +} diff --git a/tests/integration/test_cas_janitor_drain/__init__.py b/tests/integration/test_cas_janitor_drain/__init__.py new file mode 100644 index 000000000000..e69de29bb2d1 diff --git a/tests/integration/test_cas_janitor_drain/configs/storage_conf.xml b/tests/integration/test_cas_janitor_drain/configs/storage_conf.xml new file mode 100644 index 000000000000..28aedcb75703 --- /dev/null +++ b/tests/integration/test_cas_janitor_drain/configs/storage_conf.xml @@ -0,0 +1,36 @@ + + + + + object_storage + s3 + cas + itest-janitor-batch + http://rustfs1:11121/test/cas_janitor_batch/ + clickhouse + clickhouse + 1 + 86400 + + 1 + + + object_storage + s3 + cas + itest-janitor-no-batch + http://rustfs1:11121/test/cas_janitor_no_batch/ + clickhouse + clickhouse + false + 1 + 86400 + 1 + + + +
cas_batch
+
cas_no_batch
+
+
+
diff --git a/tests/integration/test_cas_janitor_drain/test.py b/tests/integration/test_cas_janitor_drain/test.py new file mode 100644 index 000000000000..8526355320b2 --- /dev/null +++ b/tests/integration/test_cas_janitor_drain/test.py @@ -0,0 +1,90 @@ +""" +A dropped table's ref stream drains in one deleting round of the namespace janitor on RustFS, on a store +with a batch delete and on one without. +""" +import pytest + +from helpers.cluster import ClickHouseCluster + +cluster = ClickHouseCluster(__file__) + +INSERTS = 2500 +MAX_ROUNDS = 3 + + +@pytest.fixture(scope="module", autouse=True) +def start_cluster(): + cluster.add_instance( + "node", + main_configs=["configs/storage_conf.xml"], + with_rustfs=True, + stay_alive=True, + ) + try: + cluster.start() + yield cluster + finally: + cluster.shutdown() + + +def stream_keys(disk): + prefix = f"cas_janitor_{disk.removeprefix('cas_')}/cas/ns/stream/" + return [ + o.object_name + for o in cluster.rustfs_client.list_objects(cluster.rustfs_bucket, prefix, recursive=True) + if "/_log/" in o.object_name or "/_snap/" in o.object_name + ] + + +def seed_and_drop(node, disk): + node.query(f"DROP TABLE IF EXISTS t_{disk} SYNC") + node.query( + f"CREATE TABLE t_{disk} (k UInt64) ENGINE = MergeTree ORDER BY k " + f"SETTINGS storage_policy = '{disk}', parts_to_delay_insert = 100000, parts_to_throw_insert = 100000" + ) + node.query(f"SYSTEM STOP MERGES t_{disk}") + # One row per block gives one part per row without INSERTS client round trips. + node.query( + f"INSERT INTO t_{disk} SELECT number FROM numbers({INSERTS}) SETTINGS max_block_size = 1, " + "max_insert_block_size = 1, min_insert_block_size_rows = 1, min_insert_block_size_bytes = 0, " + "max_insert_threads = 1" + ) + assert len(stream_keys(disk)) >= INSERTS + node.query(f"DROP TABLE t_{disk} SYNC") + return len(stream_keys(disk)) + + +def janitor_rows(node, disk, since): + node.query("SYSTEM FLUSH LOGS") + columns = ["janitor_deleted", "janitor_pages", "leaked", "delete_requests", "budget_exhausted"] + select = ", ".join(f"phase_metrics['{c}']" for c in columns) + out = node.query( + f"SELECT {select} FROM system.cas_gc_log WHERE event_type = 'Phase' AND phase = 'namespace_cleanup' " + f"AND disk_name = '{disk}' AND event_time_microseconds >= toDateTime64('{since}', 6) " + "ORDER BY event_time_microseconds FORMAT TSV" + ) + return [dict(zip(columns, map(int, line.split("\t")))) for line in out.strip().splitlines()] + + +@pytest.mark.parametrize("disk", ["cas_batch", "cas_no_batch"]) +def test_dropped_table_drains_in_one_round(disk): + node = cluster.instances["node"] + since = node.query("SELECT now64(6)").strip() + dead_keys = seed_and_drop(node, disk) + for _ in range(MAX_ROUNDS): + node.query(f"SYSTEM CAS GC RUN {disk}") + if not stream_keys(disk): + break + assert not stream_keys(disk), f"{disk}: dead ref stream keys left after {MAX_ROUNDS} rounds" + rows = janitor_rows(node, disk, since) + deleting = [r for r in rows if r["janitor_deleted"] > 0] + assert len(deleting) == 1, f"{disk}: expected one deleting round, got {rows}" + row = deleting[0] + assert row["janitor_deleted"] >= dead_keys, f"{disk}: the round did not drain all {dead_keys} dead keys: {rows}" + assert row["janitor_pages"] >= 3 + assert row["leaked"] == 0 + assert row["budget_exhausted"] == 0 + if disk == "cas_batch": + assert row["delete_requests"] * 10 < dead_keys, f"{disk}: the dead keys did not go in batches: {rows}" + else: + assert row["delete_requests"] >= dead_keys, f"{disk}: expected one request per key: {rows}"