Selaa lähdekoodia

feat(storage): txn-accepting overloads on LmdbDocumentStore

Add put(WriteTxn&, ...), del(WriteTxn&, ...), put_vector(WriteTxn&, ...)
and del_vector(WriteTxn&, ...), extending the pattern open_for_write /
maintainIndexes / maintainRelations already use to the public surface.
Task 12's atomic cascade needs one transaction spanning the parent
document, each affected child, and every index/relation sub-db they
touch - the no-txn forms each open and commit their own WriteTxn, so
they cannot compose.

Each overload takes a trailing to_cache vector and appends every MDB_dbi
it opens, but does not commit and does not cache - both stay the
caller's job, since caching before commit is the exact v2.8.0 bug (an
abort of the caller's own still-open transaction would leave a closed
handle behind, poisoning the collection with EINVAL for the process's
life). The no-txn forms become thin wrappers: open, call, commit, then
cache every returned handle.

Extracted cachedDbi() - a cache-only existence probe with no transaction
argument - to replace the ReadTxn+try_open_for_read probe del()/
del_vector() used only to avoid creating an empty sub-db on a no-op
delete; try_open_for_read never actually touches its ReadTxn on a cache
hit or miss, so the probe was cache-only already.

Test: test_subdb_identity gains
test_txn_accepting_writes_are_atomic_across_collections - two documents
in different collections written in one WriteTxn are all-or-nothing on
abort (no index entry survives either), the collections are not
poisoned afterward, and a committed-but-uncached handle reads back as
absent until primed - pinning both known outage shapes from the
CLAUDE.md LMDB incident history.
fszontagh 1 kuukausi sitten
vanhempi
sitoutus
95977820b7

+ 84 - 38
service/src/storage/document_store_lmdb.cpp

@@ -390,6 +390,14 @@ void LmdbDocumentStore::cacheCommittedDbi(std::string_view collection,
     dbi_cache_[std::string(collection)] = dbi;
 }
 
+std::optional<unsigned int>
+LmdbDocumentStore::cachedDbi(std::string_view collection) {
+    std::lock_guard<std::mutex> lock(cache_mutex_);
+    auto it = dbi_cache_.find(std::string(collection));
+    if (it == dbi_cache_.end()) return std::nullopt;
+    return it->second;
+}
+
 std::optional<unsigned int>
 LmdbDocumentStore::try_open_for_read(ReadTxn& rtxn,
                                       std::string_view collection) {
@@ -500,8 +508,21 @@ size_t LmdbDocumentStore::prime_dbi_cache() {
 void LmdbDocumentStore::put(std::string_view collection,
                              std::string_view id,
                              const smartbotic::database::Document& doc) {
-    std::string payload = encode_document(doc);
     WriteTxn wtxn(env_);
+    std::vector<std::pair<std::string, unsigned int>> to_cache;
+    put(wtxn, collection, id, doc, to_cache);
+    // Commit BEFORE any handle is cached - the v2.8.0 lesson. Every handle in
+    // to_cache is private to wtxn until this succeeds.
+    wtxn.commit();
+    for (const auto& [sub, d] : to_cache) cacheCommittedDbi(sub, d);
+}
+
+void LmdbDocumentStore::put(WriteTxn& wtxn,
+                             std::string_view collection,
+                             std::string_view id,
+                             const smartbotic::database::Document& doc,
+                             std::vector<std::pair<std::string, unsigned int>>& to_cache) {
+    std::string payload = encode_document(doc);
     unsigned int dbi = open_for_write(wtxn, collection);
     MDB_val k = to_val(id);
 
@@ -509,7 +530,6 @@ void LmdbDocumentStore::put(std::string_view collection,
     // throws after this point rolls the index back with the document. An
     // asynchronously-maintained index would let a query read entries for a row
     // that was never stored, and return silently wrong rows rather than an error.
-    std::vector<std::pair<std::string, unsigned int>> index_dbis;
     const bool needs_index = !indexed_fields(collection).empty();
     const bool needs_relations = !relations(collection).empty();
     if (needs_index || needs_relations) {
@@ -519,18 +539,18 @@ void LmdbDocumentStore::put(std::string_view collection,
         const std::string_view old_payload =
             rc == MDB_SUCCESS ? to_sv(old) : std::string_view{};
         if (needs_index) {
-            maintainIndexes(wtxn, collection, id, old_payload, &doc, index_dbis);
+            maintainIndexes(wtxn, collection, id, old_payload, &doc, to_cache);
         }
         if (needs_relations) {
-            maintainRelations(wtxn, collection, id, old_payload, &doc, index_dbis);
+            maintainRelations(wtxn, collection, id, old_payload, &doc, to_cache);
         }
     }
 
     MDB_val v = to_val(payload);
     mdb_check(mdb_put(wtxn.raw(), dbi, &k, &v, 0), "put");
-    wtxn.commit();
-    cacheCommittedDbi(collection, dbi);
-    for (const auto& [sub, d] : index_dbis) cacheCommittedDbi(sub, d);
+    // NOT cached here - see the header contract. The caller caches every
+    // entry in to_cache (this one included) only after ITS commit succeeds.
+    to_cache.emplace_back(std::string(collection), dbi);
 }
 
 
@@ -1428,23 +1448,34 @@ bool LmdbDocumentStore::del(std::string_view collection, std::string_view id) {
     // Same reasoning as get(): a zero-length key cannot exist, so there is
     // nothing to delete rather than an error to raise.
     if (id.empty()) return false;
-    // Open as a write txn unconditionally so we have MDB_CREATE available
-    // if the collection doesn't exist yet — but in that case there's
-    // nothing to delete; just probe with a read txn first to avoid
-    // accidentally creating an empty sub-db on a no-op delete.
-    {
-        ReadTxn rtxn(env_);
-        auto dbi_opt = try_open_for_read(rtxn, collection);
-        if (!dbi_opt) return false;
-    }
+    // Probe the cache first to avoid accidentally creating an empty sub-db on
+    // a no-op delete - no transaction needed for this, see cachedDbi().
+    if (!cachedDbi(collection)) return false;
 
     WriteTxn wtxn(env_);
+    std::vector<std::pair<std::string, unsigned int>> to_cache;
+    const bool existed = del(wtxn, collection, id, to_cache);
+    // Commit BEFORE any handle is cached - the v2.8.0 lesson.
+    wtxn.commit();
+    for (const auto& [sub, d] : to_cache) cacheCommittedDbi(sub, d);
+    return existed;
+}
+
+bool LmdbDocumentStore::del(WriteTxn& wtxn,
+                             std::string_view collection,
+                             std::string_view id,
+                             std::vector<std::pair<std::string, unsigned int>>& to_cache) {
+    if (id.empty()) return false;
+    // Same no-op-avoidance as the no-txn form above: don't spring an empty
+    // sub-db into existence for a collection that has never been written.
+    // Cache-only, so this costs nothing extra inside the caller's txn.
+    if (!cachedDbi(collection)) return false;
+
     unsigned int dbi = open_for_write(wtxn, collection);
     MDB_val k = to_val(id);
 
     // Remove index entries before the row goes, while its stored bytes are
     // still readable - they are the only record of which index keys it owns.
-    std::vector<std::pair<std::string, unsigned int>> index_dbis;
     const bool needs_index = !indexed_fields(collection).empty();
     const bool needs_relations = !relations(collection).empty();
     if (needs_index || needs_relations) {
@@ -1452,10 +1483,10 @@ bool LmdbDocumentStore::del(std::string_view collection, std::string_view id) {
         const int grc = mdb_get(wtxn.raw(), dbi, &k, &old);
         if (grc == MDB_SUCCESS) {
             if (needs_index) {
-                maintainIndexes(wtxn, collection, id, to_sv(old), nullptr, index_dbis);
+                maintainIndexes(wtxn, collection, id, to_sv(old), nullptr, to_cache);
             }
             if (needs_relations) {
-                maintainRelations(wtxn, collection, id, to_sv(old), nullptr, index_dbis);
+                maintainRelations(wtxn, collection, id, to_sv(old), nullptr, to_cache);
             }
         } else if (grc != MDB_NOTFOUND) {
             throw_mdb(grc, "get (pre-index del)");
@@ -1463,15 +1494,11 @@ bool LmdbDocumentStore::del(std::string_view collection, std::string_view id) {
     }
 
     int rc = mdb_del(wtxn.raw(), dbi, &k, nullptr);
-    if (rc == MDB_NOTFOUND) {
-        wtxn.commit();
-        cacheCommittedDbi(collection, dbi);
-        return false;
-    }
+    // NOT cached here - the caller caches every entry in to_cache (this one
+    // included) only after ITS commit succeeds.
+    to_cache.emplace_back(std::string(collection), dbi);
+    if (rc == MDB_NOTFOUND) return false;
     if (rc != MDB_SUCCESS) throw_mdb(rc, "del");
-    wtxn.commit();
-    cacheCommittedDbi(collection, dbi);
-    for (const auto& [sub, d] : index_dbis) cacheCommittedDbi(sub, d);
     return true;
 }
 
@@ -2385,34 +2412,53 @@ void LmdbDocumentStore::put_vector(std::string_view collection,
                                     std::string_view id,
                                     const std::vector<float>& vec) {
     if (vec.empty()) return;  // mirror the migration tool's no-op semantics
-    const std::string subdb = vector_subdb_name(collection);
     WriteTxn wtxn(env_);
+    std::vector<std::pair<std::string, unsigned int>> to_cache;
+    put_vector(wtxn, collection, id, vec, to_cache);
+    wtxn.commit();
+    for (const auto& [sub, d] : to_cache) cacheCommittedDbi(sub, d);
+}
+
+void LmdbDocumentStore::put_vector(WriteTxn& wtxn,
+                                    std::string_view collection,
+                                    std::string_view id,
+                                    const std::vector<float>& vec,
+                                    std::vector<std::pair<std::string, unsigned int>>& to_cache) {
+    if (vec.empty()) return;  // mirror the migration tool's no-op semantics
+    const std::string subdb = vector_subdb_name(collection);
     unsigned int dbi = open_for_write(wtxn, subdb);
     MDB_val k = to_val(id);
     MDB_val v{vec.size() * sizeof(float),
               const_cast<void*>(static_cast<const void*>(vec.data()))};
     mdb_check(mdb_put(wtxn.raw(), dbi, &k, &v, 0), "put_vector");
-    wtxn.commit();
-    cacheCommittedDbi(subdb, dbi);
+    to_cache.emplace_back(subdb, dbi);
 }
 
 bool LmdbDocumentStore::del_vector(std::string_view collection, std::string_view id) {
     const std::string subdb = vector_subdb_name(collection);
     // Probe first so we don't create an empty vectors sub-db just to
-    // discover the vector isn't there.
-    {
-        ReadTxn rtxn(env_);
-        auto dbi_opt = try_open_for_read(rtxn, subdb);
-        if (!dbi_opt) return false;
-    }
+    // discover the vector isn't there. Cache-only, see cachedDbi().
+    if (!cachedDbi(subdb)) return false;
     WriteTxn wtxn(env_);
+    std::vector<std::pair<std::string, unsigned int>> to_cache;
+    const bool existed = del_vector(wtxn, collection, id, to_cache);
+    wtxn.commit();
+    for (const auto& [sub, d] : to_cache) cacheCommittedDbi(sub, d);
+    return existed;
+}
+
+bool LmdbDocumentStore::del_vector(WriteTxn& wtxn,
+                                    std::string_view collection,
+                                    std::string_view id,
+                                    std::vector<std::pair<std::string, unsigned int>>& to_cache) {
+    const std::string subdb = vector_subdb_name(collection);
+    if (!cachedDbi(subdb)) return false;
     unsigned int dbi = open_for_write(wtxn, subdb);
     MDB_val k = to_val(id);
     int rc = mdb_del(wtxn.raw(), dbi, &k, nullptr);
+    to_cache.emplace_back(subdb, dbi);
     if (rc == MDB_NOTFOUND) return false;
     if (rc != MDB_SUCCESS) throw_mdb(rc, "del_vector");
-    wtxn.commit();
-    cacheCommittedDbi(subdb, dbi);
     return true;
 }
 

+ 54 - 0
service/src/storage/document_store_lmdb.hpp

@@ -305,6 +305,37 @@ public:
 
     bool del(std::string_view collection, std::string_view id) override;
 
+    // v2.11.0 T10 — txn-accepting overloads. Task 12's atomic cascade needs
+    // ONE transaction spanning the parent document, each affected child, and
+    // every index/relation sub-db a write touches - the no-txn forms above
+    // each open and commit their own WriteTxn, so they cannot compose into
+    // one atomic unit. Internal helpers already take a WriteTxn& (
+    // open_for_write, maintainIndexes, maintainRelations); this extends that
+    // pattern to the public surface rather than inventing a new one.
+    //
+    // ⚠ These do NOT commit and do NOT cache - both stay the caller's job.
+    // Caching before commit is the v2.8.0 bug: a later abort of the SAME
+    // caller-owned transaction would leave a closed handle in the cache,
+    // poisoning the collection with EINVAL for the life of the process. Every
+    // MDB_dbi this call opens (the collection's own handle, plus any
+    // index/relation sub-db maintenance touched) is appended to `to_cache`;
+    // call cacheCommittedDbi() for each entry ONLY after wtxn.commit()
+    // succeeds. Never skip that step either - see try_open_for_read: reads
+    // are cache-only, so a committed-but-uncached sub-db reads as absent for
+    // the rest of the process's life (this is how relations would silently
+    // stop enforcing). The no-txn forms above are thin wrappers that do
+    // exactly this: open, call, commit, then cache every returned handle.
+    void put(class WriteTxn& wtxn,
+             std::string_view collection,
+             std::string_view id,
+             const smartbotic::database::Document& doc,
+             std::vector<std::pair<std::string, unsigned int>>& to_cache);
+
+    bool del(class WriteTxn& wtxn,
+             std::string_view collection,
+             std::string_view id,
+             std::vector<std::pair<std::string, unsigned int>>& to_cache);
+
     uint64_t count(std::string_view collection) override;
 
     ScanResult scan(std::string_view collection,
@@ -319,6 +350,19 @@ public:
                     std::string_view id,
                     const std::vector<float>& vec) override;
     bool del_vector(std::string_view collection, std::string_view id) override;
+
+    // v2.11.0 T10 — txn-accepting vector equivalents. Same contract as
+    // put(WriteTxn&, ...)/del(WriteTxn&, ...) above: no commit, no cache,
+    // handles opened are appended to `to_cache` for the caller.
+    void put_vector(class WriteTxn& wtxn,
+                    std::string_view collection,
+                    std::string_view id,
+                    const std::vector<float>& vec,
+                    std::vector<std::pair<std::string, unsigned int>>& to_cache);
+    bool del_vector(class WriteTxn& wtxn,
+                    std::string_view collection,
+                    std::string_view id,
+                    std::vector<std::pair<std::string, unsigned int>>& to_cache);
     std::optional<std::vector<float>>
     get_vector(std::string_view collection, std::string_view id) override;
     void scan_vectors(
@@ -520,6 +564,16 @@ private:
     void cacheCommittedDbi(std::string_view collection, unsigned int dbi);
     std::optional<unsigned int> try_open_for_read(class ReadTxn& rtxn,
                                                    std::string_view collection);
+
+    // v2.11.0 T10 — cache-only existence probe, no transaction required.
+    // Mirrors try_open_for_read's cache-hit path exactly (which never
+    // touches its ReadTxn& argument either — the cache is complete by
+    // construction once primed, see that function's long comment) and its
+    // cache-miss path (nullopt). Lets a txn-accepting overload ask "does
+    // this collection already exist" without opening a second transaction
+    // inside the caller's already-open WriteTxn — nesting a ReadTxn there
+    // would be needless and this answers exactly the same question.
+    std::optional<unsigned int> cachedDbi(std::string_view collection);
 };
 
 }  // namespace smartbotic::db::storage

+ 136 - 0
tests/test_subdb_identity.cpp

@@ -1716,6 +1716,141 @@ void test_index_values() {
           "values, and returning the array would misstate what the key means");
 }
 
+// v2.11.0 T10 — the whole point of the txn-accepting overloads: one
+// transaction spanning TWO DIFFERENT collections (plus an index sub-db) must
+// be all-or-nothing. Task 12's atomic cascade depends on this - a parent
+// delete and its children's cleanup have to share one commit/abort, and the
+// no-txn put()/del() each open and commit their own WriteTxn, so they cannot
+// compose into one atomic unit at all.
+//
+// Written BEFORE the overloads exist, per the brief: this must fail to
+// compile/link until put(WriteTxn&, ...) and friends are added.
+void test_txn_accepting_writes_are_atomic_across_collections() {
+    TmpEnv t("txn-atomic");
+    LmdbDocumentStore store(t.env);
+    store.set_indexed_fields("parents", {"name"});
+
+    Document pd;
+    pd.id = "p1";
+    pd.collection = "parents";
+    pd.set_data(nlohmann::json{{"name", "alice"}});
+
+    Document cd;
+    cd.id = "c1";
+    cd.collection = "children";
+    cd.set_data(nlohmann::json{{"parentId", "p1"}});
+
+    // --- Abort path: neither document, nor the index entry, may survive. ---
+    {
+        WriteTxn wtxn(t.env);
+        std::vector<std::pair<std::string, unsigned int>> to_cache;
+        store.put(wtxn, "parents", "p1", pd, to_cache);
+        store.put(wtxn, "children", "c1", cd, to_cache);
+        wtxn.abort();
+        // to_cache is deliberately NOT applied to the process-wide dbi cache -
+        // see the header comment: caching before commit is exactly the v2.8.0
+        // bug. Nothing here should call cacheCommittedDbi.
+    }
+    check(!store.get("parents", "p1").has_value(),
+          "aborted txn: the parent document does not exist");
+    check(!store.get("children", "c1").has_value(),
+          "aborted txn: the child document does not exist");
+    auto idx = store.index_lookup_eq("parents", "name", nlohmann::json("alice"));
+    check(!idx.has_value() || idx->empty(),
+          "aborted txn: no index entry survives for the never-committed parent");
+
+    // --- The collection must not be poisoned by the aborted transaction: a
+    // later write, scan, count and get on EITHER collection must still work.
+    // This is precisely the v2.8.0 failure mode - caching a handle before
+    // commit leaves a closed handle behind an aborted caller transaction. ---
+    bool ok_put = true;
+    try {
+        store.put("parents", "p2", pd);
+    } catch (const std::exception&) {
+        ok_put = false;
+    }
+    check(ok_put, "a later write to 'parents' after the abort still works");
+
+    bool ok_put2 = true;
+    try {
+        store.put("children", "c2", cd);
+    } catch (const std::exception&) {
+        ok_put2 = false;
+    }
+    check(ok_put2, "a later write to 'children' after the abort still works");
+
+    check(store.get("parents", "p2").has_value(), "get() on 'parents' works after the abort");
+    check(store.get("children", "c2").has_value(), "get() on 'children' works after the abort");
+
+    smartbotic::database::Query q;
+    q.limit = 10;
+    check(store.scan("parents", q).documents.size() == 1,
+          "scan() on 'parents' works after the abort and sees only the post-abort write");
+    check(store.scan("children", q).documents.size() == 1,
+          "scan() on 'children' works after the abort and sees only the post-abort write");
+    check(store.count("parents") == 1, "count() on 'parents' is right after the abort");
+    check(store.count("children") == 1, "count() on 'children' is right after the abort");
+
+    // --- Commit path: both documents AND the index entry land together,
+    // once the caller re-primes to pick up the handles it chose not to
+    // cache itself (see the next block for why that step is mandatory in
+    // real callers). ---
+    {
+        WriteTxn wtxn(t.env);
+        std::vector<std::pair<std::string, unsigned int>> to_cache;
+        Document pd2;
+        pd2.id = "p3";
+        pd2.collection = "parents";
+        pd2.set_data(nlohmann::json{{"name", "bob"}});
+        Document cd2;
+        cd2.id = "c3";
+        cd2.collection = "children";
+        cd2.set_data(nlohmann::json{{"parentId", "p3"}});
+        store.put(wtxn, "parents", "p3", pd2, to_cache);
+        store.put(wtxn, "children", "c3", cd2, to_cache);
+        wtxn.commit();
+        check(to_cache.size() >= 2,
+              "put(wtxn,...) reports every handle it opened (collection dbis, "
+              "plus the index dbi since 'parents' declares one) for the caller "
+              "to cache after its own commit");
+    }
+    store.prime_dbi_cache();  // stand-in for the caching step a real caller does
+    check(store.get("parents", "p3").has_value(), "committed txn: parent document exists");
+    check(store.get("children", "c3").has_value(), "committed txn: child document exists");
+    auto idx2 = store.index_lookup_eq("parents", "name", nlohmann::json("bob"));
+    check(idx2.has_value() && idx2->size() == 1 && (*idx2)[0] == "p3",
+          "committed txn: the index entry for the new parent is there");
+
+    // --- The contract that makes handle-return necessary: after a commit
+    // through the txn-accepting overload, the caller MUST cache the handles
+    // it was given, because reads are cache-only (see try_open_for_read's
+    // comment: a cache miss reads as "no such collection" by design, so
+    // that read transactions never have to call mdb_dbi_open). Skip that
+    // step and freshly-committed data becomes unreadable - not corrupted,
+    // just invisible - for the life of the process. Demonstrated on a
+    // brand-new collection name never touched before this point. ---
+    {
+        WriteTxn wtxn(t.env);
+        std::vector<std::pair<std::string, unsigned int>> to_cache;
+        Document nd;
+        nd.id = "n1";
+        nd.collection = "brandnew";
+        nd.set_data(nlohmann::json{{"x", 1}});
+        store.put(wtxn, "brandnew", "n1", nd, to_cache);
+        wtxn.commit();
+        check(!to_cache.empty(),
+              "put(wtxn,...) hands back the handle it opened for 'brandnew'");
+        // Deliberately do NOT cache it - that is the point of this block.
+    }
+    check(!store.get("brandnew", "n1").has_value(),
+          "without caching after commit, the freshly-committed row reads as "
+          "absent - a cold cache, not lost data (proven next)");
+    store.prime_dbi_cache();
+    check(store.get("brandnew", "n1").has_value(),
+          "after priming (which caches it), the same row is visible - "
+          "confirming the earlier miss was purely a cold cache");
+}
+
 }  // namespace
 
 int main() {
@@ -1729,6 +1864,7 @@ int main() {
     test_scan_limit_zero_reports_total();
     test_scan_fast_path_matches_general_path();
     test_aborted_write_does_not_poison_the_collection();
+    test_txn_accepting_writes_are_atomic_across_collections();
     test_filtered_scan_operator_matrix();
     test_concurrent_reads_do_not_rebind_cached_handles();
     test_existing_collection_readable_without_writing_first();