瀏覽代碼

feat(relations): maintain the reverse index on child writes (T3)

put() and del() now maintain the relation reverse index declared via
set_relations(), inside the same WriteTxn as the document write, using
the same set-difference discipline as maintainIndexes(): old references
come from the stored bytes via resolve_from_yyjson(), new references
from filter_eval::resolveFilterValue() on the incoming Document, diffed
so an unchanged reference costs no index write.

- RelationRef{name, childField} is a new light struct on the storage
  layer, deliberately independent of relations/relation_manager.hpp's
  RelationInfo (project qualification, OnDelete policy, etc. stay out
  of this layer, per relation_index.hpp's existing design note).
- extract_relation_ids() treats absent/null as no reference and an
  array value as one posting per string element; non-string elements
  are skipped rather than stringified, since ids must round-trip
  exactly for relation_index_remove to find them again.
- Sub-dbs created at write time are registered via cacheCommittedDbi()
  only after the caller's commit, matching every other runtime-created
  sub-db in this file.
- Collections with no relations declared pay nothing (checked before
  touching the stored bytes).

tests/test_relation_enforcement.cpp: the brief's test plus three more
(undeclared collection stays untouched, unrelated field update leaves
the posting alone, posting survives a fresh LmdbDocumentStore over the
same env). 11/11 passing.
fszontagh 1 月之前
父節點
當前提交
e9d17c7454

+ 142 - 4
service/src/storage/document_store_lmdb.cpp

@@ -195,6 +195,31 @@ std::optional<nlohmann::json> resolve_from_yyjson(yyjson_val* root,
     return to_json(cur);
 }
 
+// v2.11.0 T3 — the set of parent ids a reference field value names.
+//
+// Relation ids are raw strings (unlike secondary index keys, they need no
+// type-tagged encoding - see relation_index.hpp), so this is simpler than
+// encode_index_keys: a string value names one id, an array value names one
+// id per STRING element (non-string elements are not valid ids and are
+// skipped rather than stringified, since an id must round-trip exactly for
+// relation_index_remove to find it again). Absent/null must never reach
+// here as a reference - callers pass std::nullopt for those, which this
+// treats the same as "no ids".
+std::vector<std::string> extract_relation_ids(const std::optional<nlohmann::json>& v) {
+    std::vector<std::string> out;
+    if (!v || v->is_null()) return out;
+    if (v->is_string()) {
+        out.push_back(v->get<std::string>());
+    } else if (v->is_array()) {
+        for (const auto& el : *v) {
+            if (el.is_string()) out.push_back(el.get<std::string>());
+        }
+    }
+    std::sort(out.begin(), out.end());
+    out.erase(std::unique(out.begin(), out.end()), out.end());
+    return out;
+}
+
 void sort_documents(std::vector<smartbotic::database::Document>& docs,
                     const smartbotic::database::Sort& sort) {
     std::sort(docs.begin(), docs.end(),
@@ -485,13 +510,20 @@ void LmdbDocumentStore::put(std::string_view collection,
     // 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;
-    if (!indexed_fields(collection).empty()) {
+    const bool needs_index = !indexed_fields(collection).empty();
+    const bool needs_relations = !relations(collection).empty();
+    if (needs_index || needs_relations) {
         MDB_val old{0, nullptr};
         const int rc = mdb_get(wtxn.raw(), dbi, &k, &old);
         if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) throw_mdb(rc, "get (pre-index)");
         const std::string_view old_payload =
             rc == MDB_SUCCESS ? to_sv(old) : std::string_view{};
-        maintainIndexes(wtxn, collection, id, old_payload, &doc, index_dbis);
+        if (needs_index) {
+            maintainIndexes(wtxn, collection, id, old_payload, &doc, index_dbis);
+        }
+        if (needs_relations) {
+            maintainRelations(wtxn, collection, id, old_payload, &doc, index_dbis);
+        }
     }
 
     MDB_val v = to_val(payload);
@@ -642,6 +674,105 @@ void LmdbDocumentStore::maintainIndexes(
     }
 }
 
+// -------------------------------------------------------------------------
+// v2.11.0 T3 — relation reverse-index maintenance on child writes.
+// -------------------------------------------------------------------------
+
+void LmdbDocumentStore::set_relations(std::string_view collection,
+                                       std::vector<RelationRef> rels) {
+    std::lock_guard<std::mutex> lock(relations_mutex_);
+    if (rels.empty()) {
+        relations_.erase(std::string(collection));
+    } else {
+        relations_[std::string(collection)] = std::move(rels);
+    }
+}
+
+std::vector<RelationRef> LmdbDocumentStore::relations(std::string_view collection) {
+    std::lock_guard<std::mutex> lock(relations_mutex_);
+    auto it = relations_.find(std::string(collection));
+    if (it == relations_.end()) return {};
+    return it->second;
+}
+
+void LmdbDocumentStore::maintainRelations(
+    WriteTxn& wtxn,
+    std::string_view collection,
+    std::string_view id,
+    std::string_view old_payload,
+    const smartbotic::database::Document* new_doc,
+    std::vector<std::pair<std::string, unsigned int>>& to_cache) {
+
+    const auto rels = relations(collection);
+    if (rels.empty()) return;   // collections with no relations pay nothing
+
+    // Old references come from the STORED bytes, parsed with yyjson and
+    // resolved one field at a time - never materialised into a Document,
+    // for the same reason maintainIndexes avoids it (decode_document was
+    // measured at 88% of write cost on a large row).
+    std::unordered_map<std::string, std::vector<std::string>> old_ids_by_field;
+    if (!old_payload.empty()) {
+        yyjson_doc* d = yyjson_read(old_payload.data(), old_payload.size(), 0);
+        if (d) {
+            yyjson_val* root = yyjson_doc_get_root(d);
+            for (const auto& r : rels) {
+                if (old_ids_by_field.count(r.childField)) continue;
+                old_ids_by_field[r.childField] =
+                    extract_relation_ids(resolve_from_yyjson(root, r.childField));
+            }
+            yyjson_doc_free(d);
+        }
+    }
+
+    for (const auto& r : rels) {
+        std::vector<std::string> new_ids;
+        if (new_doc != nullptr) {
+            // resolveFilterValue is the SAME resolver queries/filters use, so a
+            // relation and a query can never disagree about which field a
+            // dot-path names.
+            new_ids = extract_relation_ids(
+                filter_eval::resolveFilterValue(*new_doc, r.childField));
+        }
+        std::vector<std::string> old_ids;
+        if (auto it = old_ids_by_field.find(r.childField); it != old_ids_by_field.end()) {
+            old_ids = it->second;
+        }
+
+        // Unchanged reference costs no index write - the common case for an
+        // update that touches other fields. Both sides are sorted/deduped by
+        // extract_relation_ids, so this comparison is exact.
+        if (old_ids == new_ids) continue;
+
+        std::vector<std::string> to_remove;
+        std::vector<std::string> to_add;
+        std::set_difference(old_ids.begin(), old_ids.end(), new_ids.begin(), new_ids.end(),
+                            std::back_inserter(to_remove));
+        std::set_difference(new_ids.begin(), new_ids.end(), old_ids.begin(), old_ids.end(),
+                            std::back_inserter(to_add));
+        if (to_remove.empty() && to_add.empty()) continue;
+
+        const std::string sub = relation_index_subdb(r.name);
+        const unsigned int dbi = open_for_write(wtxn, sub, MDB_DUPSORT);
+        to_cache.emplace_back(sub, dbi);
+
+        for (const auto& parentId : to_remove) {
+            MDB_val pk = to_val(parentId);
+            MDB_val cv = to_val(id);
+            // DUPSORT: passing the data removes just this (parent, child) pair -
+            // siblings of the same parent are untouched.
+            const int rc = mdb_del(wtxn.raw(), dbi, &pk, &cv);
+            if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) throw_mdb(rc, "relation index del");
+        }
+        for (const auto& parentId : to_add) {
+            MDB_val pk = to_val(parentId);
+            MDB_val cv = to_val(id);
+            // MDB_NODUPDATA makes a repeat put a no-op rather than an error.
+            const int rc = mdb_put(wtxn.raw(), dbi, &pk, &cv, MDB_NODUPDATA);
+            if (rc != MDB_SUCCESS && rc != MDB_KEYEXIST) throw_mdb(rc, "relation index put");
+        }
+    }
+}
+
 std::optional<uint64_t>
 LmdbDocumentStore::index_count_eq(std::string_view collection,
                                    const std::string& field,
@@ -1154,11 +1285,18 @@ bool LmdbDocumentStore::del(std::string_view collection, std::string_view 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;
-    if (!indexed_fields(collection).empty()) {
+    const bool needs_index = !indexed_fields(collection).empty();
+    const bool needs_relations = !relations(collection).empty();
+    if (needs_index || needs_relations) {
         MDB_val old{0, nullptr};
         const int grc = mdb_get(wtxn.raw(), dbi, &k, &old);
         if (grc == MDB_SUCCESS) {
-            maintainIndexes(wtxn, collection, id, to_sv(old), nullptr, index_dbis);
+            if (needs_index) {
+                maintainIndexes(wtxn, collection, id, to_sv(old), nullptr, index_dbis);
+            }
+            if (needs_relations) {
+                maintainRelations(wtxn, collection, id, to_sv(old), nullptr, index_dbis);
+            }
         } else if (grc != MDB_NOTFOUND) {
             throw_mdb(grc, "get (pre-index del)");
         }

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

@@ -46,6 +46,19 @@ public:
     std::string existing_id;
 };
 
+// v2.11.0 T3 — a child-side relation reference, as the storage layer needs
+// it to maintain the reverse index on a write. Deliberately NOT
+// RelationInfo (relations/relation_manager.hpp): this layer stays free of
+// project qualification, OnDelete policy, etc. - see relation_index.hpp's
+// file comment. `name` is the BARE relation name (the caller strips the
+// project); `childField` is the dot-path on the CHILD document (the
+// collection this ref is declared against) that holds the parent id, or an
+// array of them.
+struct RelationRef {
+    std::string name;
+    std::string childField;
+};
+
 class LmdbDocumentStore : public DocumentStore {
 public:
     explicit LmdbDocumentStore(LmdbEnv& env);
@@ -211,6 +224,15 @@ public:
                                                       std::string_view parentId,
                                                       size_t limit);
 
+    // v2.11.0 T3 — declare which relations have `collection` as their CHILD
+    // side, so put()/del() know to maintain the reverse index. Empty removes
+    // the declaration. Mirrors set_indexed_fields: maintenance runs inside
+    // the same write transaction as the document, and a collection with no
+    // relations declared pays nothing (the write path checks this map and
+    // returns immediately).
+    void set_relations(std::string_view collection, std::vector<RelationRef> rels);
+    std::vector<RelationRef> relations(std::string_view collection);
+
     // Outcome of the constructor's priming pass, for the owner to report.
     // prime_error() is empty on success.
     size_t primed_count() const noexcept { return primed_count_; }
@@ -268,6 +290,12 @@ private:
     std::mutex index_mutex_;
     std::unordered_map<std::string, std::vector<std::string>> indexed_fields_;
     std::unordered_map<std::string, std::vector<std::string>> unique_fields_;
+
+    // v2.11.0 T3 — collection (as CHILD) -> declared relation refs. Its own
+    // mutex for the same reason indexed_fields_ has one: the write path
+    // consults it on every put/del and must not contend with the dbi cache.
+    std::mutex relations_mutex_;
+    std::unordered_map<std::string, std::vector<RelationRef>> relations_;
     mutable std::atomic<uint64_t> indexed_scans_{0};
     mutable std::atomic<uint64_t> full_scans_{0};
     mutable std::atomic<uint64_t> declined_unselective_{0};
@@ -417,6 +445,20 @@ private:
                          const smartbotic::database::Document* new_doc,
                          std::vector<std::pair<std::string, unsigned int>>& to_cache);
 
+    // v2.11.0 T3 — bring every relation declared with `collection` as CHILD
+    // into line with a write, inside `wtxn`. Same shape as maintainIndexes:
+    // old references come from the stored bytes (yyjson only), new
+    // references from filter_eval::resolveFilterValue() on `new_doc` (null
+    // on delete), diffed with std::set_difference so an unchanged reference
+    // costs no write. Handles opened here are appended to `to_cache` for the
+    // caller to cache AFTER its commit.
+    void maintainRelations(class WriteTxn& wtxn,
+                           std::string_view collection,
+                           std::string_view id,
+                           std::string_view old_payload,
+                           const smartbotic::database::Document* new_doc,
+                           std::vector<std::pair<std::string, unsigned int>>& to_cache);
+
     // v2.8.0 — record a handle in the cache, to be called ONLY after the
     // transaction that opened it has committed. LMDB closes a handle whose
     // opening transaction aborts, so caching any earlier leaves a closed handle

+ 35 - 0
tests/CMakeLists.txt

@@ -594,6 +594,41 @@ endif()
 
 add_test(NAME test_relation_index COMMAND test_relation_index)
 
+# v2.11.0 T3 — relation reverse-index maintenance on child writes.
+# Exercises LmdbDocumentStore::put()/del() maintaining the reverse index
+# declared via set_relations(), inside the same write transaction as the
+# document. Same source list as test_relation_index: no relations/
+# relation_manager.cpp - the storage layer stays free of that dependency.
+add_executable(test_relation_enforcement
+    test_relation_enforcement.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/storage/lmdb_env.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/storage/lmdb_txn.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/storage/lmdb_dbi.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/storage/subdb_identity.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/storage/document_store_lmdb.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/storage/secondary_index.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/relations/relation_index.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/json_parse.cpp
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src/doc_binary.cpp
+)
+
+target_include_directories(test_relation_enforcement PRIVATE
+    ${CMAKE_CURRENT_SOURCE_DIR}/../service/src
+    ${LMDB_INCLUDE_DIR}
+    ${yyjson_INCLUDE_DIRS}
+)
+
+target_link_libraries(test_relation_enforcement PRIVATE ${LMDB_LIBRARY})
+target_link_libraries(test_relation_enforcement PRIVATE ${yyjson_LIBRARIES})
+
+if(TARGET nlohmann_json::nlohmann_json)
+    target_link_libraries(test_relation_enforcement PRIVATE nlohmann_json::nlohmann_json)
+else()
+    target_include_directories(test_relation_enforcement PRIVATE ${NLOHMANN_JSON_INCLUDE_DIRS})
+endif()
+
+add_test(NAME test_relation_enforcement COMMAND test_relation_enforcement)
+
 # v2.9.0 — secondary index key encoding. Pins the one property that matters:
 # an index key comparison must mean the same thing as the scan's comparison,
 # because two paths answering one question that disagree return wrong data

+ 164 - 0
tests/test_relation_enforcement.cpp

@@ -0,0 +1,164 @@
+// v2.11.0 T3 — maintain the relation reverse index on child writes.
+//
+// Storage-only: exercises LmdbDocumentStore::put()/del() maintaining the
+// reverse index declared via set_relations(), using the same set-difference
+// discipline as maintainIndexes(). No enforcement (Task 4) is involved here
+// - relation_index_child_count is used purely as the observation point.
+
+#include <atomic>
+#include <cstdio>
+#include <filesystem>
+#include <iostream>
+#include <string>
+#include <unistd.h>
+#include <vector>
+
+#include <nlohmann/json.hpp>
+
+#include "document.hpp"
+#include "storage/document_store_lmdb.hpp"
+#include "storage/lmdb_env.hpp"
+
+namespace fs = std::filesystem;
+
+using smartbotic::database::Document;
+using smartbotic::db::storage::LmdbDocumentStore;
+using smartbotic::db::storage::LmdbEnv;
+using smartbotic::db::storage::LmdbEnvOpts;
+using smartbotic::db::storage::RelationRef;
+
+namespace {
+
+int g_pass = 0;
+int g_fail = 0;
+
+void check(bool cond, const char* msg) {
+    if (cond) {
+        ++g_pass;
+    } else {
+        ++g_fail;
+        std::cerr << "FAIL: " << msg << "\n";
+    }
+}
+
+std::string make_tmpdir(const char* tag) {
+    static std::atomic<int> counter{0};
+    std::string path = "/tmp/relmaint-test-" + std::to_string(::getpid()) + "-" +
+                       std::to_string(counter.fetch_add(1)) + "-" + tag;
+    std::error_code ec;
+    fs::remove_all(path, ec);
+    return path;
+}
+
+struct TmpEnv {
+    std::string path;
+    LmdbEnv env;
+    explicit TmpEnv(const char* tag)
+        : path(make_tmpdir(tag)),
+          env(LmdbEnvOpts{path, 64ULL << 20, 256, 126, false}) {}
+    ~TmpEnv() {
+        std::error_code ec;
+        fs::remove_all(path, ec);
+    }
+    TmpEnv(const TmpEnv&) = delete;
+    TmpEnv& operator=(const TmpEnv&) = delete;
+};
+
+// -------------------------------------------------------------------------
+
+void test_index_follows_the_child_field() {
+    TmpEnv t("rel-maint");
+    LmdbDocumentStore store(t.env);
+    store.set_relations("executions", {{"exec_wf", "workflowId"}});
+
+    auto put = [&](const std::string& id, const nlohmann::json& data) {
+        Document d; d.id = id; d.collection = "executions"; d.set_data(data);
+        store.put("executions", id, d);
+    };
+
+    put("e1", {{"workflowId", "wf-1"}});
+    put("e2", {{"workflowId", "wf-1"}});
+    check(store.relation_index_child_count("exec_wf", "wf-1") == 2, "two children");
+
+    put("e2", {{"workflowId", "wf-2"}});          // re-point
+    check(store.relation_index_child_count("exec_wf", "wf-1") == 1, "left the old parent");
+    check(store.relation_index_child_count("exec_wf", "wf-2") == 1, "joined the new one");
+
+    store.del("executions", "e1");
+    check(store.relation_index_child_count("exec_wf", "wf-1") == 0, "delete removes it");
+
+    // Absent and null are NOT references: they never block a delete and never
+    // count as dangling.
+    put("e3", {{"other", 1}});
+    put("e4", {{"workflowId", nullptr}});
+    check(store.relation_index_child_count("exec_wf", "wf-1") == 0, "absent adds nothing");
+
+    // An ARRAY-valued reference contributes one posting per element.
+    store.set_relations("nodes", {{"node_creds", "config.credentialIds"}});
+    Document n; n.id = "n1"; n.collection = "nodes";
+    n.set_data({{"config", {{"credentialIds", {"c1", "c2"}}}}});
+    store.put("nodes", "n1", n);
+    check(store.relation_index_child_count("node_creds", "c1") == 1, "array element 1");
+    check(store.relation_index_child_count("node_creds", "c2") == 1, "array element 2");
+}
+
+// A collection with no relations declared must pay nothing and never touch
+// any relation sub-db, mirroring how maintainIndexes short-circuits.
+void test_undeclared_collection_maintains_nothing() {
+    TmpEnv t("rel-none");
+    LmdbDocumentStore store(t.env);
+
+    Document d; d.id = "x1"; d.collection = "misc"; d.set_data({{"workflowId", "wf-9"}});
+    store.put("misc", "x1", d);
+    check(store.relation_index_child_count("exec_wf", "wf-9") == 0,
+          "no relation declared on this collection means no posting written");
+}
+
+// Unrelated field updates must not touch the index - the set-difference is
+// exact, not "rewrite unconditionally".
+void test_unrelated_update_leaves_the_posting_alone() {
+    TmpEnv t("rel-stable");
+    LmdbDocumentStore store(t.env);
+    store.set_relations("executions", {{"exec_wf", "workflowId"}});
+
+    Document d1; d1.id = "e1"; d1.collection = "executions";
+    d1.set_data({{"workflowId", "wf-1"}, {"status", "running"}});
+    store.put("executions", "e1", d1);
+    check(store.relation_index_child_count("exec_wf", "wf-1") == 1, "initial posting");
+
+    Document d2; d2.id = "e1"; d2.collection = "executions";
+    d2.set_data({{"workflowId", "wf-1"}, {"status", "completed"}});
+    store.put("executions", "e1", d2);
+    check(store.relation_index_child_count("exec_wf", "wf-1") == 1,
+          "still one posting after an unrelated field changed");
+}
+
+// Survives across a fresh store instance over the same env - i.e. the reverse
+// index sub-db created at write time was registered via cacheCommittedDbi
+// after commit, not merely usable within the same process instance.
+void test_posting_visible_after_reopen() {
+    TmpEnv t("rel-reopen");
+    {
+        LmdbDocumentStore store(t.env);
+        store.set_relations("executions", {{"exec_wf", "workflowId"}});
+        Document d; d.id = "e1"; d.collection = "executions";
+        d.set_data({{"workflowId", "wf-1"}});
+        store.put("executions", "e1", d);
+    }
+    LmdbDocumentStore reopened(t.env);
+    check(reopened.relation_index_child_count("exec_wf", "wf-1") == 1,
+          "posting survives a fresh LmdbDocumentStore over the same env");
+}
+
+}  // namespace
+
+int main() {
+    std::cout << "=== test_relation_enforcement ===\n";
+    test_index_follows_the_child_field();
+    test_undeclared_collection_maintains_nothing();
+    test_unrelated_update_leaves_the_posting_alone();
+    test_posting_visible_after_reopen();
+
+    std::cout << "passed: " << g_pass << ", failed: " << g_fail << "\n";
+    return g_fail == 0 ? 0 : 1;
+}