test_subdb_identity.cpp 42 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035
  1. // v2.4.4 — sub-db identity sentinel tests.
  2. //
  3. // Regression cover for the production incident in which an invalidated
  4. // MDB_dbi was reused after LMDB reassigned its slot, so writes aimed at
  5. // `image_hashes` landed in `executions` and succeeded silently. 31 documents
  6. // across smartbotic-automation ended up in a sub-db other than the one they
  7. // declared. See storage/subdb_identity.hpp for the full mechanism.
  8. //
  9. // The core test is `misbound_handle_is_refused`: it reproduces the misbinding
  10. // directly by handing the verifier a handle for a different sub-db, which is
  11. // what a stale cache entry amounts to.
  12. #include <cassert>
  13. #include <atomic>
  14. #include <filesystem>
  15. #include <iostream>
  16. #include <cstdio>
  17. #include <algorithm>
  18. #include <thread>
  19. #include <set>
  20. #include <string>
  21. #include <unistd.h>
  22. #include <lmdb.h>
  23. #include <nlohmann/json.hpp>
  24. #include "document.hpp"
  25. #include "storage/document_store_lmdb.hpp"
  26. #include "storage/lmdb_env.hpp"
  27. #include "storage/lmdb_txn.hpp"
  28. #include "storage/subdb_identity.hpp"
  29. namespace fs = std::filesystem;
  30. using smartbotic::database::Document;
  31. using smartbotic::db::storage::is_identity_key;
  32. using smartbotic::db::storage::kSubdbIdentityKey;
  33. using smartbotic::db::storage::LmdbDocumentStore;
  34. using smartbotic::db::storage::LmdbEnv;
  35. using smartbotic::db::storage::LmdbEnvOpts;
  36. using smartbotic::db::storage::read_subdb_identity;
  37. using smartbotic::db::storage::ReadTxn;
  38. using smartbotic::db::storage::verify_subdb_identity;
  39. using smartbotic::db::storage::write_subdb_identity;
  40. using smartbotic::db::storage::WriteTxn;
  41. namespace {
  42. int g_pass = 0;
  43. int g_fail = 0;
  44. void check(bool cond, const char* msg) {
  45. if (cond) {
  46. ++g_pass;
  47. } else {
  48. ++g_fail;
  49. std::cerr << "FAIL: " << msg << "\n";
  50. }
  51. }
  52. std::string make_tmpdir(const char* tag) {
  53. static std::atomic<int> counter{0};
  54. std::string path = "/tmp/subdb-identity-test-" + std::to_string(::getpid()) +
  55. "-" + std::to_string(counter.fetch_add(1)) + "-" + tag;
  56. std::error_code ec;
  57. fs::remove_all(path, ec);
  58. return path;
  59. }
  60. struct TmpEnv {
  61. std::string path;
  62. LmdbEnv env;
  63. explicit TmpEnv(const char* tag)
  64. : path(make_tmpdir(tag)),
  65. env(LmdbEnvOpts{path, 64ULL << 20, 256, 126, false}) {}
  66. ~TmpEnv() {
  67. std::error_code ec;
  68. fs::remove_all(path, ec);
  69. }
  70. TmpEnv(const TmpEnv&) = delete;
  71. TmpEnv& operator=(const TmpEnv&) = delete;
  72. };
  73. // Open (creating) a named sub-db inside a write txn and return its handle.
  74. unsigned int open_subdb(WriteTxn& txn, const char* name) {
  75. MDB_dbi dbi = 0;
  76. int rc = mdb_dbi_open(txn.raw(), name, MDB_CREATE, &dbi);
  77. assert(rc == MDB_SUCCESS);
  78. (void)rc;
  79. return dbi;
  80. }
  81. Document make_doc(const std::string& id, const std::string& collection) {
  82. Document d;
  83. d.id = id;
  84. d.collection = collection;
  85. d.set_data(nlohmann::json{{"seenCount", 1}, {"who", collection}});
  86. return d;
  87. }
  88. // -------------------------------------------------------------------------
  89. void test_sentinel_roundtrip() {
  90. TmpEnv t("roundtrip");
  91. {
  92. WriteTxn w(t.env);
  93. unsigned int dbi = open_subdb(w, "image_hashes");
  94. write_subdb_identity(w, dbi, "image_hashes");
  95. w.commit();
  96. }
  97. {
  98. WriteTxn w(t.env);
  99. unsigned int dbi = open_subdb(w, "image_hashes");
  100. bool threw = false;
  101. try {
  102. verify_subdb_identity(w, dbi, "image_hashes");
  103. } catch (const std::exception&) {
  104. threw = true;
  105. }
  106. check(!threw, "matching sentinel must verify without throwing");
  107. w.commit();
  108. }
  109. {
  110. ReadTxn r(t.env);
  111. MDB_dbi dbi = 0;
  112. mdb_dbi_open(r.raw(), "image_hashes", 0, &dbi);
  113. check(read_subdb_identity(r, dbi) == "image_hashes",
  114. "read_subdb_identity returns the stamped name");
  115. }
  116. }
  117. // THE regression test. A stale cache entry is, in effect, a handle that
  118. // addresses someone else's sub-db. Hand the verifier exactly that.
  119. void test_misbound_handle_is_refused() {
  120. TmpEnv t("misbound");
  121. unsigned int executions_dbi = 0;
  122. {
  123. WriteTxn w(t.env);
  124. unsigned int ih = open_subdb(w, "image_hashes");
  125. write_subdb_identity(w, ih, "image_hashes");
  126. executions_dbi = open_subdb(w, "executions");
  127. write_subdb_identity(w, executions_dbi, "executions");
  128. w.commit();
  129. }
  130. WriteTxn w(t.env);
  131. // Re-open so the handle is valid in this txn, then deliberately verify it
  132. // under the WRONG name — the production misbinding, reproduced.
  133. unsigned int exec = open_subdb(w, "executions");
  134. bool threw = false;
  135. std::string msg;
  136. try {
  137. verify_subdb_identity(w, exec, "image_hashes");
  138. } catch (const std::exception& e) {
  139. threw = true;
  140. msg = e.what();
  141. }
  142. check(threw, "handle for 'executions' verified as 'image_hashes' must throw");
  143. check(msg.find("image_hashes") != std::string::npos &&
  144. msg.find("executions") != std::string::npos,
  145. "misbinding error names both the requested and actual sub-db");
  146. w.abort();
  147. }
  148. // Existing deployments have sub-dbs with no sentinel. Those must keep working.
  149. void test_unstamped_subdb_is_permitted() {
  150. TmpEnv t("unstamped");
  151. WriteTxn w(t.env);
  152. unsigned int dbi = open_subdb(w, "legacy");
  153. bool threw = false;
  154. try {
  155. verify_subdb_identity(w, dbi, "legacy");
  156. } catch (const std::exception&) {
  157. threw = true;
  158. }
  159. check(!threw, "sub-db without a sentinel must verify (absence is unknown, not wrong)");
  160. w.commit();
  161. }
  162. void test_identity_key_predicate() {
  163. check(is_identity_key(kSubdbIdentityKey), "sentinel key recognised");
  164. check(!is_identity_key("__subdb_identity__"),
  165. "same text without the leading NUL is NOT the sentinel");
  166. check(!is_identity_key("e63c1b90"), "a document id is not the sentinel");
  167. check(kSubdbIdentityKey[0] == '\0',
  168. "sentinel must start with NUL so it cannot collide with a doc id");
  169. }
  170. // The sentinel is an implementation detail: it must never surface through the
  171. // DocumentStore API as a document, nor inflate a count.
  172. void test_sentinel_invisible_through_store() {
  173. TmpEnv t("invisible");
  174. LmdbDocumentStore store(t.env);
  175. store.put("image_hashes", "aaa", make_doc("aaa", "image_hashes"));
  176. store.put("image_hashes", "bbb", make_doc("bbb", "image_hashes"));
  177. check(store.count("image_hashes") == 2,
  178. "count() must exclude the identity sentinel");
  179. smartbotic::database::Query q;
  180. q.limit = 100;
  181. auto res = store.scan("image_hashes", q);
  182. check(res.documents.size() == 2, "scan() must exclude the identity sentinel");
  183. check(res.total_matched == 2, "scan() total_matched must exclude the sentinel");
  184. for (const auto& d : res.documents) {
  185. check(d.id == "aaa" || d.id == "bbb",
  186. "scan() must not surface the sentinel as a document");
  187. }
  188. // And it really is on disk.
  189. ReadTxn r(t.env);
  190. MDB_dbi dbi = 0;
  191. int rc = mdb_dbi_open(r.raw(), "image_hashes", 0, &dbi);
  192. check(rc == MDB_SUCCESS, "sub-db exists");
  193. check(read_subdb_identity(r, dbi) == "image_hashes",
  194. "store.put() stamps the sentinel on first write");
  195. }
  196. // A vector sub-db gets stamped with its own (prefixed) name, and scan_vectors
  197. // must skip the sentinel rather than trying to read it as float32 bytes.
  198. void test_vector_subdb_sentinel() {
  199. TmpEnv t("vectors");
  200. LmdbDocumentStore store(t.env);
  201. store.put_vector("emb", "v1", {1.0f, 2.0f, 3.0f});
  202. store.put_vector("emb", "v2", {4.0f, 5.0f, 6.0f});
  203. int seen = 0;
  204. bool bad = false;
  205. store.scan_vectors("emb", [&](std::string_view id, const float*, size_t n) {
  206. ++seen;
  207. if (n != 3) bad = true;
  208. if (is_identity_key(id)) bad = true;
  209. });
  210. check(seen == 2, "scan_vectors must skip the sentinel");
  211. check(!bad, "scan_vectors must not decode the sentinel as float data");
  212. ReadTxn r(t.env);
  213. MDB_dbi dbi = 0;
  214. mdb_dbi_open(r.raw(), "_vectors_emb", 0, &dbi);
  215. check(read_subdb_identity(r, dbi) == "_vectors_emb",
  216. "vector sub-db is stamped with its prefixed name");
  217. }
  218. // v2.4.4 Count is LMDB-first. Unfiltered it uses count(); filtered it uses
  219. // scan() with limit=0 and reads total_matched. Pin that contract: total_matched
  220. // is computed BEFORE pagination, so limit=0 must still report the true total
  221. // while returning no documents.
  222. void test_scan_limit_zero_reports_total() {
  223. TmpEnv t("counting");
  224. LmdbDocumentStore store(t.env);
  225. for (int i = 0; i < 5; ++i) {
  226. Document d;
  227. d.id = "id" + std::to_string(i);
  228. d.collection = "things";
  229. d.set_data(nlohmann::json{{"kind", i < 3 ? "alpha" : "beta"}});
  230. store.put("things", d.id, d);
  231. }
  232. smartbotic::database::Query q;
  233. q.limit = 0;
  234. auto all = store.scan("things", q);
  235. check(all.total_matched == 5, "limit=0 reports the full total");
  236. check(all.documents.empty(), "limit=0 returns no documents");
  237. check(store.count("things") == 5, "unfiltered count matches");
  238. smartbotic::database::Query fq;
  239. fq.limit = 0;
  240. smartbotic::database::Filter f;
  241. f.field = "kind";
  242. f.op = smartbotic::database::FilterOp::EQ;
  243. f.value = "alpha";
  244. fq.filters.push_back(f);
  245. auto filtered = store.scan("things", fq);
  246. check(filtered.total_matched == 3,
  247. "filtered limit=0 reports the matching total, not the collection size");
  248. }
  249. // v2.7.1 — the unfiltered/unsorted fast path in scan() must agree with the
  250. // general path exactly. It exists because the general path decoded every
  251. // document in the collection to return `limit` of them, so cost tracked total
  252. // bytes rather than page size (382ms to return one 510-byte document from a
  253. // 414 MB collection on a live instance). Any divergence here is a paging bug.
  254. void test_scan_fast_path_matches_general_path() {
  255. TmpEnv t("fastpath");
  256. LmdbDocumentStore store(t.env);
  257. for (int i = 0; i < 25; ++i) {
  258. Document d;
  259. char buf[16];
  260. std::snprintf(buf, sizeof(buf), "id%02d", i);
  261. d.id = buf;
  262. d.collection = "things";
  263. d.set_data(nlohmann::json{{"n", i}, {"kind", i % 2 ? "odd" : "even"}});
  264. store.put("things", d.id, d);
  265. }
  266. // total_matched and has_more must match what a full count says.
  267. smartbotic::database::Query page;
  268. page.limit = 10;
  269. page.offset = 0;
  270. auto p0 = store.scan("things", page);
  271. check(p0.documents.size() == 10, "fast path returns exactly `limit` docs");
  272. check(p0.total_matched == 25, "fast path total_matched excludes the sentinel");
  273. check(p0.has_more, "has_more true when more remain");
  274. page.offset = 20;
  275. auto p2 = store.scan("things", page);
  276. check(p2.documents.size() == 5, "final page returns the remainder");
  277. check(p2.total_matched == 25, "total_matched stable across pages");
  278. check(!p2.has_more, "has_more false on the last page");
  279. page.offset = 25;
  280. auto p3 = store.scan("things", page);
  281. check(p3.documents.empty(), "offset past the end returns nothing");
  282. check(p3.total_matched == 25, "and still reports the true total");
  283. // Paging must cover every document exactly once, in a stable order.
  284. std::set<std::string> seen;
  285. for (uint32_t off = 0; off < 25; off += 7) {
  286. smartbotic::database::Query q;
  287. q.limit = 7;
  288. q.offset = off;
  289. for (const auto& d : store.scan("things", q).documents) seen.insert(d.id);
  290. }
  291. check(seen.size() == 25, "paging the whole collection yields every document once");
  292. // A filter forces the general path; it must still be correct.
  293. smartbotic::database::Query fq;
  294. fq.limit = 100;
  295. smartbotic::database::Filter f;
  296. f.field = "kind";
  297. f.op = smartbotic::database::FilterOp::EQ;
  298. f.value = "odd";
  299. fq.filters.push_back(f);
  300. auto filtered = store.scan("things", fq);
  301. check(filtered.total_matched == 12, "filtered path still counts matches, not rows");
  302. // A sort also forces the general path.
  303. smartbotic::database::Query sq;
  304. sq.limit = 3;
  305. sq.sort = smartbotic::database::Sort{"n", true};
  306. auto sorted = store.scan("things", sq);
  307. check(sorted.documents.size() == 3, "sorted path paginates");
  308. check(sorted.total_matched == 25, "sorted path totals all rows");
  309. check(sorted.documents[0].data().value("n", -1) == 24,
  310. "descending sort really sorted (fast path must not swallow sorts)");
  311. // limit=0 keeps meaning "no documents, but a true total" - the contract
  312. // Count depends on (see test_scan_limit_zero_reports_total).
  313. smartbotic::database::Query zq;
  314. zq.limit = 0;
  315. auto z = store.scan("things", zq);
  316. check(z.documents.empty(), "limit=0 returns no documents on the fast path");
  317. check(z.total_matched == 25, "limit=0 still reports the true total");
  318. }
  319. // v2.8.0 — a WRITE that aborts must not poison the collection.
  320. //
  321. // This is the v2.4.3 EINVAL bug in a third failure mode, observed live in
  322. // production on 2.7.1: `find` on smartbotic-automation:workflows failed with
  323. // "LMDB cursor_open: Invalid argument" on every attempt while every other
  324. // collection was fine, and a restart was the only cure.
  325. //
  326. // Cause: open_for_write() cached the MDB_dbi immediately after mdb_dbi_open,
  327. // BEFORE the caller committed. LMDB keeps a handle private to the opening
  328. // transaction until it commits and CLOSES it if that transaction aborts - so any
  329. // write that threw after the handle was cached (a failed mdb_put, a sentinel
  330. // mismatch, a WriteTxn destructing uncommitted) left a closed handle in the
  331. // cache, and every later operation on that collection failed EINVAL for the rest
  332. // of the process's life.
  333. //
  334. // The abort is induced honestly here, with a key past LMDB's 511-byte limit, so
  335. // the test exercises the same path a real failed write takes.
  336. void test_aborted_write_does_not_poison_the_collection() {
  337. TmpEnv t("abortpoison");
  338. LmdbDocumentStore store(t.env);
  339. // Force a write that opens the sub-db and then fails: an oversized key makes
  340. // mdb_put return MDB_BAD_VALSIZE, which throws, so the WriteTxn aborts.
  341. const std::string huge_id(600, 'k');
  342. bool threw = false;
  343. try {
  344. store.put("poisoned", huge_id, make_doc(huge_id, "poisoned"));
  345. } catch (const std::exception&) {
  346. threw = true;
  347. }
  348. check(threw, "an oversized key really does fail the write");
  349. // The collection must still be usable. Before the fix, every one of these
  350. // failed with EINVAL because the cache held a handle LMDB had closed.
  351. bool ok_put = true;
  352. try {
  353. store.put("poisoned", "good", make_doc("good", "poisoned"));
  354. } catch (const std::exception&) {
  355. ok_put = false;
  356. }
  357. check(ok_put, "a later WRITE to the same collection still works");
  358. bool ok_read = true;
  359. try {
  360. smartbotic::database::Query q;
  361. q.limit = 10;
  362. auto res = store.scan("poisoned", q);
  363. check(res.documents.size() == 1, "and the document written after the abort is there");
  364. } catch (const std::exception&) {
  365. ok_read = false;
  366. }
  367. check(ok_read, "a later SCAN of the same collection still works (cursor_open)");
  368. bool ok_count = true;
  369. try {
  370. check(store.count("poisoned") == 1, "count is right after the abort");
  371. } catch (const std::exception&) {
  372. ok_count = false;
  373. }
  374. check(ok_count, "and count() does not throw");
  375. check(store.get("poisoned", "good").has_value(), "get() works after the abort");
  376. }
  377. // v2.8.0 — the two-pass filtered scan must agree with the old row-at-a-time path
  378. // on every operator, not just the common ones.
  379. //
  380. // scan() now evaluates predicates against a yyjson tree via a field resolver and
  381. // materialises only the returned page, because building a Document per row was
  382. // 88% of a filtered query's cost (3168ms vs 369ms for the parse alone over 193 MB
  383. // of real rows). A resolver that mishandles one operator returns silently wrong
  384. // data, so this walks the matrix.
  385. void test_filtered_scan_operator_matrix() {
  386. TmpEnv t("filtermatrix");
  387. LmdbDocumentStore store(t.env);
  388. auto put = [&](const std::string& id, const nlohmann::json& data) {
  389. Document d;
  390. d.id = id;
  391. d.collection = "m";
  392. d.version = 3;
  393. d.createdAt = 1000;
  394. d.updatedAt = 2000;
  395. d.set_data(data);
  396. store.put("m", id, d);
  397. };
  398. put("a", {{"n", 1}, {"kind", "odd"}, {"tags", {"x", "y"}}, {"nest", {{"deep", "hit"}}}});
  399. put("b", {{"n", 2}, {"kind", "even"}, {"tags", {"y"}}, {"nest", {{"deep", "miss"}}}});
  400. put("c", {{"n", 3}, {"kind", "odd"}, {"tags", nlohmann::json::array()}});
  401. put("d", {{"n", 4}, {"kind", "even"}, {"extra", "present"}});
  402. auto ids = [&](const smartbotic::database::Query& q) {
  403. std::vector<std::string> out;
  404. for (const auto& d : store.scan("m", q).documents) out.push_back(d.id);
  405. std::sort(out.begin(), out.end());
  406. return out;
  407. };
  408. auto q1 = [&](const char* field, smartbotic::database::FilterOp op,
  409. const nlohmann::json& val) {
  410. smartbotic::database::Query q;
  411. q.limit = 100;
  412. smartbotic::database::Filter f;
  413. f.field = field; f.op = op; f.value = val;
  414. q.filters.push_back(f);
  415. return q;
  416. };
  417. using Op = smartbotic::database::FilterOp;
  418. check(ids(q1("kind", Op::EQ, "odd")) == (std::vector<std::string>{"a", "c"}),
  419. "EQ on a data field");
  420. check(ids(q1("kind", Op::NE, "odd")) == (std::vector<std::string>{"b", "d"}),
  421. "NE on a data field");
  422. check(ids(q1("n", Op::GT, 2)) == (std::vector<std::string>{"c", "d"}), "GT numeric");
  423. check(ids(q1("n", Op::GTE, 3)) == (std::vector<std::string>{"c", "d"}), "GTE numeric");
  424. check(ids(q1("n", Op::LT, 2)) == (std::vector<std::string>{"a"}), "LT numeric");
  425. check(ids(q1("n", Op::LTE, 2)) == (std::vector<std::string>{"a", "b"}), "LTE numeric");
  426. check(ids(q1("n", Op::IN, nlohmann::json::array({1, 4}))) ==
  427. (std::vector<std::string>{"a", "d"}), "IN");
  428. check(ids(q1("tags", Op::CONTAINS, "x")) == (std::vector<std::string>{"a"}),
  429. "CONTAINS descends into an array value");
  430. check(ids(q1("extra", Op::EXISTS, true)) == (std::vector<std::string>{"d"}),
  431. "EXISTS true");
  432. check(ids(q1("extra", Op::EXISTS, false)) ==
  433. (std::vector<std::string>{"a", "b", "c"}), "EXISTS false");
  434. check(ids(q1("kind", Op::REGEX, "^od")) == (std::vector<std::string>{"a", "c"}),
  435. "REGEX");
  436. check(ids(q1("nest.deep", Op::EQ, "hit")) == (std::vector<std::string>{"a"}),
  437. "dotted path descends into data");
  438. check(ids(q1("nest.missing", Op::EXISTS, true)).empty(),
  439. "a dotted path that does not resolve matches nothing");
  440. // Document metadata, which lives at the top level of the stored JSON rather
  441. // than inside "data".
  442. check(ids(q1("_id", Op::EQ, "b")) == (std::vector<std::string>{"b"}), "_id");
  443. check(ids(q1("_version", Op::EQ, 3)).size() == 4, "_version");
  444. check(ids(q1("_created_at", Op::GTE, 1000)).size() == 4, "_created_at");
  445. check(ids(q1("_updated_at", Op::LT, 2000)).empty(), "_updated_at");
  446. // SEARCH must still work - it needs the whole document, so it takes the old
  447. // path.
  448. check(ids(q1("", Op::SEARCH, "present")) == (std::vector<std::string>{"d"}),
  449. "SEARCH still matches (routed to the whole-document path)");
  450. check(ids(q1("", Op::SEARCH, "nothinghere")).empty(), "SEARCH non-match");
  451. // Sorting, pagination and total_matched over a filtered set.
  452. {
  453. smartbotic::database::Query q;
  454. q.limit = 1;
  455. smartbotic::database::Filter f;
  456. f.field = "kind"; f.op = Op::EQ; f.value = "odd";
  457. q.filters.push_back(f);
  458. q.sort = smartbotic::database::Sort{"n", true}; // descending
  459. auto page0 = store.scan("m", q);
  460. check(page0.total_matched == 2, "total_matched counts matches, not rows");
  461. check(page0.documents.size() == 1, "limit honoured");
  462. check(page0.documents[0].id == "c", "descending sort picks the highest first");
  463. check(page0.has_more, "has_more true mid-set");
  464. q.offset = 1;
  465. auto page1 = store.scan("m", q);
  466. check(page1.documents.size() == 1 && page1.documents[0].id == "a",
  467. "second page continues the sort order");
  468. check(!page1.has_more, "has_more false on the last page");
  469. q.offset = 5;
  470. check(store.scan("m", q).documents.empty(), "offset past the end is empty");
  471. }
  472. // Ascending, and a sort field that is missing from some documents.
  473. {
  474. smartbotic::database::Query q;
  475. q.limit = 10;
  476. q.sort = smartbotic::database::Sort{"extra", false};
  477. auto res = store.scan("m", q);
  478. check(res.total_matched == 4, "no filter plus a sort still totals every row");
  479. check(res.documents.size() == 4, "and returns them all");
  480. // sort_documents returns `descending` when the LEFT value is missing, so
  481. // ascending puts documents that HAVE the field first and the ones missing
  482. // it last. The two-pass path copies that rule rather than inventing one.
  483. check(res.documents.front().id == "d",
  484. "ascending: the document that has the sort field comes first");
  485. std::vector<std::string> tail;
  486. for (size_t i = 1; i < res.documents.size(); ++i) tail.push_back(res.documents[i].id);
  487. check(tail == (std::vector<std::string>{"a", "b", "c"}),
  488. "and the ones missing it follow, tie-broken by id");
  489. }
  490. }
  491. // v2.8.1 — concurrent reads must not rebind a cached handle.
  492. //
  493. // lmdb.h: "This function [mdb_dbi_open] must not be called from multiple
  494. // concurrent transactions in the same process. A transaction that uses this
  495. // function must finish (either commit or abort) before any other transaction in
  496. // the process may use this function."
  497. //
  498. // The old read path violated that on every read of a collection this process had
  499. // not yet written, from every gRPC thread at once. MDB_dbi is an index into the
  500. // env's shared handle table, so churning it silently REBOUND handles that
  501. // committed writes had already cached. Live on 2.8.0-3 that produced 41 refusals
  502. // in 45 minutes, all "caller asked for 'image_hashes' but the handle addresses
  503. // 'nsfw_images'" - and because each refusal bumped mirror drift, every read in
  504. // the process fell back to MemoryStore permanently.
  505. //
  506. // The fix primes all handles in one committed write txn at construction, so only
  507. // write txns ever call mdb_dbi_open and LMDB serialises those itself.
  508. //
  509. // HONEST LIMITATION: this test does NOT reproduce the race. Reverting the fix
  510. // leaves it green apart from the primed_count assertion - the torn slot-table
  511. // update needs an interleaving of concurrent mdb_dbi_open calls that 8 threads
  512. // over 10 collections does not reliably hit. What the test does pin is the
  513. // invariant the race violates (no read returns another collection's document, no
  514. // write is refused by the sentinel) plus the mechanism that removes the race:
  515. // handles are opened up front, so the read path has no mdb_dbi_open left to call.
  516. // The deterministic evidence is the documented contract in lmdb.h and the
  517. // production log.
  518. void test_concurrent_reads_do_not_rebind_cached_handles() {
  519. const std::string path = make_tmpdir("dbi-concurrent");
  520. const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
  521. struct Cleanup {
  522. const std::string& p;
  523. ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
  524. } cleanup{path};
  525. // Enough collections that slot churn has somewhere to go.
  526. const std::vector<std::string> colls = {
  527. "alpha", "beta", "gamma", "delta", "epsilon",
  528. "zeta", "eta", "theta", "iota", "kappa"};
  529. {
  530. LmdbEnv env(opts);
  531. LmdbDocumentStore writer(env);
  532. for (const auto& c : colls) {
  533. Document d;
  534. d.id = "seed";
  535. d.collection = c;
  536. d.set_data(nlohmann::json{{"who", c}});
  537. writer.put(c, "seed", d);
  538. }
  539. }
  540. // Restart: fresh env, so the shared handle table starts empty and the store
  541. // must prime it.
  542. LmdbEnv env2(opts);
  543. LmdbDocumentStore store(env2);
  544. check(store.prime_error().empty(), "priming succeeded on reopen");
  545. check(store.primed_count() >= colls.size(),
  546. "priming opened a handle for every existing sub-db");
  547. // Hammer reads from many threads. Every one of these used to call
  548. // mdb_dbi_open inside its own read txn.
  549. std::atomic<int> read_failures{0};
  550. std::atomic<int> wrong_data{0};
  551. {
  552. std::vector<std::thread> threads;
  553. for (int t = 0; t < 8; ++t) {
  554. threads.emplace_back([&, t]() {
  555. for (int i = 0; i < 40; ++i) {
  556. const auto& c = colls[(t + i) % colls.size()];
  557. try {
  558. auto got = store.get(c, "seed");
  559. if (!got) { ++read_failures; continue; }
  560. // A rebound handle reads a STRANGER's sub-db, so the
  561. // document that comes back belongs to another collection.
  562. if (got->data().value("who", std::string{}) != c) ++wrong_data;
  563. } catch (const std::exception&) {
  564. ++read_failures;
  565. }
  566. }
  567. });
  568. }
  569. for (auto& th : threads) th.join();
  570. }
  571. check(read_failures.load() == 0, "concurrent reads all succeeded");
  572. check(wrong_data.load() == 0,
  573. "no read returned another collection's document - a rebound handle "
  574. "addresses whichever sub-db now occupies its slot");
  575. // Now write to every collection through the cached handles. This is where
  576. // the live failure surfaced: the sentinel refused the write.
  577. int write_failures = 0;
  578. for (const auto& c : colls) {
  579. try {
  580. Document d;
  581. d.id = "after";
  582. d.collection = c;
  583. d.set_data(nlohmann::json{{"who", c}});
  584. store.put(c, "after", d);
  585. } catch (const std::exception&) {
  586. ++write_failures;
  587. }
  588. }
  589. check(write_failures == 0,
  590. "writes through primed handles are not refused by the identity "
  591. "sentinel - the refusal is what production saw");
  592. for (const auto& c : colls) {
  593. check(store.count(c) == 2, ("both documents readable in " + c).c_str());
  594. }
  595. }
  596. // A collection that exists on disk but has NOT been written by this process must
  597. // still be readable. The read path no longer opens handles on demand, so if
  598. // priming missed anything a read would report the collection as EMPTY - a
  599. // silent-wrong-data failure worse than the bug being fixed.
  600. void test_existing_collection_readable_without_writing_first() {
  601. const std::string path = make_tmpdir("dbi-prime-read");
  602. const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
  603. struct Cleanup {
  604. const std::string& p;
  605. ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
  606. } cleanup{path};
  607. {
  608. LmdbEnv env(opts);
  609. LmdbDocumentStore writer(env);
  610. Document d;
  611. d.id = "only";
  612. d.collection = "archive";
  613. d.set_data(nlohmann::json{{"kept", true}});
  614. writer.put("archive", "only", d);
  615. }
  616. LmdbEnv env2(opts);
  617. LmdbDocumentStore reader(env2);
  618. // Read-only access, never a write on this collection in this process.
  619. auto got = reader.get("archive", "only");
  620. check(got.has_value(), "a never-written-here collection is still readable");
  621. check(reader.count("archive") == 1, "count sees it");
  622. smartbotic::database::Query q; q.limit = 10;
  623. check(reader.scan("archive", q).documents.size() == 1, "scan sees it");
  624. // And a collection that genuinely does not exist still reads as absent.
  625. check(!reader.get("nosuch", "x").has_value(),
  626. "a missing collection is still absent, not an error");
  627. check(reader.count("nosuch") == 0, "and counts zero");
  628. }
  629. // v2.9.0 — a secondary index must stay exactly in step with the documents.
  630. //
  631. // The index is a second copy of a fact already stored in the row. Every way the
  632. // two can diverge is a silent-wrong-data bug: a stale entry returns a row that
  633. // no longer matches, a missing entry hides a row that does. So this walks the
  634. // full lifecycle - insert, update the indexed field, update something else,
  635. // delete, re-insert - and after every step asserts the index agrees with a
  636. // brute-force scan of the collection.
  637. void test_index_tracks_documents_through_every_write() {
  638. TmpEnv t("idx-maint");
  639. LmdbDocumentStore store(t.env);
  640. store.set_indexed_fields("execs", {"workflowId"});
  641. auto put = [&](const std::string& id, const std::string& wf, int n) {
  642. Document d;
  643. d.id = id;
  644. d.collection = "execs";
  645. d.set_data(nlohmann::json{{"workflowId", wf}, {"n", n}});
  646. store.put("execs", id, d);
  647. };
  648. // What the index SHOULD say, computed by scanning every row - the oracle.
  649. auto truth = [&](const std::string& wf) {
  650. smartbotic::database::Query q;
  651. q.limit = 10000;
  652. std::vector<std::string> ids;
  653. for (const auto& d : store.scan("execs", q).documents) {
  654. if (d.data().value("workflowId", std::string{}) == wf) ids.push_back(d.id);
  655. }
  656. std::sort(ids.begin(), ids.end());
  657. return ids;
  658. };
  659. auto indexed = [&](const std::string& wf) {
  660. auto got = store.index_lookup_eq("execs", "workflowId", nlohmann::json(wf));
  661. std::vector<std::string> ids = got.value_or(std::vector<std::string>{});
  662. std::sort(ids.begin(), ids.end());
  663. return ids;
  664. };
  665. auto agree = [&](const std::string& wf, const char* stage) {
  666. const bool ok = indexed(wf) == truth(wf);
  667. const std::string msg = "index agrees with a full scan for " + wf +
  668. " after " + stage;
  669. check(ok, msg.c_str());
  670. };
  671. // Insert
  672. put("e1", "wf-a", 1);
  673. put("e2", "wf-a", 2);
  674. put("e3", "wf-b", 3);
  675. agree("wf-a", "inserts");
  676. agree("wf-b", "inserts");
  677. check(indexed("wf-a").size() == 2, "two rows under wf-a");
  678. // Update the INDEXED field: the old posting must go, the new one appear.
  679. put("e2", "wf-b", 2);
  680. agree("wf-a", "moving e2 to wf-b");
  681. agree("wf-b", "moving e2 to wf-b");
  682. check(indexed("wf-a").size() == 1, "wf-a lost e2");
  683. check(indexed("wf-b").size() == 2, "wf-b gained it");
  684. // Update an UNindexed field: the index must be untouched, not duplicated.
  685. put("e1", "wf-a", 99);
  686. agree("wf-a", "updating an unindexed field");
  687. check(indexed("wf-a").size() == 1,
  688. "no duplicate posting from re-writing the same indexed value");
  689. // Delete
  690. check(store.del("execs", "e3"), "deleted e3");
  691. agree("wf-b", "deleting e3");
  692. check(indexed("wf-b").size() == 1, "e3 is gone from the index");
  693. // Re-insert the same id
  694. put("e3", "wf-b", 7);
  695. agree("wf-b", "re-inserting e3");
  696. check(indexed("wf-b").size() == 2, "e3 is back exactly once");
  697. // A value with no rows is an empty result, NOT "no index".
  698. auto none = store.index_lookup_eq("execs", "workflowId", nlohmann::json("wf-zzz"));
  699. check(none.has_value() && none->empty(),
  700. "an indexed field with no matching rows returns empty, not nullopt - "
  701. "nullopt means 'no index' and would send the caller to a scan");
  702. // An undeclared field has no index, and must say so rather than say 'none'.
  703. auto unindexed = store.index_lookup_eq("execs", "n", nlohmann::json(1));
  704. check(!unindexed.has_value(),
  705. "an unindexed field returns nullopt so the caller falls back to a scan "
  706. "instead of concluding there are no matches");
  707. }
  708. // A collection with no declared index must behave exactly as before, and pay
  709. // nothing. Also: declaring an index later must pick up the rows already there.
  710. void test_build_index_over_existing_rows() {
  711. TmpEnv t("idx-build");
  712. LmdbDocumentStore store(t.env);
  713. // Write BEFORE declaring the index.
  714. for (int i = 0; i < 20; ++i) {
  715. Document d;
  716. d.id = "d" + std::to_string(i);
  717. d.collection = "c";
  718. d.set_data(nlohmann::json{{"grp", i % 4 == 0 ? "hot" : "cold"}});
  719. store.put("c", d.id, d);
  720. }
  721. check(!store.index_lookup_eq("c", "grp", nlohmann::json("hot")).has_value(),
  722. "no index exists before it is declared");
  723. const uint64_t built = store.build_index("c", "grp");
  724. check(built == 20, "the backfill indexed every existing row");
  725. store.set_indexed_fields("c", {"grp"});
  726. auto hot = store.index_lookup_eq("c", "grp", nlohmann::json("hot"));
  727. check(hot.has_value() && hot->size() == 5,
  728. "the backfilled index finds the pre-existing rows (d0,d4,d8,d12,d16)");
  729. // Idempotent: a second build must not double the postings.
  730. store.build_index("c", "grp");
  731. hot = store.index_lookup_eq("c", "grp", nlohmann::json("hot"));
  732. check(hot.has_value() && hot->size() == 5,
  733. "re-running the backfill does not duplicate postings");
  734. // The count guard must see the same number without reading the ids.
  735. auto n = store.index_count_eq("c", "grp", nlohmann::json("hot"));
  736. check(n.has_value() && *n == 5, "index_count_eq agrees with the lookup");
  737. auto cold = store.index_count_eq("c", "grp", nlohmann::json("cold"));
  738. check(cold.has_value() && *cold == 15, "and counts the larger group");
  739. check(store.drop_index("c", "grp"), "the index drops");
  740. check(!store.index_lookup_eq("c", "grp", nlohmann::json("hot")).has_value(),
  741. "after dropping, lookups report no index rather than no rows");
  742. }
  743. // Numbers are where an index most easily disagrees with a scan, because the scan
  744. // compares across numeric subtypes. A doc stored with 5 must be found by a
  745. // filter asking for 5.0 through EITHER path.
  746. void test_index_numeric_equality_matches_scan() {
  747. TmpEnv t("idx-num");
  748. LmdbDocumentStore store(t.env);
  749. store.set_indexed_fields("m", {"code"});
  750. Document a;
  751. a.id = "a"; a.collection = "m";
  752. a.set_data(nlohmann::json{{"code", 5}}); // integer
  753. store.put("m", "a", a);
  754. Document b;
  755. b.id = "b"; b.collection = "m";
  756. b.set_data(nlohmann::json{{"code", 5.0}}); // integral double
  757. store.put("m", "b", b);
  758. auto by_int = store.index_lookup_eq("m", "code", nlohmann::json(5));
  759. auto by_dbl = store.index_lookup_eq("m", "code", nlohmann::json(5.0));
  760. check(by_int.has_value() && by_int->size() == 2,
  761. "asking for 5 finds BOTH the int and the integral-double row");
  762. check(by_dbl == by_int,
  763. "and asking for 5.0 returns exactly the same rows - the scan's EQ "
  764. "compares numbers across subtypes, so the index must too");
  765. // A big integer must not collide with its neighbour via double precision.
  766. Document c;
  767. c.id = "c"; c.collection = "m";
  768. c.set_data(nlohmann::json{{"code", 1786263002080195076LL}});
  769. store.put("m", "c", c);
  770. auto near = store.index_lookup_eq("m", "code",
  771. nlohmann::json(1786263002080195077LL));
  772. check(near.has_value() && near->empty(),
  773. "a neighbouring ns-scale integer does not collide - these two ARE "
  774. "equal as doubles, so routing through double would false-match");
  775. }
  776. // v2.9.0 — an indexed plan must return EXACTLY what the unindexed plan returns.
  777. //
  778. // This is the whole safety argument for indexing. The index is an optimisation,
  779. // so any observable difference is a bug, and the interesting failures are silent:
  780. // a missing posting drops a row, a stale one adds a row that no longer matches,
  781. // and a different code path can disagree about ordering or total_matched.
  782. //
  783. // The test runs each query twice against the same data - once with the field
  784. // declared indexed, once not - and compares the complete result: ids in order,
  785. // total_matched, and has_more.
  786. void test_indexed_and_unindexed_plans_agree() {
  787. TmpEnv t("idx-equiv");
  788. LmdbDocumentStore store(t.env);
  789. // 300 rows: `grp` is selective enough to use the index (10 groups of 30 =
  790. // 10%), `bucket` deliberately is NOT (2 values, 50% each) so the guard has
  791. // something to decline.
  792. for (int i = 0; i < 300; ++i) {
  793. Document d;
  794. d.id = "r" + std::string(i < 10 ? "00" : (i < 100 ? "0" : "")) +
  795. std::to_string(i);
  796. d.collection = "c";
  797. d.set_data(nlohmann::json{
  798. {"grp", "g" + std::to_string(i % 10)},
  799. {"bucket", (i % 2 == 0) ? "even" : "odd"},
  800. {"n", i},
  801. {"nest", {{"deep", "d" + std::to_string(i % 10)}}},
  802. });
  803. store.put("c", d.id, d);
  804. }
  805. using Op = smartbotic::database::FilterOp;
  806. struct Case {
  807. const char* name;
  808. std::vector<smartbotic::database::Filter> filters;
  809. std::optional<smartbotic::database::Sort> sort;
  810. uint32_t limit;
  811. uint32_t offset;
  812. };
  813. auto F = [](const char* f, Op op, const nlohmann::json& v) {
  814. smartbotic::database::Filter x;
  815. x.field = f; x.op = op; x.value = v;
  816. return x;
  817. };
  818. const std::vector<Case> cases = {
  819. {"eq indexed field", {F("grp", Op::EQ, "g3")}, std::nullopt, 100, 0},
  820. {"eq + second predicate", {F("grp", Op::EQ, "g3"), F("bucket", Op::EQ, "even")}, std::nullopt, 100, 0},
  821. {"eq + range on another field", {F("grp", Op::EQ, "g3"), F("n", Op::GT, 100)}, std::nullopt, 100, 0},
  822. {"eq with sort asc", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", false}, 100, 0},
  823. {"eq with sort desc", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", true}, 100, 0},
  824. {"eq paginated", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", false}, 7, 10},
  825. {"eq offset past end", {F("grp", Op::EQ, "g3")}, std::nullopt, 10, 999},
  826. {"eq matching nothing", {F("grp", Op::EQ, "nope")}, std::nullopt, 100, 0},
  827. {"eq on unselective field", {F("bucket", Op::EQ, "even")}, std::nullopt, 100, 0},
  828. {"eq plus SEARCH", {F("grp", Op::EQ, "g3"), F("", Op::SEARCH, "g3")}, std::nullopt, 100, 0},
  829. {"ne on indexed field", {F("grp", Op::NE, "g3")}, std::nullopt, 100, 0},
  830. {"eq on nested path", {F("nest.deep", Op::EQ, "d4")}, std::nullopt, 100, 0},
  831. {"limit zero", {F("grp", Op::EQ, "g3")}, std::nullopt, 0, 0},
  832. };
  833. auto run = [&](const Case& c) {
  834. smartbotic::database::Query q;
  835. q.filters = c.filters;
  836. q.sort = c.sort;
  837. q.limit = c.limit;
  838. q.offset = c.offset;
  839. auto r = store.scan("c", q);
  840. std::string sig = "total=" + std::to_string(r.total_matched) +
  841. " more=" + std::to_string(r.has_more ? 1 : 0) + " [";
  842. for (const auto& d : r.documents) { sig += d.id; sig += ","; }
  843. sig += "]";
  844. return sig;
  845. };
  846. for (const auto& c : cases) {
  847. store.set_indexed_fields("c", {}); // no index
  848. const std::string without = run(c);
  849. store.set_indexed_fields("c", {"grp", "bucket", "nest.deep"});
  850. store.build_index("c", "grp");
  851. store.build_index("c", "bucket");
  852. store.build_index("c", "nest.deep");
  853. const std::string with = run(c);
  854. const std::string msg = std::string("indexed and unindexed plans agree: ")
  855. + c.name;
  856. if (with != without) {
  857. std::cerr << " without index: " << without << "\n"
  858. << " with index: " << with << "\n";
  859. }
  860. check(with == without, msg.c_str());
  861. }
  862. // The agreement above is only meaningful if the index plan was actually
  863. // TAKEN for the selective cases. Otherwise the planner declined every time
  864. // and the test compared the scan against itself.
  865. store.set_indexed_fields("c", {"grp", "bucket", "nest.deep"});
  866. {
  867. store.reset_index_plan_stats();
  868. smartbotic::database::Query q;
  869. q.limit = 100;
  870. q.filters.push_back(F("grp", Op::EQ, "g3"));
  871. auto r = store.scan("c", q);
  872. auto st = store.index_plan_stats();
  873. check(r.total_matched == 30, "the selective query matches 30 of 300 rows");
  874. check(st.indexed_scans == 1 && st.full_scans == 0,
  875. "a selective EQ on an indexed field TAKES the index plan - without "
  876. "this the equivalence cases above would prove nothing");
  877. }
  878. {
  879. store.reset_index_plan_stats();
  880. smartbotic::database::Query q;
  881. q.limit = 100;
  882. q.filters.push_back(F("bucket", Op::EQ, "even"));
  883. auto r = store.scan("c", q);
  884. auto st = store.index_plan_stats();
  885. check(r.total_matched == 150, "the unselective query matches half the rows");
  886. check(st.indexed_scans == 0 && st.declined_unselective == 1,
  887. "and the guard DECLINES its index - 150 of 300 rows would cost more "
  888. "through the index than a scan, so declaring an index must not be "
  889. "able to pessimise a query");
  890. }
  891. {
  892. // An index on a field the query does not filter on must not be consulted.
  893. store.reset_index_plan_stats();
  894. smartbotic::database::Query q;
  895. q.limit = 100;
  896. q.filters.push_back(F("n", Op::GT, 250));
  897. store.scan("c", q);
  898. auto st = store.index_plan_stats();
  899. check(st.full_scans == 1 && st.indexed_scans == 0,
  900. "a query whose predicates name no indexed field scans");
  901. }
  902. auto n = store.index_count_eq("c", "bucket", nlohmann::json("even"));
  903. check(n.has_value() && *n == 150,
  904. "the unselective index does exist and holds 150 of 300 rows");
  905. }
  906. } // namespace
  907. int main() {
  908. std::cout << "=== test_subdb_identity ===\n";
  909. test_sentinel_roundtrip();
  910. test_misbound_handle_is_refused();
  911. test_unstamped_subdb_is_permitted();
  912. test_identity_key_predicate();
  913. test_sentinel_invisible_through_store();
  914. test_vector_subdb_sentinel();
  915. test_scan_limit_zero_reports_total();
  916. test_scan_fast_path_matches_general_path();
  917. test_aborted_write_does_not_poison_the_collection();
  918. test_filtered_scan_operator_matrix();
  919. test_concurrent_reads_do_not_rebind_cached_handles();
  920. test_existing_collection_readable_without_writing_first();
  921. test_index_tracks_documents_through_every_write();
  922. test_build_index_over_existing_rows();
  923. test_index_numeric_equality_matches_scan();
  924. test_indexed_and_unindexed_plans_agree();
  925. std::cout << "passed: " << g_pass << ", failed: " << g_fail << "\n";
  926. return g_fail == 0 ? 0 : 1;
  927. }