|
@@ -180,6 +180,7 @@ bool MemoryStore::createCollection(const std::string& name, const CollectionOpti
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
bool MemoryStore::dropCollection(const std::string& name) {
|
|
bool MemoryStore::dropCollection(const std::string& name) {
|
|
|
|
|
+ {
|
|
|
std::unique_lock<std::shared_mutex> lock(globalMutex_);
|
|
std::unique_lock<std::shared_mutex> lock(globalMutex_);
|
|
|
|
|
|
|
|
auto it = collections_.find(name);
|
|
auto it = collections_.find(name);
|
|
@@ -225,6 +226,25 @@ bool MemoryStore::dropCollection(const std::string& name) {
|
|
|
|
|
|
|
|
// Update memory tracking atomically
|
|
// Update memory tracking atomically
|
|
|
estimatedMemoryBytes_.fetch_sub(totalSize, std::memory_order_relaxed);
|
|
estimatedMemoryBytes_.fetch_sub(totalSize, std::memory_order_relaxed);
|
|
|
|
|
+ } // globalMutex_ released here — see below.
|
|
|
|
|
+
|
|
|
|
|
+ // v2.11.0 close-out (round-3 review, item 2) — release any TTL blocked-set
|
|
|
|
|
+ // entries for this collection.
|
|
|
|
|
+ //
|
|
|
|
|
+ // Necessary because expireDocuments() cannot clean these up itself: phase 1
|
|
|
|
|
+ // builds candidates by walking `collections_`, so once the collection is gone
|
|
|
|
|
+ // no candidate is ever produced for those ids again and phase 2 never sees
|
|
|
|
|
+ // them. They would sit in ttlBlocked_ for the process lifetime, inflating
|
|
|
|
|
+ // GetMemoryStats.ttl_blocked_documents with documents that no longer exist.
|
|
|
|
|
+ //
|
|
|
|
|
+ // ⚠ AFTER the globalMutex_ scope, and both halves of that matter.
|
|
|
|
|
+ // AFTER the drop, so a concurrent sweep cannot re-add an entry between the
|
|
|
|
|
+ // clear and the erase (once the collection is gone nothing can re-add).
|
|
|
|
|
+ // OUTSIDE the lock, because ttlBlockedMutex_ is ordered BEFORE globalMutex_
|
|
|
|
|
+ // (phase 1 holds it across that acquisition) and taking it while holding
|
|
|
|
|
+ // globalMutex_ exclusively would close a deadlock cycle - the same inversion
|
|
|
|
|
+ // this review found in phase 2.
|
|
|
|
|
+ clearTtlBlockedForCollection(name);
|
|
|
|
|
|
|
|
return true;
|
|
return true;
|
|
|
}
|
|
}
|
|
@@ -2212,7 +2232,7 @@ void MemoryStore::setTtlExpiryRelationHook(TtlExpiryRelationHook hook) {
|
|
|
uint64_t MemoryStore::expireDocuments() {
|
|
uint64_t MemoryStore::expireDocuments() {
|
|
|
uint64_t expired = 0;
|
|
uint64_t expired = 0;
|
|
|
auto now = currentTimeMs();
|
|
auto now = currentTimeMs();
|
|
|
- const uint64_t sweep = ++ttlSweepCounter_;
|
|
|
|
|
|
|
+ const uint64_t sweep = ttlSweepCounter_.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
|
|
|
|
|
|
// ---- Phase 1: collect candidates under the locks, mutate nothing. -----
|
|
// ---- Phase 1: collect candidates under the locks, mutate nothing. -----
|
|
|
// See the header comment on expireDocuments() for why the phases are split:
|
|
// See the header comment on expireDocuments() for why the phases are split:
|
|
@@ -2294,8 +2314,9 @@ uint64_t MemoryStore::expireDocuments() {
|
|
|
// Periodic summary instead of a per-document WARN per sweep. A flood is not
|
|
// Periodic summary instead of a per-document WARN per sweep. A flood is not
|
|
|
// a signal, and at the 1s default a persistent block was ~one WARN per
|
|
// a signal, and at the 1s default a persistent block was ~one WARN per
|
|
|
// blocked document per second, indefinitely.
|
|
// blocked document per second, indefinitely.
|
|
|
- if (blockedNotDue > 0 && sweep - ttlBlockedSummarySweep_ >= kBlockedSummarySweeps) {
|
|
|
|
|
- ttlBlockedSummarySweep_ = sweep;
|
|
|
|
|
|
|
+ if (blockedNotDue > 0 &&
|
|
|
|
|
+ sweep - ttlBlockedSummarySweep_.load(std::memory_order_relaxed) >= kBlockedSummarySweeps) {
|
|
|
|
|
+ ttlBlockedSummarySweep_.store(sweep, std::memory_order_relaxed);
|
|
|
spdlog::warn("TTL expiry: {} document(s) are still past their TTL and cannot be "
|
|
spdlog::warn("TTL expiry: {} document(s) are still past their TTL and cannot be "
|
|
|
"expired because a relation blocks the delete (each was logged once "
|
|
"expired because a relation blocks the delete (each was logged once "
|
|
|
"when first blocked; they are retried every {} sweeps). They will "
|
|
"when first blocked; they are retried every {} sweeps). They will "
|
|
@@ -2320,26 +2341,55 @@ uint64_t MemoryStore::expireDocuments() {
|
|
|
// See the header comment for the RESIDUAL: a write landing during the
|
|
// See the header comment for the RESIDUAL: a write landing during the
|
|
|
// hook's own cascade (a WAL fsync plus an LMDB commit) is still
|
|
// hook's own cascade (a WAL fsync plus an LMDB commit) is still
|
|
|
// possible. Narrowing, not eliminating.
|
|
// possible. Narrowing, not eliminating.
|
|
|
- uint64_t currentExpiresAt = 0;
|
|
|
|
|
- {
|
|
|
|
|
|
|
+ //
|
|
|
|
|
+ // ⚠ SKIPPED ENTIRELY WITH NO HOOK INSTALLED (round-3 review, item 4).
|
|
|
|
|
+ // Without a hook every candidate takes the Proceed path, whose own
|
|
|
|
|
+ // re-check shares one critical section with its erase and is therefore
|
|
|
|
|
+ // airtight - so this pre-check would be a second EXCLUSIVE collection
|
|
|
|
|
+ // lock per candidate (up to ttlMaxCandidatesPerSweep per sweep) to guard
|
|
|
|
|
+ // a window that path does not have. Every install that declares no
|
|
|
|
|
+ // relation stays exactly as fast as before v2.11.0.
|
|
|
|
|
+ //
|
|
|
|
|
+ // ⚠ LOCK ORDER (round-3 review, item 1): clearTtlBlocked() MUST NOT be
|
|
|
|
|
+ // called from inside this scope. ttlBlockedMutex_ is ordered BEFORE
|
|
|
|
|
+ // globalMutex_/coll->mutex (phase 1 holds it across both), so taking it
|
|
|
|
|
+ // while holding either closes a cycle: one sweeper holding
|
|
|
|
|
+ // ttlBlockedMutex_ and waiting for coll->mutex against another holding
|
|
|
|
|
+ // coll->mutex and waiting for ttlBlockedMutex_ - and both then pin
|
|
|
|
|
+ // globalMutex_ shared, so the next getOrCreateCollection() for a new
|
|
|
|
|
+ // collection blocks on the exclusive acquire and the whole store hangs.
|
|
|
|
|
+ // The outcome is therefore recorded and acted on AFTER the scope closes.
|
|
|
|
|
+ enum class PreCheck { Ok, Vanished, NoLongerExpired };
|
|
|
|
|
+ PreCheck pre = PreCheck::Ok;
|
|
|
|
|
+ if (ttlExpiryRelationHook_) {
|
|
|
std::shared_lock<std::shared_mutex> globalLock(globalMutex_);
|
|
std::shared_lock<std::shared_mutex> globalLock(globalMutex_);
|
|
|
auto collIt = collections_.find(cand.collection);
|
|
auto collIt = collections_.find(cand.collection);
|
|
|
- if (collIt == collections_.end()) continue; // collection dropped
|
|
|
|
|
- auto* coll = collIt->second.get();
|
|
|
|
|
- std::unique_lock<std::shared_mutex> collLock(coll->mutex);
|
|
|
|
|
- auto docIt = coll->documents.find(cand.id);
|
|
|
|
|
- if (docIt == coll->documents.end()) {
|
|
|
|
|
- // Gone (deleted, or evicted). Drop the stale index entry so it
|
|
|
|
|
- // is not re-examined every sweep forever.
|
|
|
|
|
- removeFromExpirationIndex(*coll, cand.id, cand.expiresAt);
|
|
|
|
|
- clearTtlBlocked(cand.collection, cand.id);
|
|
|
|
|
- continue;
|
|
|
|
|
|
|
+ if (collIt == collections_.end()) {
|
|
|
|
|
+ // Collection dropped. Nothing to unindex (it went with the
|
|
|
|
|
+ // collection), but the blocked entry has to go or it inflates
|
|
|
|
|
+ // ttlBlockedDocumentCount() for the process lifetime - phase 1
|
|
|
|
|
+ // never revisits a key whose collection is gone, and
|
|
|
|
|
+ // dropCollection() knows nothing about ttlBlocked_ (round-3
|
|
|
|
|
+ // review, item 2).
|
|
|
|
|
+ pre = PreCheck::Vanished;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ auto* coll = collIt->second.get();
|
|
|
|
|
+ std::unique_lock<std::shared_mutex> collLock(coll->mutex);
|
|
|
|
|
+ auto docIt = coll->documents.find(cand.id);
|
|
|
|
|
+ if (docIt == coll->documents.end()) {
|
|
|
|
|
+ // Gone (deleted, or evicted). Drop the stale index entry so
|
|
|
|
|
+ // it is not re-examined every sweep forever.
|
|
|
|
|
+ removeFromExpirationIndex(*coll, cand.id, cand.expiresAt);
|
|
|
|
|
+ pre = PreCheck::Vanished;
|
|
|
|
|
+ } else if (docIt->second.expiresAt == 0 || docIt->second.expiresAt > now) {
|
|
|
|
|
+ pre = PreCheck::NoLongerExpired;
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
- currentExpiresAt = docIt->second.expiresAt;
|
|
|
|
|
}
|
|
}
|
|
|
- if (currentExpiresAt == 0 || currentExpiresAt > now) {
|
|
|
|
|
- // The TTL was extended or cleared while this sweep was running. Not
|
|
|
|
|
- // expired any more - and no longer blocked either, if it was.
|
|
|
|
|
|
|
+ if (pre != PreCheck::Ok) {
|
|
|
|
|
+ // Vanished, or its TTL was extended/cleared while this sweep was
|
|
|
|
|
+ // running. Either way it is not expiring now, and it is no longer
|
|
|
|
|
+ // blocked either. Called with NO other lock held - see above.
|
|
|
clearTtlBlocked(cand.collection, cand.id);
|
|
clearTtlBlocked(cand.collection, cand.id);
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
@@ -2485,6 +2535,19 @@ void MemoryStore::clearTtlBlocked(const std::string& collection, const std::stri
|
|
|
ttlBlocked_.erase(ttlBlockedKey(collection, id));
|
|
ttlBlocked_.erase(ttlBlockedKey(collection, id));
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+void MemoryStore::clearTtlBlockedForCollection(const std::string& collection) {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(ttlBlockedMutex_);
|
|
|
|
|
+ if (ttlBlocked_.empty()) return;
|
|
|
|
|
+ // Keys are "<collection>\0<id>" (see ttlBlockedKey), so a prefix match on
|
|
|
|
|
+ // the collection plus the separator is exact - it cannot match a different
|
|
|
|
|
+ // collection whose name merely starts with this one.
|
|
|
|
|
+ const std::string prefix = collection + std::string(1, '\0');
|
|
|
|
|
+ for (auto it = ttlBlocked_.begin(); it != ttlBlocked_.end();) {
|
|
|
|
|
+ it = (it->first.compare(0, prefix.size(), prefix) == 0) ? ttlBlocked_.erase(it)
|
|
|
|
|
+ : std::next(it);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
size_t MemoryStore::ttlBlockedDocumentCount() const {
|
|
size_t MemoryStore::ttlBlockedDocumentCount() const {
|
|
|
std::lock_guard<std::mutex> lock(ttlBlockedMutex_);
|
|
std::lock_guard<std::mutex> lock(ttlBlockedMutex_);
|
|
|
return ttlBlocked_.size();
|
|
return ttlBlocked_.size();
|