|
|
@@ -2115,8 +2115,10 @@ void installTtlHook(MemoryStore& mstore, RelationManager& rm, LmdbDocumentStore&
|
|
|
const smartbotic::database::CascadeNotifyFn& notify = nullptr) {
|
|
|
mstore.setTtlExpiryRelationHook(
|
|
|
[&rm, &store, &pm, &mstore, &cfgManager, notify](const std::string& coll,
|
|
|
- const std::string& id) {
|
|
|
- return ttlExpiryDecision(rm, store, pm, mstore, cfgManager, coll, id, notify);
|
|
|
+ const std::string& id,
|
|
|
+ bool firstAttempt) {
|
|
|
+ return ttlExpiryDecision(rm, store, pm, mstore, cfgManager, coll, id, notify,
|
|
|
+ firstAttempt);
|
|
|
});
|
|
|
}
|
|
|
|
|
|
@@ -2168,7 +2170,10 @@ void test_ttl_restrict_blocks_the_expiry() {
|
|
|
TmpPersistence p("ttl-restrict-wal");
|
|
|
check(p.pm.start(), "persistence manager started");
|
|
|
|
|
|
- MemoryStore mstore(MemoryStore::Config{});
|
|
|
+ // Small retry cadence so the test does not have to run 60 sweeps.
|
|
|
+ MemoryStore::Config cfg;
|
|
|
+ cfg.ttlBlockedRetrySweeps = 3;
|
|
|
+ MemoryStore mstore(cfg);
|
|
|
mstore.start();
|
|
|
RelationManager rm(mstore);
|
|
|
rm.loadFromStore();
|
|
|
@@ -2188,21 +2193,30 @@ void test_ttl_restrict_blocks_the_expiry() {
|
|
|
"the skip is counted, so an operator can see a document is stuck past its TTL");
|
|
|
check(mstore.getStats().expiredCount == 0, "and it is not counted as expired");
|
|
|
|
|
|
- // The expiry must stay ARMED so the next sweep retries - a skip that also
|
|
|
- // dropped the expiration index entry would leave the document permanently
|
|
|
- // unexpirable even after the last child went away.
|
|
|
+ check(mstore.ttlBlockedDocumentCount() == 1,
|
|
|
+ "and it is tracked as ONE stuck document - the gauge an operator wants");
|
|
|
+
|
|
|
+ // ⚠ The next sweep must NOT re-examine it (close-out review, finding 2). A
|
|
|
+ // blocked document costs one LMDB read txn + cursor scan + one log line every
|
|
|
+ // time it is examined, and at the 1s default that is a permanent flood plus a
|
|
|
+ // permanent index-lookup load. It is skipped for free until its retry is due.
|
|
|
const uint64_t again = mstore.expireDocuments();
|
|
|
- check(again == 0, "still blocked on the second sweep");
|
|
|
- check(mstore.getStats().ttlExpiryBlockedByRelation == 2,
|
|
|
- "the second sweep re-examined it, so the expiry is still armed");
|
|
|
+ check(again == 0, "still not expired on the next sweep");
|
|
|
+ check(mstore.getStats().ttlExpiryBlockedByRelation == 1,
|
|
|
+ "and the relation was NOT re-consulted - no second refusal event");
|
|
|
+ check(residentInMemory(mstore, "default:workflows", "wf-1"), "the parent is still there");
|
|
|
|
|
|
- // And once the child is gone, the same sweep expires it - the block is the
|
|
|
- // relation's, not a permanent quarantine.
|
|
|
+ // The expiry must stay ARMED, so once the retry comes due AND the child is
|
|
|
+ // gone, it expires - the block is the relation's, not a permanent quarantine.
|
|
|
store.del("executions", "ex-1");
|
|
|
check(mstore.remove("default:executions", "ex-1"), "child removed");
|
|
|
- const uint64_t third = mstore.expireDocuments();
|
|
|
- check(third == 1, "with no children left, the retry finally expires the parent");
|
|
|
- check(!mstore.get("default:workflows", "wf-1").has_value(), "the parent is gone now");
|
|
|
+ uint64_t total = 0;
|
|
|
+ for (uint32_t i = 0; i < cfg.ttlBlockedRetrySweeps + 1; ++i) {
|
|
|
+ total += mstore.expireDocuments();
|
|
|
+ }
|
|
|
+ check(total == 1, "when the retry comes due, the parent finally expires");
|
|
|
+ check(!residentInMemory(mstore, "default:workflows", "wf-1"), "the parent is gone now");
|
|
|
+ check(mstore.ttlBlockedDocumentCount() == 0, "and it is no longer counted as stuck");
|
|
|
|
|
|
p.pm.stop();
|
|
|
mstore.stop();
|
|
|
@@ -2250,6 +2264,15 @@ void test_ttl_cascade_deletes_children_like_a_manual_delete() {
|
|
|
check(notifications.size() == 2,
|
|
|
"replication/events fired once per child plus once for the parent");
|
|
|
|
|
|
+ // Accounting (close-out review, minor 1): the parent is ONE expiry, and the
|
|
|
+ // one child is ONE delete. The cascade removes the parent through
|
|
|
+ // unloadDocument(), which counts a DELETE, so without compensating for that
|
|
|
+ // the parent was counted as both an expiry and a delete.
|
|
|
+ const auto st = mstore.getStats();
|
|
|
+ check(st.expiredCount == 1, "the parent counted as exactly one expiry");
|
|
|
+ check(st.deleteCount == 1,
|
|
|
+ "and the delete count covers the CHILD only - the parent is not counted twice");
|
|
|
+
|
|
|
// The property a hand-rolled sweeper cascade breaks: MemoryStore is rebuilt
|
|
|
// from snapshot + WAL, NEVER from LMDB, so a cascade whose child deletions
|
|
|
// exist only in LMDB has them RESURRECTED on the next boot.
|
|
|
@@ -2497,6 +2520,229 @@ void test_ttl_cascade_wal_is_durable_before_the_lmdb_commit() {
|
|
|
mstore.stop();
|
|
|
}
|
|
|
|
|
|
+
|
|
|
+// =========================================================================
|
|
|
+// v2.11.0 close-out review, FINDING 1 — a concurrent TTL change must not have
|
|
|
+// its document (and its children) deleted by an in-flight sweep.
|
|
|
+//
|
|
|
+// Phase 1 collects candidates and releases every lock; the hook may then DELETE
|
|
|
+// the document and cascade its children, and the hook has no view of
|
|
|
+// `expiresAt`, so `Handled` cannot re-check the way `Proceed` does. An
|
|
|
+// update()/patch() that extends or clears the TTL in that window would
|
|
|
+// otherwise destroy a document that is no longer expired - and its children with
|
|
|
+// it. The pre-v2.11.0 loop held the collection lock across the whole erase and
|
|
|
+// had no such window, so this is a cost of the two-phase split.
|
|
|
+//
|
|
|
+// ⚠ HOW THIS IS INDUCED, and why it is honest rather than staged: a
|
|
|
+// single-threaded test cannot interleave a real writer, so the write is made to
|
|
|
+// land at exactly the wrong moment from INSIDE the sweep. Two documents expire
|
|
|
+// in the same sweep, wf-early before wf-late (the expiration index is keyed by
|
|
|
+// expiry time and walked in ascending order, so the order is deterministic).
|
|
|
+// The hook, while handling wf-early, extends wf-late's TTL - which is precisely
|
|
|
+// "a write landed after phase 1 collected wf-late and before the sweeper got to
|
|
|
+// it". Everything about wf-late's path after that point is the production path.
|
|
|
+void test_ttl_concurrent_ttl_extension_is_not_expired() {
|
|
|
+ TmpEnv t("ttl-race");
|
|
|
+ LmdbDocumentStore store(t.env);
|
|
|
+ TmpPersistence p("ttl-race-wal");
|
|
|
+ check(p.pm.start(), "persistence manager started");
|
|
|
+
|
|
|
+ MemoryStore mstore(MemoryStore::Config{});
|
|
|
+ mstore.start();
|
|
|
+ RelationManager rm(mstore);
|
|
|
+ rm.loadFromStore();
|
|
|
+ CollectionConfigManager cfgManager(mstore);
|
|
|
+
|
|
|
+ RelationInfo rel;
|
|
|
+ rel.name = "default:exec_wf";
|
|
|
+ rel.child = "default:executions";
|
|
|
+ rel.childField = "workflowId";
|
|
|
+ rel.parent = "default:workflows";
|
|
|
+ rel.onDelete = OnDelete::Cascade;
|
|
|
+ std::string mgrErr;
|
|
|
+ check(rm.createRelation(rel, mgrErr), "declared the cascade relation");
|
|
|
+ store.set_relations("executions",
|
|
|
+ {RelationRef{"exec_wf", "workflowId", "workflows", false, true}});
|
|
|
+
|
|
|
+ const uint64_t farFuture = static_cast<uint64_t>(1) << 62;
|
|
|
+
|
|
|
+ auto seedParent = [&](const std::string& id, uint64_t expiresAt) {
|
|
|
+ Document d; d.id = id; d.collection = "workflows";
|
|
|
+ d.set_data({{"name", id}});
|
|
|
+ d.expiresAt = expiresAt;
|
|
|
+ store.put("workflows", id, d);
|
|
|
+ mstore.loadDocument("default:workflows", d);
|
|
|
+ };
|
|
|
+ auto seedChild = [&](const std::string& id, const std::string& parentId) {
|
|
|
+ Document d; d.id = id; d.collection = "executions";
|
|
|
+ d.set_data({{"workflowId", parentId}});
|
|
|
+ store.put("executions", id, d);
|
|
|
+ mstore.loadDocument("default:executions", d);
|
|
|
+ };
|
|
|
+
|
|
|
+ // expiresAt 1 sorts before 2, so wf-early is handled first in the sweep.
|
|
|
+ seedParent("wf-early", 1);
|
|
|
+ seedParent("wf-late", 2);
|
|
|
+ seedChild("ex-early", "wf-early");
|
|
|
+ seedChild("ex-late", "wf-late");
|
|
|
+
|
|
|
+ // The injected "concurrent" write: while the sweeper is busy expiring
|
|
|
+ // wf-early (a real cascade - WAL fsync plus an LMDB commit, which is what
|
|
|
+ // makes the window wide), wf-late's TTL is extended.
|
|
|
+ bool injected = false;
|
|
|
+ mstore.setTtlExpiryRelationHook(
|
|
|
+ [&](const std::string& coll, const std::string& id, bool firstAttempt) {
|
|
|
+ if (!injected && id == "wf-early") {
|
|
|
+ injected = true;
|
|
|
+ Document renewed;
|
|
|
+ renewed.id = "wf-late";
|
|
|
+ renewed.collection = "workflows";
|
|
|
+ renewed.set_data({{"name", "wf-late"}});
|
|
|
+ renewed.expiresAt = farFuture; // TTL extended
|
|
|
+ mstore.loadDocument("default:workflows", renewed);
|
|
|
+ }
|
|
|
+ return ttlExpiryDecision(rm, store, p.pm, mstore, cfgManager, coll, id,
|
|
|
+ nullptr, firstAttempt);
|
|
|
+ });
|
|
|
+
|
|
|
+ const uint64_t expired = mstore.expireDocuments();
|
|
|
+
|
|
|
+ check(injected, "the injected concurrent TTL extension actually ran");
|
|
|
+ check(expired == 1, "exactly ONE document expired - wf-early only");
|
|
|
+
|
|
|
+ // wf-early: genuinely expired, cascaded as normal. Proves the sweep worked.
|
|
|
+ check(!residentInMemory(mstore, "default:workflows", "wf-early"), "wf-early expired");
|
|
|
+ check(!store.get("executions", "ex-early").has_value(), "its child cascaded away");
|
|
|
+
|
|
|
+ // wf-late: the whole point. Not expired, and - the data-loss half - its
|
|
|
+ // child was not cascaded either.
|
|
|
+ check(residentInMemory(mstore, "default:workflows", "wf-late"),
|
|
|
+ "wf-late was NOT expired - its TTL was extended after phase 1 collected it");
|
|
|
+ check(store.get("workflows", "wf-late").has_value(), "and it is intact in LMDB");
|
|
|
+ check(store.get("executions", "ex-late").has_value(),
|
|
|
+ "and ITS CHILD was not cascaded - this is the data-loss half of the race");
|
|
|
+ check(store.relation_index_child_count("exec_wf", "wf-late") == 1,
|
|
|
+ "the reverse-index posting survived too");
|
|
|
+ check(mstore.getStats().ttlExpiryBlockedByRelation == 0,
|
|
|
+ "and it was not recorded as relation-blocked - it simply is not expired");
|
|
|
+
|
|
|
+ p.pm.stop();
|
|
|
+ mstore.stop();
|
|
|
+}
|
|
|
+
|
|
|
+// =========================================================================
|
|
|
+// v2.11.0 close-out review, FINDING 2 — a relation-blocked document must not
|
|
|
+// consume the sweep's candidate budget.
|
|
|
+//
|
|
|
+// Phase 1 never erases a blocked document's expiration index entry (deliberately
|
|
|
+// - it has to stay armed so the block can lift), so with one shared budget those
|
|
|
+// documents refill the cap on every sweep. `collections_` iteration order is
|
|
|
+// stable, so every collection after them SILENTLY STOPS BEING SWEPT - and TTL'd
|
|
|
+// parents under `restrict` is exactly the combination this feature creates, so
|
|
|
+// that is the expected steady state, not an edge case.
|
|
|
+//
|
|
|
+// Budgets come from Config so this can be shown at 3 documents instead of 10001.
|
|
|
+void test_ttl_blocked_documents_do_not_starve_the_sweep() {
|
|
|
+ TmpEnv t("ttl-starve");
|
|
|
+ LmdbDocumentStore store(t.env);
|
|
|
+ TmpPersistence p("ttl-starve-wal");
|
|
|
+ check(p.pm.start(), "persistence manager started");
|
|
|
+
|
|
|
+ MemoryStore::Config cfg;
|
|
|
+ cfg.ttlMaxCandidatesPerSweep = 3; // exactly filled by the blocked parents
|
|
|
+ cfg.ttlBlockedRetrySweeps = 100; // far enough away not to interfere
|
|
|
+ MemoryStore mstore(cfg);
|
|
|
+ mstore.start();
|
|
|
+ RelationManager rm(mstore);
|
|
|
+ rm.loadFromStore();
|
|
|
+ CollectionConfigManager cfgManager(mstore);
|
|
|
+
|
|
|
+ std::atomic<bool> healthy{true};
|
|
|
+ std::atomic<uint64_t> drift{0};
|
|
|
+ mstore.setDocumentStoreMirror(
|
|
|
+ [&store](std::string_view) -> smartbotic::db::storage::DocumentStore* { return &store; },
|
|
|
+ &healthy, &drift);
|
|
|
+
|
|
|
+ RelationInfo rel;
|
|
|
+ rel.name = "default:exec_wf";
|
|
|
+ rel.child = "default:executions";
|
|
|
+ rel.childField = "workflowId";
|
|
|
+ rel.parent = "default:workflows";
|
|
|
+ rel.onDelete = OnDelete::Restrict;
|
|
|
+ std::string mgrErr;
|
|
|
+ check(rm.createRelation(rel, mgrErr), "declared the restrict relation");
|
|
|
+ store.set_relations("executions",
|
|
|
+ {RelationRef{"exec_wf", "workflowId", "workflows", false, true}});
|
|
|
+
|
|
|
+ // Three restrict-blocked TTL'd parents. Their expiry times are the earliest,
|
|
|
+ // so they are collected first and fill the budget on the first sweep.
|
|
|
+ for (int i = 0; i < 3; ++i) {
|
|
|
+ const std::string pid = "wf-" + std::to_string(i);
|
|
|
+ Document parent; parent.id = pid; parent.collection = "workflows";
|
|
|
+ parent.set_data({{"name", pid}});
|
|
|
+ parent.expiresAt = 1 + static_cast<uint64_t>(i);
|
|
|
+ store.put("workflows", pid, parent);
|
|
|
+ mstore.loadDocument("default:workflows", parent);
|
|
|
+
|
|
|
+ Document child; child.id = "ex-" + std::to_string(i);
|
|
|
+ child.collection = "executions";
|
|
|
+ child.set_data({{"workflowId", pid}});
|
|
|
+ store.put("executions", child.id, child);
|
|
|
+ mstore.loadDocument("default:executions", child);
|
|
|
+ }
|
|
|
+
|
|
|
+ // The victim: a TTL'd document with NO children, expiring after the three
|
|
|
+ // blocked parents. Deliberately in the SAME collection, because that makes
|
|
|
+ // the ordering deterministic - `collections_` is an unordered_map, but a
|
|
|
+ // collection's expirationIndex is a std::map walked in ascending expiry
|
|
|
+ // order, so wf-0/1/2 (1, 2, 3) are always collected before wf-victim (100).
|
|
|
+ // Starvation inside one collection is the same defect as starvation across
|
|
|
+ // collections, and this way the test cannot pass or fail on hash order.
|
|
|
+ Document victim;
|
|
|
+ victim.id = "wf-victim";
|
|
|
+ victim.collection = "workflows";
|
|
|
+ victim.set_data({{"name", "wf-victim"}});
|
|
|
+ victim.expiresAt = 100;
|
|
|
+ store.put("workflows", "wf-victim", victim);
|
|
|
+ mstore.loadDocument("default:workflows", victim);
|
|
|
+
|
|
|
+ installTtlHook(mstore, rm, store, p.pm, cfgManager);
|
|
|
+
|
|
|
+ // Sweep 1: the three blocked parents are fresh, so they legitimately take
|
|
|
+ // the whole budget and the victim is not reached at all.
|
|
|
+ const uint64_t first = mstore.expireDocuments();
|
|
|
+ check(first == 0, "sweep 1 expired nothing - the budget went to the blocked parents");
|
|
|
+ check(mstore.getStats().ttlExpiryBlockedByRelation == 3, "all three parents were refused");
|
|
|
+ check(mstore.ttlBlockedDocumentCount() == 3, "and all three are tracked as stuck");
|
|
|
+ check(residentInMemory(mstore, "default:workflows", "wf-victim"),
|
|
|
+ "and the victim has not been reached yet");
|
|
|
+
|
|
|
+ // Sweep 2: the blocked parents are known and not due, so they must consume
|
|
|
+ // NOTHING and the victim finally gets in. With one shared budget it never
|
|
|
+ // would - the three blocked parents refill the cap on every sweep, forever.
|
|
|
+ uint64_t total = first;
|
|
|
+ total += mstore.expireDocuments();
|
|
|
+ check(total == 1,
|
|
|
+ "the childless document was expired despite three blocked parents filling "
|
|
|
+ "the budget - blocked documents no longer starve the sweep");
|
|
|
+ check(!residentInMemory(mstore, "default:workflows", "wf-victim"), "the victim is gone");
|
|
|
+ check(!store.get("workflows", "wf-victim").has_value(), "and its DELETE reached LMDB");
|
|
|
+ check(mstore.getStats().ttlExpiryBlockedByRelation == 3,
|
|
|
+ "and the blocked parents were not re-consulted - still three refusal events, "
|
|
|
+ "not three more per sweep");
|
|
|
+
|
|
|
+ // Ten more sweeps: still nothing re-examined, so the steady-state cost of a
|
|
|
+ // persistent block is zero rather than one LMDB scan + one WARN per document
|
|
|
+ // per second.
|
|
|
+ for (int i = 0; i < 10; ++i) mstore.expireDocuments();
|
|
|
+ check(mstore.getStats().ttlExpiryBlockedByRelation == 3,
|
|
|
+ "ten further sweeps consulted the relation zero times");
|
|
|
+ check(mstore.ttlBlockedDocumentCount() == 3, "and the three are still tracked as stuck");
|
|
|
+
|
|
|
+ p.pm.stop();
|
|
|
+ mstore.stop();
|
|
|
+}
|
|
|
+
|
|
|
// =========================================================================
|
|
|
// v2.11.0 close-out — BOOT-PASS DRIFT MUST NOT DISABLE DESTRUCTIVE POLICIES.
|
|
|
//
|
|
|
@@ -2643,6 +2889,8 @@ int main() {
|
|
|
test_ttl_no_action_expires_and_leaves_the_reference_dangling();
|
|
|
test_ttl_with_no_relations_expires_exactly_as_before();
|
|
|
test_ttl_cascade_wal_is_durable_before_the_lmdb_commit();
|
|
|
+ test_ttl_concurrent_ttl_extension_is_not_expired();
|
|
|
+ test_ttl_blocked_documents_do_not_starve_the_sweep();
|
|
|
test_boot_pass_drift_does_not_disable_the_cascade();
|
|
|
test_pending_remirror_list_is_released_after_the_pass();
|
|
|
|