| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070 |
- // v2.4.4 — sub-db identity sentinel tests.
- //
- // Regression cover for the production incident in which an invalidated
- // MDB_dbi was reused after LMDB reassigned its slot, so writes aimed at
- // `image_hashes` landed in `executions` and succeeded silently. 31 documents
- // across smartbotic-automation ended up in a sub-db other than the one they
- // declared. See storage/subdb_identity.hpp for the full mechanism.
- //
- // The core test is `misbound_handle_is_refused`: it reproduces the misbinding
- // directly by handing the verifier a handle for a different sub-db, which is
- // what a stale cache entry amounts to.
- #include <cassert>
- #include <atomic>
- #include <filesystem>
- #include <iostream>
- #include <cstdio>
- #include <algorithm>
- #include <thread>
- #include <set>
- #include <string>
- #include <unistd.h>
- #include <lmdb.h>
- #include <nlohmann/json.hpp>
- #include "document.hpp"
- #include "storage/document_store_lmdb.hpp"
- #include "storage/lmdb_env.hpp"
- #include "storage/lmdb_txn.hpp"
- #include "storage/subdb_identity.hpp"
- namespace fs = std::filesystem;
- using smartbotic::database::Document;
- using smartbotic::db::storage::is_identity_key;
- using smartbotic::db::storage::kSubdbIdentityKey;
- using smartbotic::db::storage::LmdbDocumentStore;
- using smartbotic::db::storage::LmdbEnv;
- using smartbotic::db::storage::LmdbEnvOpts;
- using smartbotic::db::storage::read_subdb_identity;
- using smartbotic::db::storage::ReadTxn;
- using smartbotic::db::storage::verify_subdb_identity;
- using smartbotic::db::storage::write_subdb_identity;
- using smartbotic::db::storage::WriteTxn;
- 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/subdb-identity-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;
- };
- // Open (creating) a named sub-db inside a write txn and return its handle.
- unsigned int open_subdb(WriteTxn& txn, const char* name) {
- MDB_dbi dbi = 0;
- int rc = mdb_dbi_open(txn.raw(), name, MDB_CREATE, &dbi);
- assert(rc == MDB_SUCCESS);
- (void)rc;
- return dbi;
- }
- Document make_doc(const std::string& id, const std::string& collection) {
- Document d;
- d.id = id;
- d.collection = collection;
- d.set_data(nlohmann::json{{"seenCount", 1}, {"who", collection}});
- return d;
- }
- // -------------------------------------------------------------------------
- void test_sentinel_roundtrip() {
- TmpEnv t("roundtrip");
- {
- WriteTxn w(t.env);
- unsigned int dbi = open_subdb(w, "image_hashes");
- write_subdb_identity(w, dbi, "image_hashes");
- w.commit();
- }
- {
- WriteTxn w(t.env);
- unsigned int dbi = open_subdb(w, "image_hashes");
- bool threw = false;
- try {
- verify_subdb_identity(w, dbi, "image_hashes");
- } catch (const std::exception&) {
- threw = true;
- }
- check(!threw, "matching sentinel must verify without throwing");
- w.commit();
- }
- {
- ReadTxn r(t.env);
- MDB_dbi dbi = 0;
- mdb_dbi_open(r.raw(), "image_hashes", 0, &dbi);
- check(read_subdb_identity(r, dbi) == "image_hashes",
- "read_subdb_identity returns the stamped name");
- }
- }
- // THE regression test. A stale cache entry is, in effect, a handle that
- // addresses someone else's sub-db. Hand the verifier exactly that.
- void test_misbound_handle_is_refused() {
- TmpEnv t("misbound");
- unsigned int executions_dbi = 0;
- {
- WriteTxn w(t.env);
- unsigned int ih = open_subdb(w, "image_hashes");
- write_subdb_identity(w, ih, "image_hashes");
- executions_dbi = open_subdb(w, "executions");
- write_subdb_identity(w, executions_dbi, "executions");
- w.commit();
- }
- WriteTxn w(t.env);
- // Re-open so the handle is valid in this txn, then deliberately verify it
- // under the WRONG name — the production misbinding, reproduced.
- unsigned int exec = open_subdb(w, "executions");
- bool threw = false;
- std::string msg;
- try {
- verify_subdb_identity(w, exec, "image_hashes");
- } catch (const std::exception& e) {
- threw = true;
- msg = e.what();
- }
- check(threw, "handle for 'executions' verified as 'image_hashes' must throw");
- check(msg.find("image_hashes") != std::string::npos &&
- msg.find("executions") != std::string::npos,
- "misbinding error names both the requested and actual sub-db");
- w.abort();
- }
- // Existing deployments have sub-dbs with no sentinel. Those must keep working.
- void test_unstamped_subdb_is_permitted() {
- TmpEnv t("unstamped");
- WriteTxn w(t.env);
- unsigned int dbi = open_subdb(w, "legacy");
- bool threw = false;
- try {
- verify_subdb_identity(w, dbi, "legacy");
- } catch (const std::exception&) {
- threw = true;
- }
- check(!threw, "sub-db without a sentinel must verify (absence is unknown, not wrong)");
- w.commit();
- }
- void test_identity_key_predicate() {
- check(is_identity_key(kSubdbIdentityKey), "sentinel key recognised");
- check(!is_identity_key("__subdb_identity__"),
- "same text without the leading NUL is NOT the sentinel");
- check(!is_identity_key("e63c1b90"), "a document id is not the sentinel");
- check(kSubdbIdentityKey[0] == '\0',
- "sentinel must start with NUL so it cannot collide with a doc id");
- }
- // The sentinel is an implementation detail: it must never surface through the
- // DocumentStore API as a document, nor inflate a count.
- void test_sentinel_invisible_through_store() {
- TmpEnv t("invisible");
- LmdbDocumentStore store(t.env);
- store.put("image_hashes", "aaa", make_doc("aaa", "image_hashes"));
- store.put("image_hashes", "bbb", make_doc("bbb", "image_hashes"));
- check(store.count("image_hashes") == 2,
- "count() must exclude the identity sentinel");
- smartbotic::database::Query q;
- q.limit = 100;
- auto res = store.scan("image_hashes", q);
- check(res.documents.size() == 2, "scan() must exclude the identity sentinel");
- check(res.total_matched == 2, "scan() total_matched must exclude the sentinel");
- for (const auto& d : res.documents) {
- check(d.id == "aaa" || d.id == "bbb",
- "scan() must not surface the sentinel as a document");
- }
- // And it really is on disk.
- ReadTxn r(t.env);
- MDB_dbi dbi = 0;
- int rc = mdb_dbi_open(r.raw(), "image_hashes", 0, &dbi);
- check(rc == MDB_SUCCESS, "sub-db exists");
- check(read_subdb_identity(r, dbi) == "image_hashes",
- "store.put() stamps the sentinel on first write");
- }
- // A vector sub-db gets stamped with its own (prefixed) name, and scan_vectors
- // must skip the sentinel rather than trying to read it as float32 bytes.
- void test_vector_subdb_sentinel() {
- TmpEnv t("vectors");
- LmdbDocumentStore store(t.env);
- store.put_vector("emb", "v1", {1.0f, 2.0f, 3.0f});
- store.put_vector("emb", "v2", {4.0f, 5.0f, 6.0f});
- int seen = 0;
- bool bad = false;
- store.scan_vectors("emb", [&](std::string_view id, const float*, size_t n) {
- ++seen;
- if (n != 3) bad = true;
- if (is_identity_key(id)) bad = true;
- });
- check(seen == 2, "scan_vectors must skip the sentinel");
- check(!bad, "scan_vectors must not decode the sentinel as float data");
- ReadTxn r(t.env);
- MDB_dbi dbi = 0;
- mdb_dbi_open(r.raw(), "_vectors_emb", 0, &dbi);
- check(read_subdb_identity(r, dbi) == "_vectors_emb",
- "vector sub-db is stamped with its prefixed name");
- }
- // v2.4.4 Count is LMDB-first. Unfiltered it uses count(); filtered it uses
- // scan() with limit=0 and reads total_matched. Pin that contract: total_matched
- // is computed BEFORE pagination, so limit=0 must still report the true total
- // while returning no documents.
- void test_scan_limit_zero_reports_total() {
- TmpEnv t("counting");
- LmdbDocumentStore store(t.env);
- for (int i = 0; i < 5; ++i) {
- Document d;
- d.id = "id" + std::to_string(i);
- d.collection = "things";
- d.set_data(nlohmann::json{{"kind", i < 3 ? "alpha" : "beta"}});
- store.put("things", d.id, d);
- }
- smartbotic::database::Query q;
- q.limit = 0;
- auto all = store.scan("things", q);
- check(all.total_matched == 5, "limit=0 reports the full total");
- check(all.documents.empty(), "limit=0 returns no documents");
- check(store.count("things") == 5, "unfiltered count matches");
- smartbotic::database::Query fq;
- fq.limit = 0;
- smartbotic::database::Filter f;
- f.field = "kind";
- f.op = smartbotic::database::FilterOp::EQ;
- f.value = "alpha";
- fq.filters.push_back(f);
- auto filtered = store.scan("things", fq);
- check(filtered.total_matched == 3,
- "filtered limit=0 reports the matching total, not the collection size");
- }
- // v2.7.1 — the unfiltered/unsorted fast path in scan() must agree with the
- // general path exactly. It exists because the general path decoded every
- // document in the collection to return `limit` of them, so cost tracked total
- // bytes rather than page size (382ms to return one 510-byte document from a
- // 414 MB collection on a live instance). Any divergence here is a paging bug.
- void test_scan_fast_path_matches_general_path() {
- TmpEnv t("fastpath");
- LmdbDocumentStore store(t.env);
- for (int i = 0; i < 25; ++i) {
- Document d;
- char buf[16];
- std::snprintf(buf, sizeof(buf), "id%02d", i);
- d.id = buf;
- d.collection = "things";
- d.set_data(nlohmann::json{{"n", i}, {"kind", i % 2 ? "odd" : "even"}});
- store.put("things", d.id, d);
- }
- // total_matched and has_more must match what a full count says.
- smartbotic::database::Query page;
- page.limit = 10;
- page.offset = 0;
- auto p0 = store.scan("things", page);
- check(p0.documents.size() == 10, "fast path returns exactly `limit` docs");
- check(p0.total_matched == 25, "fast path total_matched excludes the sentinel");
- check(p0.has_more, "has_more true when more remain");
- page.offset = 20;
- auto p2 = store.scan("things", page);
- check(p2.documents.size() == 5, "final page returns the remainder");
- check(p2.total_matched == 25, "total_matched stable across pages");
- check(!p2.has_more, "has_more false on the last page");
- page.offset = 25;
- auto p3 = store.scan("things", page);
- check(p3.documents.empty(), "offset past the end returns nothing");
- check(p3.total_matched == 25, "and still reports the true total");
- // Paging must cover every document exactly once, in a stable order.
- std::set<std::string> seen;
- for (uint32_t off = 0; off < 25; off += 7) {
- smartbotic::database::Query q;
- q.limit = 7;
- q.offset = off;
- for (const auto& d : store.scan("things", q).documents) seen.insert(d.id);
- }
- check(seen.size() == 25, "paging the whole collection yields every document once");
- // A filter forces the general path; it must still be correct.
- smartbotic::database::Query fq;
- fq.limit = 100;
- smartbotic::database::Filter f;
- f.field = "kind";
- f.op = smartbotic::database::FilterOp::EQ;
- f.value = "odd";
- fq.filters.push_back(f);
- auto filtered = store.scan("things", fq);
- check(filtered.total_matched == 12, "filtered path still counts matches, not rows");
- // A sort also forces the general path.
- smartbotic::database::Query sq;
- sq.limit = 3;
- sq.sort = smartbotic::database::Sort{"n", true};
- auto sorted = store.scan("things", sq);
- check(sorted.documents.size() == 3, "sorted path paginates");
- check(sorted.total_matched == 25, "sorted path totals all rows");
- check(sorted.documents[0].data().value("n", -1) == 24,
- "descending sort really sorted (fast path must not swallow sorts)");
- // limit=0 keeps meaning "no documents, but a true total" - the contract
- // Count depends on (see test_scan_limit_zero_reports_total).
- smartbotic::database::Query zq;
- zq.limit = 0;
- auto z = store.scan("things", zq);
- check(z.documents.empty(), "limit=0 returns no documents on the fast path");
- check(z.total_matched == 25, "limit=0 still reports the true total");
- }
- // v2.8.0 — a WRITE that aborts must not poison the collection.
- //
- // This is the v2.4.3 EINVAL bug in a third failure mode, observed live in
- // production on 2.7.1: `find` on smartbotic-automation:workflows failed with
- // "LMDB cursor_open: Invalid argument" on every attempt while every other
- // collection was fine, and a restart was the only cure.
- //
- // Cause: open_for_write() cached the MDB_dbi immediately after mdb_dbi_open,
- // BEFORE the caller committed. LMDB keeps a handle private to the opening
- // transaction until it commits and CLOSES it if that transaction aborts - so any
- // write that threw after the handle was cached (a failed mdb_put, a sentinel
- // mismatch, a WriteTxn destructing uncommitted) left a closed handle in the
- // cache, and every later operation on that collection failed EINVAL for the rest
- // of the process's life.
- //
- // The abort is induced honestly here, with a key past LMDB's 511-byte limit, so
- // the test exercises the same path a real failed write takes.
- void test_aborted_write_does_not_poison_the_collection() {
- TmpEnv t("abortpoison");
- LmdbDocumentStore store(t.env);
- // Force a write that opens the sub-db and then fails: an oversized key makes
- // mdb_put return MDB_BAD_VALSIZE, which throws, so the WriteTxn aborts.
- const std::string huge_id(600, 'k');
- bool threw = false;
- try {
- store.put("poisoned", huge_id, make_doc(huge_id, "poisoned"));
- } catch (const std::exception&) {
- threw = true;
- }
- check(threw, "an oversized key really does fail the write");
- // The collection must still be usable. Before the fix, every one of these
- // failed with EINVAL because the cache held a handle LMDB had closed.
- bool ok_put = true;
- try {
- store.put("poisoned", "good", make_doc("good", "poisoned"));
- } catch (const std::exception&) {
- ok_put = false;
- }
- check(ok_put, "a later WRITE to the same collection still works");
- bool ok_read = true;
- try {
- smartbotic::database::Query q;
- q.limit = 10;
- auto res = store.scan("poisoned", q);
- check(res.documents.size() == 1, "and the document written after the abort is there");
- } catch (const std::exception&) {
- ok_read = false;
- }
- check(ok_read, "a later SCAN of the same collection still works (cursor_open)");
- bool ok_count = true;
- try {
- check(store.count("poisoned") == 1, "count is right after the abort");
- } catch (const std::exception&) {
- ok_count = false;
- }
- check(ok_count, "and count() does not throw");
- check(store.get("poisoned", "good").has_value(), "get() works after the abort");
- }
- // v2.8.0 — the two-pass filtered scan must agree with the old row-at-a-time path
- // on every operator, not just the common ones.
- //
- // scan() now evaluates predicates against a yyjson tree via a field resolver and
- // materialises only the returned page, because building a Document per row was
- // 88% of a filtered query's cost (3168ms vs 369ms for the parse alone over 193 MB
- // of real rows). A resolver that mishandles one operator returns silently wrong
- // data, so this walks the matrix.
- void test_filtered_scan_operator_matrix() {
- TmpEnv t("filtermatrix");
- LmdbDocumentStore store(t.env);
- auto put = [&](const std::string& id, const nlohmann::json& data) {
- Document d;
- d.id = id;
- d.collection = "m";
- d.version = 3;
- d.createdAt = 1000;
- d.updatedAt = 2000;
- d.set_data(data);
- store.put("m", id, d);
- };
- put("a", {{"n", 1}, {"kind", "odd"}, {"tags", {"x", "y"}}, {"nest", {{"deep", "hit"}}}});
- put("b", {{"n", 2}, {"kind", "even"}, {"tags", {"y"}}, {"nest", {{"deep", "miss"}}}});
- put("c", {{"n", 3}, {"kind", "odd"}, {"tags", nlohmann::json::array()}});
- put("d", {{"n", 4}, {"kind", "even"}, {"extra", "present"}});
- auto ids = [&](const smartbotic::database::Query& q) {
- std::vector<std::string> out;
- for (const auto& d : store.scan("m", q).documents) out.push_back(d.id);
- std::sort(out.begin(), out.end());
- return out;
- };
- auto q1 = [&](const char* field, smartbotic::database::FilterOp op,
- const nlohmann::json& val) {
- smartbotic::database::Query q;
- q.limit = 100;
- smartbotic::database::Filter f;
- f.field = field; f.op = op; f.value = val;
- q.filters.push_back(f);
- return q;
- };
- using Op = smartbotic::database::FilterOp;
- check(ids(q1("kind", Op::EQ, "odd")) == (std::vector<std::string>{"a", "c"}),
- "EQ on a data field");
- check(ids(q1("kind", Op::NE, "odd")) == (std::vector<std::string>{"b", "d"}),
- "NE on a data field");
- check(ids(q1("n", Op::GT, 2)) == (std::vector<std::string>{"c", "d"}), "GT numeric");
- check(ids(q1("n", Op::GTE, 3)) == (std::vector<std::string>{"c", "d"}), "GTE numeric");
- check(ids(q1("n", Op::LT, 2)) == (std::vector<std::string>{"a"}), "LT numeric");
- check(ids(q1("n", Op::LTE, 2)) == (std::vector<std::string>{"a", "b"}), "LTE numeric");
- check(ids(q1("n", Op::IN, nlohmann::json::array({1, 4}))) ==
- (std::vector<std::string>{"a", "d"}), "IN");
- check(ids(q1("tags", Op::CONTAINS, "x")) == (std::vector<std::string>{"a"}),
- "CONTAINS descends into an array value");
- check(ids(q1("extra", Op::EXISTS, true)) == (std::vector<std::string>{"d"}),
- "EXISTS true");
- check(ids(q1("extra", Op::EXISTS, false)) ==
- (std::vector<std::string>{"a", "b", "c"}), "EXISTS false");
- check(ids(q1("kind", Op::REGEX, "^od")) == (std::vector<std::string>{"a", "c"}),
- "REGEX");
- check(ids(q1("nest.deep", Op::EQ, "hit")) == (std::vector<std::string>{"a"}),
- "dotted path descends into data");
- check(ids(q1("nest.missing", Op::EXISTS, true)).empty(),
- "a dotted path that does not resolve matches nothing");
- // Document metadata, which lives at the top level of the stored JSON rather
- // than inside "data".
- check(ids(q1("_id", Op::EQ, "b")) == (std::vector<std::string>{"b"}), "_id");
- check(ids(q1("_version", Op::EQ, 3)).size() == 4, "_version");
- check(ids(q1("_created_at", Op::GTE, 1000)).size() == 4, "_created_at");
- check(ids(q1("_updated_at", Op::LT, 2000)).empty(), "_updated_at");
- // SEARCH must still work - it needs the whole document, so it takes the old
- // path.
- check(ids(q1("", Op::SEARCH, "present")) == (std::vector<std::string>{"d"}),
- "SEARCH still matches (routed to the whole-document path)");
- check(ids(q1("", Op::SEARCH, "nothinghere")).empty(), "SEARCH non-match");
- // Sorting, pagination and total_matched over a filtered set.
- {
- smartbotic::database::Query q;
- q.limit = 1;
- smartbotic::database::Filter f;
- f.field = "kind"; f.op = Op::EQ; f.value = "odd";
- q.filters.push_back(f);
- q.sort = smartbotic::database::Sort{"n", true}; // descending
- auto page0 = store.scan("m", q);
- check(page0.total_matched == 2, "total_matched counts matches, not rows");
- check(page0.documents.size() == 1, "limit honoured");
- check(page0.documents[0].id == "c", "descending sort picks the highest first");
- check(page0.has_more, "has_more true mid-set");
- q.offset = 1;
- auto page1 = store.scan("m", q);
- check(page1.documents.size() == 1 && page1.documents[0].id == "a",
- "second page continues the sort order");
- check(!page1.has_more, "has_more false on the last page");
- q.offset = 5;
- check(store.scan("m", q).documents.empty(), "offset past the end is empty");
- }
- // Ascending, and a sort field that is missing from some documents.
- {
- smartbotic::database::Query q;
- q.limit = 10;
- q.sort = smartbotic::database::Sort{"extra", false};
- auto res = store.scan("m", q);
- check(res.total_matched == 4, "no filter plus a sort still totals every row");
- check(res.documents.size() == 4, "and returns them all");
- // sort_documents returns `descending` when the LEFT value is missing, so
- // ascending puts documents that HAVE the field first and the ones missing
- // it last. The two-pass path copies that rule rather than inventing one.
- check(res.documents.front().id == "d",
- "ascending: the document that has the sort field comes first");
- std::vector<std::string> tail;
- for (size_t i = 1; i < res.documents.size(); ++i) tail.push_back(res.documents[i].id);
- check(tail == (std::vector<std::string>{"a", "b", "c"}),
- "and the ones missing it follow, tie-broken by id");
- }
- }
- // v2.8.1 — concurrent reads must not rebind a cached handle.
- //
- // lmdb.h: "This function [mdb_dbi_open] must not be called from multiple
- // concurrent transactions in the same process. A transaction that uses this
- // function must finish (either commit or abort) before any other transaction in
- // the process may use this function."
- //
- // The old read path violated that on every read of a collection this process had
- // not yet written, from every gRPC thread at once. MDB_dbi is an index into the
- // env's shared handle table, so churning it silently REBOUND handles that
- // committed writes had already cached. Live on 2.8.0-3 that produced 41 refusals
- // in 45 minutes, all "caller asked for 'image_hashes' but the handle addresses
- // 'nsfw_images'" - and because each refusal bumped mirror drift, every read in
- // the process fell back to MemoryStore permanently.
- //
- // The fix primes all handles in one committed write txn at construction, so only
- // write txns ever call mdb_dbi_open and LMDB serialises those itself.
- //
- // HONEST LIMITATION: this test does NOT reproduce the race. Reverting the fix
- // leaves it green apart from the primed_count assertion - the torn slot-table
- // update needs an interleaving of concurrent mdb_dbi_open calls that 8 threads
- // over 10 collections does not reliably hit. What the test does pin is the
- // invariant the race violates (no read returns another collection's document, no
- // write is refused by the sentinel) plus the mechanism that removes the race:
- // handles are opened up front, so the read path has no mdb_dbi_open left to call.
- // The deterministic evidence is the documented contract in lmdb.h and the
- // production log.
- void test_concurrent_reads_do_not_rebind_cached_handles() {
- const std::string path = make_tmpdir("dbi-concurrent");
- const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
- struct Cleanup {
- const std::string& p;
- ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
- } cleanup{path};
- // Enough collections that slot churn has somewhere to go.
- const std::vector<std::string> colls = {
- "alpha", "beta", "gamma", "delta", "epsilon",
- "zeta", "eta", "theta", "iota", "kappa"};
- {
- LmdbEnv env(opts);
- LmdbDocumentStore writer(env);
- for (const auto& c : colls) {
- Document d;
- d.id = "seed";
- d.collection = c;
- d.set_data(nlohmann::json{{"who", c}});
- writer.put(c, "seed", d);
- }
- }
- // Restart: fresh env, so the shared handle table starts empty and the store
- // must prime it.
- LmdbEnv env2(opts);
- LmdbDocumentStore store(env2);
- check(store.prime_error().empty(), "priming succeeded on reopen");
- check(store.primed_count() >= colls.size(),
- "priming opened a handle for every existing sub-db");
- // Hammer reads from many threads. Every one of these used to call
- // mdb_dbi_open inside its own read txn.
- std::atomic<int> read_failures{0};
- std::atomic<int> wrong_data{0};
- {
- std::vector<std::thread> threads;
- for (int t = 0; t < 8; ++t) {
- threads.emplace_back([&, t]() {
- for (int i = 0; i < 40; ++i) {
- const auto& c = colls[(t + i) % colls.size()];
- try {
- auto got = store.get(c, "seed");
- if (!got) { ++read_failures; continue; }
- // A rebound handle reads a STRANGER's sub-db, so the
- // document that comes back belongs to another collection.
- if (got->data().value("who", std::string{}) != c) ++wrong_data;
- } catch (const std::exception&) {
- ++read_failures;
- }
- }
- });
- }
- for (auto& th : threads) th.join();
- }
- check(read_failures.load() == 0, "concurrent reads all succeeded");
- check(wrong_data.load() == 0,
- "no read returned another collection's document - a rebound handle "
- "addresses whichever sub-db now occupies its slot");
- // Now write to every collection through the cached handles. This is where
- // the live failure surfaced: the sentinel refused the write.
- int write_failures = 0;
- for (const auto& c : colls) {
- try {
- Document d;
- d.id = "after";
- d.collection = c;
- d.set_data(nlohmann::json{{"who", c}});
- store.put(c, "after", d);
- } catch (const std::exception&) {
- ++write_failures;
- }
- }
- check(write_failures == 0,
- "writes through primed handles are not refused by the identity "
- "sentinel - the refusal is what production saw");
- for (const auto& c : colls) {
- check(store.count(c) == 2, ("both documents readable in " + c).c_str());
- }
- }
- // A collection that exists on disk but has NOT been written by this process must
- // still be readable. The read path no longer opens handles on demand, so if
- // priming missed anything a read would report the collection as EMPTY - a
- // silent-wrong-data failure worse than the bug being fixed.
- void test_existing_collection_readable_without_writing_first() {
- const std::string path = make_tmpdir("dbi-prime-read");
- const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
- struct Cleanup {
- const std::string& p;
- ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
- } cleanup{path};
- {
- LmdbEnv env(opts);
- LmdbDocumentStore writer(env);
- Document d;
- d.id = "only";
- d.collection = "archive";
- d.set_data(nlohmann::json{{"kept", true}});
- writer.put("archive", "only", d);
- }
- LmdbEnv env2(opts);
- LmdbDocumentStore reader(env2);
- // Read-only access, never a write on this collection in this process.
- auto got = reader.get("archive", "only");
- check(got.has_value(), "a never-written-here collection is still readable");
- check(reader.count("archive") == 1, "count sees it");
- smartbotic::database::Query q; q.limit = 10;
- check(reader.scan("archive", q).documents.size() == 1, "scan sees it");
- // And a collection that genuinely does not exist still reads as absent.
- check(!reader.get("nosuch", "x").has_value(),
- "a missing collection is still absent, not an error");
- check(reader.count("nosuch") == 0, "and counts zero");
- }
- // v2.9.0 — a secondary index must stay exactly in step with the documents.
- //
- // The index is a second copy of a fact already stored in the row. Every way the
- // two can diverge is a silent-wrong-data bug: a stale entry returns a row that
- // no longer matches, a missing entry hides a row that does. So this walks the
- // full lifecycle - insert, update the indexed field, update something else,
- // delete, re-insert - and after every step asserts the index agrees with a
- // brute-force scan of the collection.
- void test_index_tracks_documents_through_every_write() {
- TmpEnv t("idx-maint");
- LmdbDocumentStore store(t.env);
- store.set_indexed_fields("execs", {"workflowId"});
- auto put = [&](const std::string& id, const std::string& wf, int n) {
- Document d;
- d.id = id;
- d.collection = "execs";
- d.set_data(nlohmann::json{{"workflowId", wf}, {"n", n}});
- store.put("execs", id, d);
- };
- // What the index SHOULD say, computed by scanning every row - the oracle.
- auto truth = [&](const std::string& wf) {
- smartbotic::database::Query q;
- q.limit = 10000;
- std::vector<std::string> ids;
- for (const auto& d : store.scan("execs", q).documents) {
- if (d.data().value("workflowId", std::string{}) == wf) ids.push_back(d.id);
- }
- std::sort(ids.begin(), ids.end());
- return ids;
- };
- auto indexed = [&](const std::string& wf) {
- auto got = store.index_lookup_eq("execs", "workflowId", nlohmann::json(wf));
- std::vector<std::string> ids = got.value_or(std::vector<std::string>{});
- std::sort(ids.begin(), ids.end());
- return ids;
- };
- auto agree = [&](const std::string& wf, const char* stage) {
- const bool ok = indexed(wf) == truth(wf);
- const std::string msg = "index agrees with a full scan for " + wf +
- " after " + stage;
- check(ok, msg.c_str());
- };
- // Insert
- put("e1", "wf-a", 1);
- put("e2", "wf-a", 2);
- put("e3", "wf-b", 3);
- agree("wf-a", "inserts");
- agree("wf-b", "inserts");
- check(indexed("wf-a").size() == 2, "two rows under wf-a");
- // Update the INDEXED field: the old posting must go, the new one appear.
- put("e2", "wf-b", 2);
- agree("wf-a", "moving e2 to wf-b");
- agree("wf-b", "moving e2 to wf-b");
- check(indexed("wf-a").size() == 1, "wf-a lost e2");
- check(indexed("wf-b").size() == 2, "wf-b gained it");
- // Update an UNindexed field: the index must be untouched, not duplicated.
- put("e1", "wf-a", 99);
- agree("wf-a", "updating an unindexed field");
- check(indexed("wf-a").size() == 1,
- "no duplicate posting from re-writing the same indexed value");
- // Delete
- check(store.del("execs", "e3"), "deleted e3");
- agree("wf-b", "deleting e3");
- check(indexed("wf-b").size() == 1, "e3 is gone from the index");
- // Re-insert the same id
- put("e3", "wf-b", 7);
- agree("wf-b", "re-inserting e3");
- check(indexed("wf-b").size() == 2, "e3 is back exactly once");
- // A value with no rows is an empty result, NOT "no index".
- auto none = store.index_lookup_eq("execs", "workflowId", nlohmann::json("wf-zzz"));
- check(none.has_value() && none->empty(),
- "an indexed field with no matching rows returns empty, not nullopt - "
- "nullopt means 'no index' and would send the caller to a scan");
- // An undeclared field has no index, and must say so rather than say 'none'.
- auto unindexed = store.index_lookup_eq("execs", "n", nlohmann::json(1));
- check(!unindexed.has_value(),
- "an unindexed field returns nullopt so the caller falls back to a scan "
- "instead of concluding there are no matches");
- }
- // A collection with no declared index must behave exactly as before, and pay
- // nothing. Also: declaring an index later must pick up the rows already there.
- void test_build_index_over_existing_rows() {
- TmpEnv t("idx-build");
- LmdbDocumentStore store(t.env);
- // Write BEFORE declaring the index.
- for (int i = 0; i < 20; ++i) {
- Document d;
- d.id = "d" + std::to_string(i);
- d.collection = "c";
- d.set_data(nlohmann::json{{"grp", i % 4 == 0 ? "hot" : "cold"}});
- store.put("c", d.id, d);
- }
- check(!store.index_lookup_eq("c", "grp", nlohmann::json("hot")).has_value(),
- "no index exists before it is declared");
- const uint64_t built = store.build_index("c", "grp");
- check(built == 20, "the backfill indexed every existing row");
- store.set_indexed_fields("c", {"grp"});
- auto hot = store.index_lookup_eq("c", "grp", nlohmann::json("hot"));
- check(hot.has_value() && hot->size() == 5,
- "the backfilled index finds the pre-existing rows (d0,d4,d8,d12,d16)");
- // Idempotent: a second build must not double the postings.
- store.build_index("c", "grp");
- hot = store.index_lookup_eq("c", "grp", nlohmann::json("hot"));
- check(hot.has_value() && hot->size() == 5,
- "re-running the backfill does not duplicate postings");
- // The count guard must see the same number without reading the ids.
- auto n = store.index_count_eq("c", "grp", nlohmann::json("hot"));
- check(n.has_value() && *n == 5, "index_count_eq agrees with the lookup");
- auto cold = store.index_count_eq("c", "grp", nlohmann::json("cold"));
- check(cold.has_value() && *cold == 15, "and counts the larger group");
- check(store.drop_index("c", "grp"), "the index drops");
- check(!store.index_lookup_eq("c", "grp", nlohmann::json("hot")).has_value(),
- "after dropping, lookups report no index rather than no rows");
- }
- // Numbers are where an index most easily disagrees with a scan, because the scan
- // compares across numeric subtypes. A doc stored with 5 must be found by a
- // filter asking for 5.0 through EITHER path.
- void test_index_numeric_equality_matches_scan() {
- TmpEnv t("idx-num");
- LmdbDocumentStore store(t.env);
- store.set_indexed_fields("m", {"code"});
- Document a;
- a.id = "a"; a.collection = "m";
- a.set_data(nlohmann::json{{"code", 5}}); // integer
- store.put("m", "a", a);
- Document b;
- b.id = "b"; b.collection = "m";
- b.set_data(nlohmann::json{{"code", 5.0}}); // integral double
- store.put("m", "b", b);
- auto by_int = store.index_lookup_eq("m", "code", nlohmann::json(5));
- auto by_dbl = store.index_lookup_eq("m", "code", nlohmann::json(5.0));
- check(by_int.has_value() && by_int->size() == 2,
- "asking for 5 finds BOTH the int and the integral-double row");
- check(by_dbl == by_int,
- "and asking for 5.0 returns exactly the same rows - the scan's EQ "
- "compares numbers across subtypes, so the index must too");
- // A big integer must not collide with its neighbour via double precision.
- Document c;
- c.id = "c"; c.collection = "m";
- c.set_data(nlohmann::json{{"code", 1786263002080195076LL}});
- store.put("m", "c", c);
- auto near = store.index_lookup_eq("m", "code",
- nlohmann::json(1786263002080195077LL));
- check(near.has_value() && near->empty(),
- "a neighbouring ns-scale integer does not collide - these two ARE "
- "equal as doubles, so routing through double would false-match");
- }
- // v2.9.0 — an indexed plan must return EXACTLY what the unindexed plan returns.
- //
- // This is the whole safety argument for indexing. The index is an optimisation,
- // so any observable difference is a bug, and the interesting failures are silent:
- // a missing posting drops a row, a stale one adds a row that no longer matches,
- // and a different code path can disagree about ordering or total_matched.
- //
- // The test runs each query twice against the same data - once with the field
- // declared indexed, once not - and compares the complete result: ids in order,
- // total_matched, and has_more.
- void test_indexed_and_unindexed_plans_agree() {
- TmpEnv t("idx-equiv");
- LmdbDocumentStore store(t.env);
- // 300 rows: `grp` is selective enough to use the index (10 groups of 30 =
- // 10%), `bucket` deliberately is NOT (2 values, 50% each) so the guard has
- // something to decline.
- for (int i = 0; i < 300; ++i) {
- Document d;
- d.id = "r" + std::string(i < 10 ? "00" : (i < 100 ? "0" : "")) +
- std::to_string(i);
- d.collection = "c";
- d.set_data(nlohmann::json{
- {"grp", "g" + std::to_string(i % 10)},
- {"bucket", (i % 2 == 0) ? "even" : "odd"},
- {"n", i},
- {"nest", {{"deep", "d" + std::to_string(i % 10)}}},
- });
- store.put("c", d.id, d);
- }
- using Op = smartbotic::database::FilterOp;
- struct Case {
- const char* name;
- std::vector<smartbotic::database::Filter> filters;
- std::optional<smartbotic::database::Sort> sort;
- uint32_t limit;
- uint32_t offset;
- };
- auto F = [](const char* f, Op op, const nlohmann::json& v) {
- smartbotic::database::Filter x;
- x.field = f; x.op = op; x.value = v;
- return x;
- };
- const std::vector<Case> cases = {
- {"eq indexed field", {F("grp", Op::EQ, "g3")}, std::nullopt, 100, 0},
- {"eq + second predicate", {F("grp", Op::EQ, "g3"), F("bucket", Op::EQ, "even")}, std::nullopt, 100, 0},
- {"eq + range on another field", {F("grp", Op::EQ, "g3"), F("n", Op::GT, 100)}, std::nullopt, 100, 0},
- {"eq with sort asc", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", false}, 100, 0},
- {"eq with sort desc", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", true}, 100, 0},
- {"eq paginated", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", false}, 7, 10},
- {"eq offset past end", {F("grp", Op::EQ, "g3")}, std::nullopt, 10, 999},
- {"eq matching nothing", {F("grp", Op::EQ, "nope")}, std::nullopt, 100, 0},
- {"eq on unselective field", {F("bucket", Op::EQ, "even")}, std::nullopt, 100, 0},
- {"eq plus SEARCH", {F("grp", Op::EQ, "g3"), F("", Op::SEARCH, "g3")}, std::nullopt, 100, 0},
- {"ne on indexed field", {F("grp", Op::NE, "g3")}, std::nullopt, 100, 0},
- {"eq on nested path", {F("nest.deep", Op::EQ, "d4")}, std::nullopt, 100, 0},
- {"limit zero", {F("grp", Op::EQ, "g3")}, std::nullopt, 0, 0},
- };
- auto run = [&](const Case& c) {
- smartbotic::database::Query q;
- q.filters = c.filters;
- q.sort = c.sort;
- q.limit = c.limit;
- q.offset = c.offset;
- auto r = store.scan("c", q);
- std::string sig = "total=" + std::to_string(r.total_matched) +
- " more=" + std::to_string(r.has_more ? 1 : 0) + " [";
- for (const auto& d : r.documents) { sig += d.id; sig += ","; }
- sig += "]";
- return sig;
- };
- for (const auto& c : cases) {
- store.set_indexed_fields("c", {}); // no index
- const std::string without = run(c);
- store.set_indexed_fields("c", {"grp", "bucket", "nest.deep"});
- store.build_index("c", "grp");
- store.build_index("c", "bucket");
- store.build_index("c", "nest.deep");
- const std::string with = run(c);
- const std::string msg = std::string("indexed and unindexed plans agree: ")
- + c.name;
- if (with != without) {
- std::cerr << " without index: " << without << "\n"
- << " with index: " << with << "\n";
- }
- check(with == without, msg.c_str());
- }
- // The agreement above is only meaningful if the index plan was actually
- // TAKEN for the selective cases. Otherwise the planner declined every time
- // and the test compared the scan against itself.
- store.set_indexed_fields("c", {"grp", "bucket", "nest.deep"});
- {
- store.reset_index_plan_stats();
- smartbotic::database::Query q;
- q.limit = 100;
- q.filters.push_back(F("grp", Op::EQ, "g3"));
- auto r = store.scan("c", q);
- auto st = store.index_plan_stats();
- check(r.total_matched == 30, "the selective query matches 30 of 300 rows");
- check(st.indexed_scans == 1 && st.full_scans == 0,
- "a selective EQ on an indexed field TAKES the index plan - without "
- "this the equivalence cases above would prove nothing");
- }
- {
- store.reset_index_plan_stats();
- smartbotic::database::Query q;
- q.limit = 100;
- q.filters.push_back(F("bucket", Op::EQ, "even"));
- auto r = store.scan("c", q);
- auto st = store.index_plan_stats();
- check(r.total_matched == 150, "the unselective query matches half the rows");
- check(st.indexed_scans == 0 && st.declined_unselective == 1,
- "and the guard DECLINES its index - 150 of 300 rows would cost more "
- "through the index than a scan, so declaring an index must not be "
- "able to pessimise a query");
- }
- {
- // An index on a field the query does not filter on must not be consulted.
- store.reset_index_plan_stats();
- smartbotic::database::Query q;
- q.limit = 100;
- q.filters.push_back(F("n", Op::GT, 250));
- store.scan("c", q);
- auto st = store.index_plan_stats();
- check(st.full_scans == 1 && st.indexed_scans == 0,
- "a query whose predicates name no indexed field scans");
- }
- auto n = store.index_count_eq("c", "bucket", nlohmann::json("even"));
- check(n.has_value() && *n == 150,
- "the unselective index does exist and holds 150 of 300 rows");
- }
- // An empty document id is a zero-length LMDB key, which mdb_get rejects with
- // MDB_BAD_VALSIZE. That surfaced in production as a gRPC INTERNAL and an ERROR
- // log line every time a consumer asked for one - noise that masks real failures.
- // A key that cannot exist is absent, not an error.
- void test_empty_id_reads_as_absent() {
- TmpEnv t("empty-id");
- LmdbDocumentStore store(t.env);
- Document d;
- d.id = "real";
- d.collection = "c";
- d.set_data(nlohmann::json{{"x", 1}});
- store.put("c", "real", d);
- bool threw = false;
- try {
- check(!store.get("c", "").has_value(), "an empty id reads as absent");
- } catch (const std::exception&) {
- threw = true;
- }
- check(!threw, "and does NOT throw - MDB_BAD_VALSIZE became a gRPC INTERNAL");
- threw = false;
- try {
- check(!store.del("c", ""), "deleting an empty id is a no-op");
- } catch (const std::exception&) {
- threw = true;
- }
- check(!threw, "and does not throw either");
- check(store.get("c", "real").has_value(), "real ids still work");
- check(store.count("c") == 1, "and nothing was disturbed");
- }
- } // namespace
- int main() {
- std::cout << "=== test_subdb_identity ===\n";
- test_sentinel_roundtrip();
- test_misbound_handle_is_refused();
- test_unstamped_subdb_is_permitted();
- test_identity_key_predicate();
- test_sentinel_invisible_through_store();
- test_vector_subdb_sentinel();
- test_scan_limit_zero_reports_total();
- test_scan_fast_path_matches_general_path();
- test_aborted_write_does_not_poison_the_collection();
- test_filtered_scan_operator_matrix();
- test_concurrent_reads_do_not_rebind_cached_handles();
- test_existing_collection_readable_without_writing_first();
- test_index_tracks_documents_through_every_write();
- test_build_index_over_existing_rows();
- test_index_numeric_equality_matches_scan();
- test_indexed_and_unindexed_plans_agree();
- test_empty_id_reads_as_absent();
- std::cout << "passed: " << g_pass << ", failed: " << g_fail << "\n";
- return g_fail == 0 ? 0 : 1;
- }
|