test_subdb_identity.cpp 60 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442
  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 <functional>
  19. #include <thread>
  20. #include <set>
  21. #include <string>
  22. #include <unistd.h>
  23. #include <lmdb.h>
  24. #include <nlohmann/json.hpp>
  25. #include "document.hpp"
  26. #include "storage/document_store_lmdb.hpp"
  27. #include "storage/lmdb_env.hpp"
  28. #include "storage/lmdb_txn.hpp"
  29. #include "storage/subdb_identity.hpp"
  30. namespace fs = std::filesystem;
  31. using smartbotic::database::Document;
  32. using smartbotic::db::storage::is_identity_key;
  33. using smartbotic::db::storage::kSubdbIdentityKey;
  34. using smartbotic::db::storage::LmdbDocumentStore;
  35. using smartbotic::db::storage::LmdbEnv;
  36. using smartbotic::db::storage::LmdbEnvOpts;
  37. using smartbotic::db::storage::read_subdb_identity;
  38. using smartbotic::db::storage::ReadTxn;
  39. using smartbotic::db::storage::verify_subdb_identity;
  40. using smartbotic::db::storage::write_subdb_identity;
  41. using smartbotic::db::storage::WriteTxn;
  42. namespace {
  43. int g_pass = 0;
  44. int g_fail = 0;
  45. void check(bool cond, const char* msg) {
  46. if (cond) {
  47. ++g_pass;
  48. } else {
  49. ++g_fail;
  50. std::cerr << "FAIL: " << msg << "\n";
  51. }
  52. }
  53. std::string make_tmpdir(const char* tag) {
  54. static std::atomic<int> counter{0};
  55. std::string path = "/tmp/subdb-identity-test-" + std::to_string(::getpid()) +
  56. "-" + std::to_string(counter.fetch_add(1)) + "-" + tag;
  57. std::error_code ec;
  58. fs::remove_all(path, ec);
  59. return path;
  60. }
  61. struct TmpEnv {
  62. std::string path;
  63. LmdbEnv env;
  64. explicit TmpEnv(const char* tag)
  65. : path(make_tmpdir(tag)),
  66. env(LmdbEnvOpts{path, 64ULL << 20, 256, 126, false}) {}
  67. ~TmpEnv() {
  68. std::error_code ec;
  69. fs::remove_all(path, ec);
  70. }
  71. TmpEnv(const TmpEnv&) = delete;
  72. TmpEnv& operator=(const TmpEnv&) = delete;
  73. };
  74. // Open (creating) a named sub-db inside a write txn and return its handle.
  75. unsigned int open_subdb(WriteTxn& txn, const char* name) {
  76. MDB_dbi dbi = 0;
  77. int rc = mdb_dbi_open(txn.raw(), name, MDB_CREATE, &dbi);
  78. assert(rc == MDB_SUCCESS);
  79. (void)rc;
  80. return dbi;
  81. }
  82. Document make_doc(const std::string& id, const std::string& collection) {
  83. Document d;
  84. d.id = id;
  85. d.collection = collection;
  86. d.set_data(nlohmann::json{{"seenCount", 1}, {"who", collection}});
  87. return d;
  88. }
  89. // -------------------------------------------------------------------------
  90. void test_sentinel_roundtrip() {
  91. TmpEnv t("roundtrip");
  92. {
  93. WriteTxn w(t.env);
  94. unsigned int dbi = open_subdb(w, "image_hashes");
  95. write_subdb_identity(w, dbi, "image_hashes");
  96. w.commit();
  97. }
  98. {
  99. WriteTxn w(t.env);
  100. unsigned int dbi = open_subdb(w, "image_hashes");
  101. bool threw = false;
  102. try {
  103. verify_subdb_identity(w, dbi, "image_hashes");
  104. } catch (const std::exception&) {
  105. threw = true;
  106. }
  107. check(!threw, "matching sentinel must verify without throwing");
  108. w.commit();
  109. }
  110. {
  111. ReadTxn r(t.env);
  112. MDB_dbi dbi = 0;
  113. mdb_dbi_open(r.raw(), "image_hashes", 0, &dbi);
  114. check(read_subdb_identity(r, dbi) == "image_hashes",
  115. "read_subdb_identity returns the stamped name");
  116. }
  117. }
  118. // THE regression test. A stale cache entry is, in effect, a handle that
  119. // addresses someone else's sub-db. Hand the verifier exactly that.
  120. void test_misbound_handle_is_refused() {
  121. TmpEnv t("misbound");
  122. unsigned int executions_dbi = 0;
  123. {
  124. WriteTxn w(t.env);
  125. unsigned int ih = open_subdb(w, "image_hashes");
  126. write_subdb_identity(w, ih, "image_hashes");
  127. executions_dbi = open_subdb(w, "executions");
  128. write_subdb_identity(w, executions_dbi, "executions");
  129. w.commit();
  130. }
  131. WriteTxn w(t.env);
  132. // Re-open so the handle is valid in this txn, then deliberately verify it
  133. // under the WRONG name — the production misbinding, reproduced.
  134. unsigned int exec = open_subdb(w, "executions");
  135. bool threw = false;
  136. std::string msg;
  137. try {
  138. verify_subdb_identity(w, exec, "image_hashes");
  139. } catch (const std::exception& e) {
  140. threw = true;
  141. msg = e.what();
  142. }
  143. check(threw, "handle for 'executions' verified as 'image_hashes' must throw");
  144. check(msg.find("image_hashes") != std::string::npos &&
  145. msg.find("executions") != std::string::npos,
  146. "misbinding error names both the requested and actual sub-db");
  147. w.abort();
  148. }
  149. // Existing deployments have sub-dbs with no sentinel. Those must keep working.
  150. void test_unstamped_subdb_is_permitted() {
  151. TmpEnv t("unstamped");
  152. WriteTxn w(t.env);
  153. unsigned int dbi = open_subdb(w, "legacy");
  154. bool threw = false;
  155. try {
  156. verify_subdb_identity(w, dbi, "legacy");
  157. } catch (const std::exception&) {
  158. threw = true;
  159. }
  160. check(!threw, "sub-db without a sentinel must verify (absence is unknown, not wrong)");
  161. w.commit();
  162. }
  163. void test_identity_key_predicate() {
  164. check(is_identity_key(kSubdbIdentityKey), "sentinel key recognised");
  165. check(!is_identity_key("__subdb_identity__"),
  166. "same text without the leading NUL is NOT the sentinel");
  167. check(!is_identity_key("e63c1b90"), "a document id is not the sentinel");
  168. check(kSubdbIdentityKey[0] == '\0',
  169. "sentinel must start with NUL so it cannot collide with a doc id");
  170. }
  171. // The sentinel is an implementation detail: it must never surface through the
  172. // DocumentStore API as a document, nor inflate a count.
  173. void test_sentinel_invisible_through_store() {
  174. TmpEnv t("invisible");
  175. LmdbDocumentStore store(t.env);
  176. store.put("image_hashes", "aaa", make_doc("aaa", "image_hashes"));
  177. store.put("image_hashes", "bbb", make_doc("bbb", "image_hashes"));
  178. check(store.count("image_hashes") == 2,
  179. "count() must exclude the identity sentinel");
  180. smartbotic::database::Query q;
  181. q.limit = 100;
  182. auto res = store.scan("image_hashes", q);
  183. check(res.documents.size() == 2, "scan() must exclude the identity sentinel");
  184. check(res.total_matched == 2, "scan() total_matched must exclude the sentinel");
  185. for (const auto& d : res.documents) {
  186. check(d.id == "aaa" || d.id == "bbb",
  187. "scan() must not surface the sentinel as a document");
  188. }
  189. // And it really is on disk.
  190. ReadTxn r(t.env);
  191. MDB_dbi dbi = 0;
  192. int rc = mdb_dbi_open(r.raw(), "image_hashes", 0, &dbi);
  193. check(rc == MDB_SUCCESS, "sub-db exists");
  194. check(read_subdb_identity(r, dbi) == "image_hashes",
  195. "store.put() stamps the sentinel on first write");
  196. }
  197. // A vector sub-db gets stamped with its own (prefixed) name, and scan_vectors
  198. // must skip the sentinel rather than trying to read it as float32 bytes.
  199. void test_vector_subdb_sentinel() {
  200. TmpEnv t("vectors");
  201. LmdbDocumentStore store(t.env);
  202. store.put_vector("emb", "v1", {1.0f, 2.0f, 3.0f});
  203. store.put_vector("emb", "v2", {4.0f, 5.0f, 6.0f});
  204. int seen = 0;
  205. bool bad = false;
  206. store.scan_vectors("emb", [&](std::string_view id, const float*, size_t n) {
  207. ++seen;
  208. if (n != 3) bad = true;
  209. if (is_identity_key(id)) bad = true;
  210. });
  211. check(seen == 2, "scan_vectors must skip the sentinel");
  212. check(!bad, "scan_vectors must not decode the sentinel as float data");
  213. ReadTxn r(t.env);
  214. MDB_dbi dbi = 0;
  215. mdb_dbi_open(r.raw(), "_vectors_emb", 0, &dbi);
  216. check(read_subdb_identity(r, dbi) == "_vectors_emb",
  217. "vector sub-db is stamped with its prefixed name");
  218. }
  219. // v2.4.4 Count is LMDB-first. Unfiltered it uses count(); filtered it uses
  220. // scan() with limit=0 and reads total_matched. Pin that contract: total_matched
  221. // is computed BEFORE pagination, so limit=0 must still report the true total
  222. // while returning no documents.
  223. void test_scan_limit_zero_reports_total() {
  224. TmpEnv t("counting");
  225. LmdbDocumentStore store(t.env);
  226. for (int i = 0; i < 5; ++i) {
  227. Document d;
  228. d.id = "id" + std::to_string(i);
  229. d.collection = "things";
  230. d.set_data(nlohmann::json{{"kind", i < 3 ? "alpha" : "beta"}});
  231. store.put("things", d.id, d);
  232. }
  233. smartbotic::database::Query q;
  234. q.limit = 0;
  235. auto all = store.scan("things", q);
  236. check(all.total_matched == 5, "limit=0 reports the full total");
  237. check(all.documents.empty(), "limit=0 returns no documents");
  238. check(store.count("things") == 5, "unfiltered count matches");
  239. smartbotic::database::Query fq;
  240. fq.limit = 0;
  241. smartbotic::database::Filter f;
  242. f.field = "kind";
  243. f.op = smartbotic::database::FilterOp::EQ;
  244. f.value = "alpha";
  245. fq.filters.push_back(f);
  246. auto filtered = store.scan("things", fq);
  247. check(filtered.total_matched == 3,
  248. "filtered limit=0 reports the matching total, not the collection size");
  249. }
  250. // v2.7.1 — the unfiltered/unsorted fast path in scan() must agree with the
  251. // general path exactly. It exists because the general path decoded every
  252. // document in the collection to return `limit` of them, so cost tracked total
  253. // bytes rather than page size (382ms to return one 510-byte document from a
  254. // 414 MB collection on a live instance). Any divergence here is a paging bug.
  255. void test_scan_fast_path_matches_general_path() {
  256. TmpEnv t("fastpath");
  257. LmdbDocumentStore store(t.env);
  258. for (int i = 0; i < 25; ++i) {
  259. Document d;
  260. char buf[16];
  261. std::snprintf(buf, sizeof(buf), "id%02d", i);
  262. d.id = buf;
  263. d.collection = "things";
  264. d.set_data(nlohmann::json{{"n", i}, {"kind", i % 2 ? "odd" : "even"}});
  265. store.put("things", d.id, d);
  266. }
  267. // total_matched and has_more must match what a full count says.
  268. smartbotic::database::Query page;
  269. page.limit = 10;
  270. page.offset = 0;
  271. auto p0 = store.scan("things", page);
  272. check(p0.documents.size() == 10, "fast path returns exactly `limit` docs");
  273. check(p0.total_matched == 25, "fast path total_matched excludes the sentinel");
  274. check(p0.has_more, "has_more true when more remain");
  275. page.offset = 20;
  276. auto p2 = store.scan("things", page);
  277. check(p2.documents.size() == 5, "final page returns the remainder");
  278. check(p2.total_matched == 25, "total_matched stable across pages");
  279. check(!p2.has_more, "has_more false on the last page");
  280. page.offset = 25;
  281. auto p3 = store.scan("things", page);
  282. check(p3.documents.empty(), "offset past the end returns nothing");
  283. check(p3.total_matched == 25, "and still reports the true total");
  284. // Paging must cover every document exactly once, in a stable order.
  285. std::set<std::string> seen;
  286. for (uint32_t off = 0; off < 25; off += 7) {
  287. smartbotic::database::Query q;
  288. q.limit = 7;
  289. q.offset = off;
  290. for (const auto& d : store.scan("things", q).documents) seen.insert(d.id);
  291. }
  292. check(seen.size() == 25, "paging the whole collection yields every document once");
  293. // A filter forces the general path; it must still be correct.
  294. smartbotic::database::Query fq;
  295. fq.limit = 100;
  296. smartbotic::database::Filter f;
  297. f.field = "kind";
  298. f.op = smartbotic::database::FilterOp::EQ;
  299. f.value = "odd";
  300. fq.filters.push_back(f);
  301. auto filtered = store.scan("things", fq);
  302. check(filtered.total_matched == 12, "filtered path still counts matches, not rows");
  303. // A sort also forces the general path.
  304. smartbotic::database::Query sq;
  305. sq.limit = 3;
  306. sq.sort = smartbotic::database::Sort{"n", true};
  307. auto sorted = store.scan("things", sq);
  308. check(sorted.documents.size() == 3, "sorted path paginates");
  309. check(sorted.total_matched == 25, "sorted path totals all rows");
  310. check(sorted.documents[0].data().value("n", -1) == 24,
  311. "descending sort really sorted (fast path must not swallow sorts)");
  312. // limit=0 keeps meaning "no documents, but a true total" - the contract
  313. // Count depends on (see test_scan_limit_zero_reports_total).
  314. smartbotic::database::Query zq;
  315. zq.limit = 0;
  316. auto z = store.scan("things", zq);
  317. check(z.documents.empty(), "limit=0 returns no documents on the fast path");
  318. check(z.total_matched == 25, "limit=0 still reports the true total");
  319. }
  320. // v2.8.0 — a WRITE that aborts must not poison the collection.
  321. //
  322. // This is the v2.4.3 EINVAL bug in a third failure mode, observed live in
  323. // production on 2.7.1: `find` on smartbotic-automation:workflows failed with
  324. // "LMDB cursor_open: Invalid argument" on every attempt while every other
  325. // collection was fine, and a restart was the only cure.
  326. //
  327. // Cause: open_for_write() cached the MDB_dbi immediately after mdb_dbi_open,
  328. // BEFORE the caller committed. LMDB keeps a handle private to the opening
  329. // transaction until it commits and CLOSES it if that transaction aborts - so any
  330. // write that threw after the handle was cached (a failed mdb_put, a sentinel
  331. // mismatch, a WriteTxn destructing uncommitted) left a closed handle in the
  332. // cache, and every later operation on that collection failed EINVAL for the rest
  333. // of the process's life.
  334. //
  335. // The abort is induced honestly here, with a key past LMDB's 511-byte limit, so
  336. // the test exercises the same path a real failed write takes.
  337. void test_aborted_write_does_not_poison_the_collection() {
  338. TmpEnv t("abortpoison");
  339. LmdbDocumentStore store(t.env);
  340. // Force a write that opens the sub-db and then fails: an oversized key makes
  341. // mdb_put return MDB_BAD_VALSIZE, which throws, so the WriteTxn aborts.
  342. const std::string huge_id(600, 'k');
  343. bool threw = false;
  344. try {
  345. store.put("poisoned", huge_id, make_doc(huge_id, "poisoned"));
  346. } catch (const std::exception&) {
  347. threw = true;
  348. }
  349. check(threw, "an oversized key really does fail the write");
  350. // The collection must still be usable. Before the fix, every one of these
  351. // failed with EINVAL because the cache held a handle LMDB had closed.
  352. bool ok_put = true;
  353. try {
  354. store.put("poisoned", "good", make_doc("good", "poisoned"));
  355. } catch (const std::exception&) {
  356. ok_put = false;
  357. }
  358. check(ok_put, "a later WRITE to the same collection still works");
  359. bool ok_read = true;
  360. try {
  361. smartbotic::database::Query q;
  362. q.limit = 10;
  363. auto res = store.scan("poisoned", q);
  364. check(res.documents.size() == 1, "and the document written after the abort is there");
  365. } catch (const std::exception&) {
  366. ok_read = false;
  367. }
  368. check(ok_read, "a later SCAN of the same collection still works (cursor_open)");
  369. bool ok_count = true;
  370. try {
  371. check(store.count("poisoned") == 1, "count is right after the abort");
  372. } catch (const std::exception&) {
  373. ok_count = false;
  374. }
  375. check(ok_count, "and count() does not throw");
  376. check(store.get("poisoned", "good").has_value(), "get() works after the abort");
  377. }
  378. // v2.8.0 — the two-pass filtered scan must agree with the old row-at-a-time path
  379. // on every operator, not just the common ones.
  380. //
  381. // scan() now evaluates predicates against a yyjson tree via a field resolver and
  382. // materialises only the returned page, because building a Document per row was
  383. // 88% of a filtered query's cost (3168ms vs 369ms for the parse alone over 193 MB
  384. // of real rows). A resolver that mishandles one operator returns silently wrong
  385. // data, so this walks the matrix.
  386. void test_filtered_scan_operator_matrix() {
  387. TmpEnv t("filtermatrix");
  388. LmdbDocumentStore store(t.env);
  389. auto put = [&](const std::string& id, const nlohmann::json& data) {
  390. Document d;
  391. d.id = id;
  392. d.collection = "m";
  393. d.version = 3;
  394. d.createdAt = 1000;
  395. d.updatedAt = 2000;
  396. d.set_data(data);
  397. store.put("m", id, d);
  398. };
  399. put("a", {{"n", 1}, {"kind", "odd"}, {"tags", {"x", "y"}}, {"nest", {{"deep", "hit"}}}});
  400. put("b", {{"n", 2}, {"kind", "even"}, {"tags", {"y"}}, {"nest", {{"deep", "miss"}}}});
  401. put("c", {{"n", 3}, {"kind", "odd"}, {"tags", nlohmann::json::array()}});
  402. put("d", {{"n", 4}, {"kind", "even"}, {"extra", "present"}});
  403. auto ids = [&](const smartbotic::database::Query& q) {
  404. std::vector<std::string> out;
  405. for (const auto& d : store.scan("m", q).documents) out.push_back(d.id);
  406. std::sort(out.begin(), out.end());
  407. return out;
  408. };
  409. auto q1 = [&](const char* field, smartbotic::database::FilterOp op,
  410. const nlohmann::json& val) {
  411. smartbotic::database::Query q;
  412. q.limit = 100;
  413. smartbotic::database::Filter f;
  414. f.field = field; f.op = op; f.value = val;
  415. q.filters.push_back(f);
  416. return q;
  417. };
  418. using Op = smartbotic::database::FilterOp;
  419. check(ids(q1("kind", Op::EQ, "odd")) == (std::vector<std::string>{"a", "c"}),
  420. "EQ on a data field");
  421. check(ids(q1("kind", Op::NE, "odd")) == (std::vector<std::string>{"b", "d"}),
  422. "NE on a data field");
  423. check(ids(q1("n", Op::GT, 2)) == (std::vector<std::string>{"c", "d"}), "GT numeric");
  424. check(ids(q1("n", Op::GTE, 3)) == (std::vector<std::string>{"c", "d"}), "GTE numeric");
  425. check(ids(q1("n", Op::LT, 2)) == (std::vector<std::string>{"a"}), "LT numeric");
  426. check(ids(q1("n", Op::LTE, 2)) == (std::vector<std::string>{"a", "b"}), "LTE numeric");
  427. check(ids(q1("n", Op::IN, nlohmann::json::array({1, 4}))) ==
  428. (std::vector<std::string>{"a", "d"}), "IN");
  429. check(ids(q1("tags", Op::CONTAINS, "x")) == (std::vector<std::string>{"a"}),
  430. "CONTAINS descends into an array value");
  431. check(ids(q1("extra", Op::EXISTS, true)) == (std::vector<std::string>{"d"}),
  432. "EXISTS true");
  433. check(ids(q1("extra", Op::EXISTS, false)) ==
  434. (std::vector<std::string>{"a", "b", "c"}), "EXISTS false");
  435. check(ids(q1("kind", Op::REGEX, "^od")) == (std::vector<std::string>{"a", "c"}),
  436. "REGEX");
  437. check(ids(q1("nest.deep", Op::EQ, "hit")) == (std::vector<std::string>{"a"}),
  438. "dotted path descends into data");
  439. check(ids(q1("nest.missing", Op::EXISTS, true)).empty(),
  440. "a dotted path that does not resolve matches nothing");
  441. // Document metadata, which lives at the top level of the stored JSON rather
  442. // than inside "data".
  443. check(ids(q1("_id", Op::EQ, "b")) == (std::vector<std::string>{"b"}), "_id");
  444. check(ids(q1("_version", Op::EQ, 3)).size() == 4, "_version");
  445. check(ids(q1("_created_at", Op::GTE, 1000)).size() == 4, "_created_at");
  446. check(ids(q1("_updated_at", Op::LT, 2000)).empty(), "_updated_at");
  447. // SEARCH must still work - it needs the whole document, so it takes the old
  448. // path.
  449. check(ids(q1("", Op::SEARCH, "present")) == (std::vector<std::string>{"d"}),
  450. "SEARCH still matches (routed to the whole-document path)");
  451. check(ids(q1("", Op::SEARCH, "nothinghere")).empty(), "SEARCH non-match");
  452. // Sorting, pagination and total_matched over a filtered set.
  453. {
  454. smartbotic::database::Query q;
  455. q.limit = 1;
  456. smartbotic::database::Filter f;
  457. f.field = "kind"; f.op = Op::EQ; f.value = "odd";
  458. q.filters.push_back(f);
  459. q.sort = smartbotic::database::Sort{"n", true}; // descending
  460. auto page0 = store.scan("m", q);
  461. check(page0.total_matched == 2, "total_matched counts matches, not rows");
  462. check(page0.documents.size() == 1, "limit honoured");
  463. check(page0.documents[0].id == "c", "descending sort picks the highest first");
  464. check(page0.has_more, "has_more true mid-set");
  465. q.offset = 1;
  466. auto page1 = store.scan("m", q);
  467. check(page1.documents.size() == 1 && page1.documents[0].id == "a",
  468. "second page continues the sort order");
  469. check(!page1.has_more, "has_more false on the last page");
  470. q.offset = 5;
  471. check(store.scan("m", q).documents.empty(), "offset past the end is empty");
  472. }
  473. // Ascending, and a sort field that is missing from some documents.
  474. {
  475. smartbotic::database::Query q;
  476. q.limit = 10;
  477. q.sort = smartbotic::database::Sort{"extra", false};
  478. auto res = store.scan("m", q);
  479. check(res.total_matched == 4, "no filter plus a sort still totals every row");
  480. check(res.documents.size() == 4, "and returns them all");
  481. // sort_documents returns `descending` when the LEFT value is missing, so
  482. // ascending puts documents that HAVE the field first and the ones missing
  483. // it last. The two-pass path copies that rule rather than inventing one.
  484. check(res.documents.front().id == "d",
  485. "ascending: the document that has the sort field comes first");
  486. std::vector<std::string> tail;
  487. for (size_t i = 1; i < res.documents.size(); ++i) tail.push_back(res.documents[i].id);
  488. check(tail == (std::vector<std::string>{"a", "b", "c"}),
  489. "and the ones missing it follow, tie-broken by id");
  490. }
  491. }
  492. // v2.8.1 — concurrent reads must not rebind a cached handle.
  493. //
  494. // lmdb.h: "This function [mdb_dbi_open] must not be called from multiple
  495. // concurrent transactions in the same process. A transaction that uses this
  496. // function must finish (either commit or abort) before any other transaction in
  497. // the process may use this function."
  498. //
  499. // The old read path violated that on every read of a collection this process had
  500. // not yet written, from every gRPC thread at once. MDB_dbi is an index into the
  501. // env's shared handle table, so churning it silently REBOUND handles that
  502. // committed writes had already cached. Live on 2.8.0-3 that produced 41 refusals
  503. // in 45 minutes, all "caller asked for 'image_hashes' but the handle addresses
  504. // 'nsfw_images'" - and because each refusal bumped mirror drift, every read in
  505. // the process fell back to MemoryStore permanently.
  506. //
  507. // The fix primes all handles in one committed write txn at construction, so only
  508. // write txns ever call mdb_dbi_open and LMDB serialises those itself.
  509. //
  510. // HONEST LIMITATION: this test does NOT reproduce the race. Reverting the fix
  511. // leaves it green apart from the primed_count assertion - the torn slot-table
  512. // update needs an interleaving of concurrent mdb_dbi_open calls that 8 threads
  513. // over 10 collections does not reliably hit. What the test does pin is the
  514. // invariant the race violates (no read returns another collection's document, no
  515. // write is refused by the sentinel) plus the mechanism that removes the race:
  516. // handles are opened up front, so the read path has no mdb_dbi_open left to call.
  517. // The deterministic evidence is the documented contract in lmdb.h and the
  518. // production log.
  519. void test_concurrent_reads_do_not_rebind_cached_handles() {
  520. const std::string path = make_tmpdir("dbi-concurrent");
  521. const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
  522. struct Cleanup {
  523. const std::string& p;
  524. ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
  525. } cleanup{path};
  526. // Enough collections that slot churn has somewhere to go.
  527. const std::vector<std::string> colls = {
  528. "alpha", "beta", "gamma", "delta", "epsilon",
  529. "zeta", "eta", "theta", "iota", "kappa"};
  530. {
  531. LmdbEnv env(opts);
  532. LmdbDocumentStore writer(env);
  533. for (const auto& c : colls) {
  534. Document d;
  535. d.id = "seed";
  536. d.collection = c;
  537. d.set_data(nlohmann::json{{"who", c}});
  538. writer.put(c, "seed", d);
  539. }
  540. }
  541. // Restart: fresh env, so the shared handle table starts empty and the store
  542. // must prime it.
  543. LmdbEnv env2(opts);
  544. LmdbDocumentStore store(env2);
  545. check(store.prime_error().empty(), "priming succeeded on reopen");
  546. check(store.primed_count() >= colls.size(),
  547. "priming opened a handle for every existing sub-db");
  548. // Hammer reads from many threads. Every one of these used to call
  549. // mdb_dbi_open inside its own read txn.
  550. std::atomic<int> read_failures{0};
  551. std::atomic<int> wrong_data{0};
  552. {
  553. std::vector<std::thread> threads;
  554. for (int t = 0; t < 8; ++t) {
  555. threads.emplace_back([&, t]() {
  556. for (int i = 0; i < 40; ++i) {
  557. const auto& c = colls[(t + i) % colls.size()];
  558. try {
  559. auto got = store.get(c, "seed");
  560. if (!got) { ++read_failures; continue; }
  561. // A rebound handle reads a STRANGER's sub-db, so the
  562. // document that comes back belongs to another collection.
  563. if (got->data().value("who", std::string{}) != c) ++wrong_data;
  564. } catch (const std::exception&) {
  565. ++read_failures;
  566. }
  567. }
  568. });
  569. }
  570. for (auto& th : threads) th.join();
  571. }
  572. check(read_failures.load() == 0, "concurrent reads all succeeded");
  573. check(wrong_data.load() == 0,
  574. "no read returned another collection's document - a rebound handle "
  575. "addresses whichever sub-db now occupies its slot");
  576. // Now write to every collection through the cached handles. This is where
  577. // the live failure surfaced: the sentinel refused the write.
  578. int write_failures = 0;
  579. for (const auto& c : colls) {
  580. try {
  581. Document d;
  582. d.id = "after";
  583. d.collection = c;
  584. d.set_data(nlohmann::json{{"who", c}});
  585. store.put(c, "after", d);
  586. } catch (const std::exception&) {
  587. ++write_failures;
  588. }
  589. }
  590. check(write_failures == 0,
  591. "writes through primed handles are not refused by the identity "
  592. "sentinel - the refusal is what production saw");
  593. for (const auto& c : colls) {
  594. check(store.count(c) == 2, ("both documents readable in " + c).c_str());
  595. }
  596. }
  597. // A collection that exists on disk but has NOT been written by this process must
  598. // still be readable. The read path no longer opens handles on demand, so if
  599. // priming missed anything a read would report the collection as EMPTY - a
  600. // silent-wrong-data failure worse than the bug being fixed.
  601. void test_existing_collection_readable_without_writing_first() {
  602. const std::string path = make_tmpdir("dbi-prime-read");
  603. const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
  604. struct Cleanup {
  605. const std::string& p;
  606. ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
  607. } cleanup{path};
  608. {
  609. LmdbEnv env(opts);
  610. LmdbDocumentStore writer(env);
  611. Document d;
  612. d.id = "only";
  613. d.collection = "archive";
  614. d.set_data(nlohmann::json{{"kept", true}});
  615. writer.put("archive", "only", d);
  616. }
  617. LmdbEnv env2(opts);
  618. LmdbDocumentStore reader(env2);
  619. // Read-only access, never a write on this collection in this process.
  620. auto got = reader.get("archive", "only");
  621. check(got.has_value(), "a never-written-here collection is still readable");
  622. check(reader.count("archive") == 1, "count sees it");
  623. smartbotic::database::Query q; q.limit = 10;
  624. check(reader.scan("archive", q).documents.size() == 1, "scan sees it");
  625. // And a collection that genuinely does not exist still reads as absent.
  626. check(!reader.get("nosuch", "x").has_value(),
  627. "a missing collection is still absent, not an error");
  628. check(reader.count("nosuch") == 0, "and counts zero");
  629. }
  630. // v2.9.0 — a secondary index must stay exactly in step with the documents.
  631. //
  632. // The index is a second copy of a fact already stored in the row. Every way the
  633. // two can diverge is a silent-wrong-data bug: a stale entry returns a row that
  634. // no longer matches, a missing entry hides a row that does. So this walks the
  635. // full lifecycle - insert, update the indexed field, update something else,
  636. // delete, re-insert - and after every step asserts the index agrees with a
  637. // brute-force scan of the collection.
  638. void test_index_tracks_documents_through_every_write() {
  639. TmpEnv t("idx-maint");
  640. LmdbDocumentStore store(t.env);
  641. store.set_indexed_fields("execs", {"workflowId"});
  642. auto put = [&](const std::string& id, const std::string& wf, int n) {
  643. Document d;
  644. d.id = id;
  645. d.collection = "execs";
  646. d.set_data(nlohmann::json{{"workflowId", wf}, {"n", n}});
  647. store.put("execs", id, d);
  648. };
  649. // What the index SHOULD say, computed by scanning every row - the oracle.
  650. auto truth = [&](const std::string& wf) {
  651. smartbotic::database::Query q;
  652. q.limit = 10000;
  653. std::vector<std::string> ids;
  654. for (const auto& d : store.scan("execs", q).documents) {
  655. if (d.data().value("workflowId", std::string{}) == wf) ids.push_back(d.id);
  656. }
  657. std::sort(ids.begin(), ids.end());
  658. return ids;
  659. };
  660. auto indexed = [&](const std::string& wf) {
  661. auto got = store.index_lookup_eq("execs", "workflowId", nlohmann::json(wf));
  662. std::vector<std::string> ids = got.value_or(std::vector<std::string>{});
  663. std::sort(ids.begin(), ids.end());
  664. return ids;
  665. };
  666. auto agree = [&](const std::string& wf, const char* stage) {
  667. const bool ok = indexed(wf) == truth(wf);
  668. const std::string msg = "index agrees with a full scan for " + wf +
  669. " after " + stage;
  670. check(ok, msg.c_str());
  671. };
  672. // Insert
  673. put("e1", "wf-a", 1);
  674. put("e2", "wf-a", 2);
  675. put("e3", "wf-b", 3);
  676. agree("wf-a", "inserts");
  677. agree("wf-b", "inserts");
  678. check(indexed("wf-a").size() == 2, "two rows under wf-a");
  679. // Update the INDEXED field: the old posting must go, the new one appear.
  680. put("e2", "wf-b", 2);
  681. agree("wf-a", "moving e2 to wf-b");
  682. agree("wf-b", "moving e2 to wf-b");
  683. check(indexed("wf-a").size() == 1, "wf-a lost e2");
  684. check(indexed("wf-b").size() == 2, "wf-b gained it");
  685. // Update an UNindexed field: the index must be untouched, not duplicated.
  686. put("e1", "wf-a", 99);
  687. agree("wf-a", "updating an unindexed field");
  688. check(indexed("wf-a").size() == 1,
  689. "no duplicate posting from re-writing the same indexed value");
  690. // Delete
  691. check(store.del("execs", "e3"), "deleted e3");
  692. agree("wf-b", "deleting e3");
  693. check(indexed("wf-b").size() == 1, "e3 is gone from the index");
  694. // Re-insert the same id
  695. put("e3", "wf-b", 7);
  696. agree("wf-b", "re-inserting e3");
  697. check(indexed("wf-b").size() == 2, "e3 is back exactly once");
  698. // A value with no rows is an empty result, NOT "no index".
  699. auto none = store.index_lookup_eq("execs", "workflowId", nlohmann::json("wf-zzz"));
  700. check(none.has_value() && none->empty(),
  701. "an indexed field with no matching rows returns empty, not nullopt - "
  702. "nullopt means 'no index' and would send the caller to a scan");
  703. // An undeclared field has no index, and must say so rather than say 'none'.
  704. auto unindexed = store.index_lookup_eq("execs", "n", nlohmann::json(1));
  705. check(!unindexed.has_value(),
  706. "an unindexed field returns nullopt so the caller falls back to a scan "
  707. "instead of concluding there are no matches");
  708. }
  709. // A collection with no declared index must behave exactly as before, and pay
  710. // nothing. Also: declaring an index later must pick up the rows already there.
  711. void test_build_index_over_existing_rows() {
  712. TmpEnv t("idx-build");
  713. LmdbDocumentStore store(t.env);
  714. // Write BEFORE declaring the index.
  715. for (int i = 0; i < 20; ++i) {
  716. Document d;
  717. d.id = "d" + std::to_string(i);
  718. d.collection = "c";
  719. d.set_data(nlohmann::json{{"grp", i % 4 == 0 ? "hot" : "cold"}});
  720. store.put("c", d.id, d);
  721. }
  722. check(!store.index_lookup_eq("c", "grp", nlohmann::json("hot")).has_value(),
  723. "no index exists before it is declared");
  724. const uint64_t built = store.build_index("c", "grp");
  725. check(built == 20, "the backfill indexed every existing row");
  726. store.set_indexed_fields("c", {"grp"});
  727. auto hot = store.index_lookup_eq("c", "grp", nlohmann::json("hot"));
  728. check(hot.has_value() && hot->size() == 5,
  729. "the backfilled index finds the pre-existing rows (d0,d4,d8,d12,d16)");
  730. // Idempotent: a second build must not double the postings.
  731. store.build_index("c", "grp");
  732. hot = store.index_lookup_eq("c", "grp", nlohmann::json("hot"));
  733. check(hot.has_value() && hot->size() == 5,
  734. "re-running the backfill does not duplicate postings");
  735. // The count guard must see the same number without reading the ids.
  736. auto n = store.index_count_eq("c", "grp", nlohmann::json("hot"));
  737. check(n.has_value() && *n == 5, "index_count_eq agrees with the lookup");
  738. auto cold = store.index_count_eq("c", "grp", nlohmann::json("cold"));
  739. check(cold.has_value() && *cold == 15, "and counts the larger group");
  740. check(store.drop_index("c", "grp"), "the index drops");
  741. check(!store.index_lookup_eq("c", "grp", nlohmann::json("hot")).has_value(),
  742. "after dropping, lookups report no index rather than no rows");
  743. }
  744. // Numbers are where an index most easily disagrees with a scan, because the scan
  745. // compares across numeric subtypes. A doc stored with 5 must be found by a
  746. // filter asking for 5.0 through EITHER path.
  747. void test_index_numeric_equality_matches_scan() {
  748. TmpEnv t("idx-num");
  749. LmdbDocumentStore store(t.env);
  750. store.set_indexed_fields("m", {"code"});
  751. Document a;
  752. a.id = "a"; a.collection = "m";
  753. a.set_data(nlohmann::json{{"code", 5}}); // integer
  754. store.put("m", "a", a);
  755. Document b;
  756. b.id = "b"; b.collection = "m";
  757. b.set_data(nlohmann::json{{"code", 5.0}}); // integral double
  758. store.put("m", "b", b);
  759. auto by_int = store.index_lookup_eq("m", "code", nlohmann::json(5));
  760. auto by_dbl = store.index_lookup_eq("m", "code", nlohmann::json(5.0));
  761. check(by_int.has_value() && by_int->size() == 2,
  762. "asking for 5 finds BOTH the int and the integral-double row");
  763. check(by_dbl == by_int,
  764. "and asking for 5.0 returns exactly the same rows - the scan's EQ "
  765. "compares numbers across subtypes, so the index must too");
  766. // A big integer must not collide with its neighbour via double precision.
  767. Document c;
  768. c.id = "c"; c.collection = "m";
  769. c.set_data(nlohmann::json{{"code", 1786263002080195076LL}});
  770. store.put("m", "c", c);
  771. auto near = store.index_lookup_eq("m", "code",
  772. nlohmann::json(1786263002080195077LL));
  773. check(near.has_value() && near->empty(),
  774. "a neighbouring ns-scale integer does not collide - these two ARE "
  775. "equal as doubles, so routing through double would false-match");
  776. }
  777. // v2.9.0 — an indexed plan must return EXACTLY what the unindexed plan returns.
  778. //
  779. // This is the whole safety argument for indexing. The index is an optimisation,
  780. // so any observable difference is a bug, and the interesting failures are silent:
  781. // a missing posting drops a row, a stale one adds a row that no longer matches,
  782. // and a different code path can disagree about ordering or total_matched.
  783. //
  784. // The test runs each query twice against the same data - once with the field
  785. // declared indexed, once not - and compares the complete result: ids in order,
  786. // total_matched, and has_more.
  787. void test_indexed_and_unindexed_plans_agree() {
  788. TmpEnv t("idx-equiv");
  789. LmdbDocumentStore store(t.env);
  790. // 300 rows: `grp` is selective enough to use the index (10 groups of 30 =
  791. // 10%), `bucket` deliberately is NOT (2 values, 50% each) so the guard has
  792. // something to decline.
  793. for (int i = 0; i < 300; ++i) {
  794. Document d;
  795. d.id = "r" + std::string(i < 10 ? "00" : (i < 100 ? "0" : "")) +
  796. std::to_string(i);
  797. d.collection = "c";
  798. d.set_data(nlohmann::json{
  799. {"grp", "g" + std::to_string(i % 10)},
  800. {"bucket", (i % 2 == 0) ? "even" : "odd"},
  801. {"n", i},
  802. {"nest", {{"deep", "d" + std::to_string(i % 10)}}},
  803. });
  804. store.put("c", d.id, d);
  805. }
  806. using Op = smartbotic::database::FilterOp;
  807. struct Case {
  808. const char* name;
  809. std::vector<smartbotic::database::Filter> filters;
  810. std::optional<smartbotic::database::Sort> sort;
  811. uint32_t limit;
  812. uint32_t offset;
  813. };
  814. auto F = [](const char* f, Op op, const nlohmann::json& v) {
  815. smartbotic::database::Filter x;
  816. x.field = f; x.op = op; x.value = v;
  817. return x;
  818. };
  819. const std::vector<Case> cases = {
  820. {"eq indexed field", {F("grp", Op::EQ, "g3")}, std::nullopt, 100, 0},
  821. {"eq + second predicate", {F("grp", Op::EQ, "g3"), F("bucket", Op::EQ, "even")}, std::nullopt, 100, 0},
  822. {"eq + range on another field", {F("grp", Op::EQ, "g3"), F("n", Op::GT, 100)}, std::nullopt, 100, 0},
  823. {"eq with sort asc", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", false}, 100, 0},
  824. {"eq with sort desc", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", true}, 100, 0},
  825. {"eq paginated", {F("grp", Op::EQ, "g3")}, smartbotic::database::Sort{"n", false}, 7, 10},
  826. {"eq offset past end", {F("grp", Op::EQ, "g3")}, std::nullopt, 10, 999},
  827. {"eq matching nothing", {F("grp", Op::EQ, "nope")}, std::nullopt, 100, 0},
  828. {"eq on unselective field", {F("bucket", Op::EQ, "even")}, std::nullopt, 100, 0},
  829. {"eq plus SEARCH", {F("grp", Op::EQ, "g3"), F("", Op::SEARCH, "g3")}, std::nullopt, 100, 0},
  830. {"ne on indexed field", {F("grp", Op::NE, "g3")}, std::nullopt, 100, 0},
  831. {"eq on nested path", {F("nest.deep", Op::EQ, "d4")}, std::nullopt, 100, 0},
  832. {"limit zero", {F("grp", Op::EQ, "g3")}, std::nullopt, 0, 0},
  833. };
  834. auto run = [&](const Case& c) {
  835. smartbotic::database::Query q;
  836. q.filters = c.filters;
  837. q.sort = c.sort;
  838. q.limit = c.limit;
  839. q.offset = c.offset;
  840. auto r = store.scan("c", q);
  841. std::string sig = "total=" + std::to_string(r.total_matched) +
  842. " more=" + std::to_string(r.has_more ? 1 : 0) + " [";
  843. for (const auto& d : r.documents) { sig += d.id; sig += ","; }
  844. sig += "]";
  845. return sig;
  846. };
  847. for (const auto& c : cases) {
  848. store.set_indexed_fields("c", {}); // no index
  849. const std::string without = run(c);
  850. store.set_indexed_fields("c", {"grp", "bucket", "nest.deep"});
  851. store.build_index("c", "grp");
  852. store.build_index("c", "bucket");
  853. store.build_index("c", "nest.deep");
  854. const std::string with = run(c);
  855. const std::string msg = std::string("indexed and unindexed plans agree: ")
  856. + c.name;
  857. if (with != without) {
  858. std::cerr << " without index: " << without << "\n"
  859. << " with index: " << with << "\n";
  860. }
  861. check(with == without, msg.c_str());
  862. }
  863. // The agreement above is only meaningful if the index plan was actually
  864. // TAKEN for the selective cases. Otherwise the planner declined every time
  865. // and the test compared the scan against itself.
  866. store.set_indexed_fields("c", {"grp", "bucket", "nest.deep"});
  867. {
  868. store.reset_index_plan_stats();
  869. smartbotic::database::Query q;
  870. q.limit = 100;
  871. q.filters.push_back(F("grp", Op::EQ, "g3"));
  872. auto r = store.scan("c", q);
  873. auto st = store.index_plan_stats();
  874. check(r.total_matched == 30, "the selective query matches 30 of 300 rows");
  875. check(st.indexed_scans == 1 && st.full_scans == 0,
  876. "a selective EQ on an indexed field TAKES the index plan - without "
  877. "this the equivalence cases above would prove nothing");
  878. }
  879. {
  880. store.reset_index_plan_stats();
  881. smartbotic::database::Query q;
  882. q.limit = 100;
  883. q.filters.push_back(F("bucket", Op::EQ, "even"));
  884. auto r = store.scan("c", q);
  885. auto st = store.index_plan_stats();
  886. check(r.total_matched == 150, "the unselective query matches half the rows");
  887. check(st.indexed_scans == 0 && st.declined_unselective == 1,
  888. "and the guard DECLINES its index - 150 of 300 rows would cost more "
  889. "through the index than a scan, so declaring an index must not be "
  890. "able to pessimise a query");
  891. }
  892. {
  893. // An index on a field the query does not filter on must not be consulted.
  894. store.reset_index_plan_stats();
  895. smartbotic::database::Query q;
  896. q.limit = 100;
  897. q.filters.push_back(F("n", Op::GT, 250));
  898. store.scan("c", q);
  899. auto st = store.index_plan_stats();
  900. check(st.full_scans == 1 && st.indexed_scans == 0,
  901. "a query whose predicates name no indexed field scans");
  902. }
  903. auto n = store.index_count_eq("c", "bucket", nlohmann::json("even"));
  904. check(n.has_value() && *n == 150,
  905. "the unselective index does exist and holds 150 of 300 rows");
  906. }
  907. // An empty document id is a zero-length LMDB key, which mdb_get rejects with
  908. // MDB_BAD_VALSIZE. That surfaced in production as a gRPC INTERNAL and an ERROR
  909. // log line every time a consumer asked for one - noise that masks real failures.
  910. // A key that cannot exist is absent, not an error.
  911. void test_empty_id_reads_as_absent() {
  912. TmpEnv t("empty-id");
  913. LmdbDocumentStore store(t.env);
  914. Document d;
  915. d.id = "real";
  916. d.collection = "c";
  917. d.set_data(nlohmann::json{{"x", 1}});
  918. store.put("c", "real", d);
  919. bool threw = false;
  920. try {
  921. check(!store.get("c", "").has_value(), "an empty id reads as absent");
  922. } catch (const std::exception&) {
  923. threw = true;
  924. }
  925. check(!threw, "and does NOT throw - MDB_BAD_VALSIZE became a gRPC INTERNAL");
  926. threw = false;
  927. try {
  928. check(!store.del("c", ""), "deleting an empty id is a no-op");
  929. } catch (const std::exception&) {
  930. threw = true;
  931. }
  932. check(!threw, "and does not throw either");
  933. check(store.get("c", "real").has_value(), "real ids still work");
  934. check(store.count("c") == 1, "and nothing was disturbed");
  935. }
  936. // v2.9.1 — ranges, CONTAINS and intersection must agree with the scan too, and
  937. // must actually be USED. Each is a new way for the index to disagree with a
  938. // brute-force answer, and each disagreement would be silent.
  939. void test_range_contains_and_intersection() {
  940. TmpEnv t("idx-v291");
  941. LmdbDocumentStore store(t.env);
  942. // n mixes integers and reals on purpose: v2.9.0 encoded those under separate
  943. // type tags, so a range spanning both could not be served at all.
  944. for (int i = 0; i < 400; ++i) {
  945. Document d;
  946. d.id = "r" + std::string(i < 10 ? "00" : (i < 100 ? "0" : "")) + std::to_string(i);
  947. d.collection = "c";
  948. nlohmann::json data{
  949. {"n", (i % 2 == 0) ? nlohmann::json(i) : nlohmann::json(i + 0.5)},
  950. {"grp", "g" + std::to_string(i % 20)},
  951. {"other", "o" + std::to_string(i % 20)},
  952. // Deliberately COARSE: 4 and 5 values, so each matches 25% and 20% of
  953. // 400 rows - both above the 10% budget alone - while their pair
  954. // narrows to i%20, i.e. 20 rows (5%). That is the only shape where
  955. // intersecting two posting lists earns its keep.
  956. {"q4", i % 4},
  957. {"q5", i % 5},
  958. {"tags", nlohmann::json::array({"t" + std::to_string(i % 25), "all"})},
  959. };
  960. d.set_data(data);
  961. store.put("c", d.id, d);
  962. }
  963. using Op = smartbotic::database::FilterOp;
  964. auto F = [](const char* f, Op op, const nlohmann::json& v) {
  965. smartbotic::database::Filter x;
  966. x.field = f; x.op = op; x.value = v;
  967. return x;
  968. };
  969. auto sig = [&](const std::vector<smartbotic::database::Filter>& fs,
  970. std::optional<smartbotic::database::Sort> so = std::nullopt) {
  971. smartbotic::database::Query q;
  972. q.filters = fs;
  973. q.sort = so;
  974. q.limit = 1000;
  975. auto r = store.scan("c", q);
  976. std::string out = "total=" + std::to_string(r.total_matched) + " [";
  977. std::vector<std::string> ids;
  978. for (const auto& d : r.documents) ids.push_back(d.id);
  979. std::sort(ids.begin(), ids.end());
  980. for (const auto& i : ids) { out += i; out += ","; }
  981. return out + "]";
  982. };
  983. const std::vector<std::pair<const char*, std::vector<smartbotic::database::Filter>>> cases = {
  984. {"GT on a mixed int/real field", {F("n", Op::GT, 380)}},
  985. {"GTE on a mixed field", {F("n", Op::GTE, 380)}},
  986. {"LT on a mixed field", {F("n", Op::LT, 12)}},
  987. {"LTE on a mixed field", {F("n", Op::LTE, 12)}},
  988. {"GT with a real bound", {F("n", Op::GT, 380.5)}},
  989. {"range that matches nothing", {F("n", Op::GT, 100000)}},
  990. {"range that matches everything", {F("n", Op::GT, -1)}},
  991. {"CONTAINS an array element", {F("tags", Op::CONTAINS, "t3")}},
  992. {"CONTAINS a common element", {F("tags", Op::CONTAINS, "all")}},
  993. {"CONTAINS a missing element", {F("tags", Op::CONTAINS, "nope")}},
  994. {"two broad EQs, narrow together", {F("q4", Op::EQ, 1), F("q5", Op::EQ, 2)}},
  995. {"two broad EQs, disjoint result", {F("q4", Op::EQ, 1), F("q5", Op::EQ, 2), F("grp", Op::EQ, "g19")}},
  996. {"range plus EQ", {F("n", Op::GT, 300), F("grp", Op::EQ, "g3")}},
  997. {"EQ on an array field (whole)", {F("tags", Op::EQ, nlohmann::json::array({"t3", "all"}))}},
  998. };
  999. for (const auto& [name, filters] : cases) {
  1000. store.set_indexed_fields("c", {});
  1001. const std::string without = sig(filters);
  1002. store.set_indexed_fields("c", {"n", "grp", "other", "tags", "q4", "q5"});
  1003. for (const char* f : {"n", "grp", "other", "tags", "q4", "q5"}) {
  1004. store.build_index("c", f);
  1005. }
  1006. const std::string with = sig(filters);
  1007. if (with != without) {
  1008. std::cerr << " without: " << without.substr(0, 200) << "\n"
  1009. << " with: " << with.substr(0, 200) << "\n";
  1010. }
  1011. const std::string msg = std::string("indexed == unindexed: ") + name;
  1012. check(with == without, msg.c_str());
  1013. }
  1014. // Now prove each new plan is actually taken.
  1015. store.set_indexed_fields("c", {"n", "grp", "other", "tags", "q4", "q5"});
  1016. {
  1017. store.reset_index_plan_stats();
  1018. (void)sig({F("n", Op::GT, 380)});
  1019. auto st = store.index_plan_stats();
  1020. check(st.range_scans == 1 && st.indexed_scans == 1,
  1021. "a selective range WALKS the index - impossible before the numeric "
  1022. "encoding was unified");
  1023. }
  1024. {
  1025. store.reset_index_plan_stats();
  1026. (void)sig({F("n", Op::GT, -1)}); // matches everything
  1027. auto st = store.index_plan_stats();
  1028. check(st.range_scans == 0 && st.full_scans == 1,
  1029. "an unselective range gives up and scans - and must return no-plan "
  1030. "rather than a truncated candidate list");
  1031. }
  1032. {
  1033. store.reset_index_plan_stats();
  1034. (void)sig({F("tags", Op::CONTAINS, "t3")});
  1035. auto st = store.index_plan_stats();
  1036. check(st.indexed_scans == 1,
  1037. "CONTAINS is served from the per-element postings an array writes");
  1038. }
  1039. {
  1040. store.reset_index_plan_stats();
  1041. (void)sig({F("q4", Op::EQ, 1), F("q5", Op::EQ, 2)});
  1042. auto st = store.index_plan_stats();
  1043. check(st.intersected_scans == 1 && st.indexed_scans == 1,
  1044. "q4 matches 25% and q5 20% - each too broad alone - so their posting "
  1045. "lists are INTERSECTED down to 5% instead of scanning");
  1046. }
  1047. // An array-valued field: EQ on a scalar may over-return from the index, which
  1048. // is safe only because every candidate is re-filtered. Pin that.
  1049. {
  1050. const std::string indexed = sig({F("tags", Op::EQ, "t3")});
  1051. store.set_indexed_fields("c", {});
  1052. const std::string scanned = sig({F("tags", Op::EQ, "t3")});
  1053. store.set_indexed_fields("c", {"n", "grp", "other", "tags", "q4", "q5"});
  1054. check(indexed == scanned,
  1055. "EQ for a scalar on an array field agrees - the index may name rows "
  1056. "whose array merely CONTAINS it, and re-filtering drops them");
  1057. }
  1058. }
  1059. // v2.9.2 — the index supplies the ORDER, not just filter candidates.
  1060. //
  1061. // "newest N, unfiltered" walked the whole collection to discover what to sort by:
  1062. // 341ms to return one row from 592 MB of real data, with the index declared and
  1063. // unused, because a Sort is not a filter. Walking the sort field's index reads the
  1064. // page and nothing else.
  1065. //
  1066. // The dangerous part is not speed, it is that sort_documents places rows MISSING
  1067. // the sort field FIRST when descending, while an index holds no posting for them.
  1068. // A walk would silently omit them from the first page - wrong rows, not slow ones.
  1069. // So every case here compares the ordered plan against the plain scan.
  1070. void test_index_supplies_ordering() {
  1071. using Op = smartbotic::database::FilterOp;
  1072. auto make = [](LmdbDocumentStore& store, int n,
  1073. const std::function<nlohmann::json(int)>& body) {
  1074. for (int i = 0; i < n; ++i) {
  1075. Document d;
  1076. d.id = "r" + std::string(i < 10 ? "00" : (i < 100 ? "0" : "")) +
  1077. std::to_string(i);
  1078. d.collection = "c";
  1079. d.set_data(body(i));
  1080. store.put("c", d.id, d);
  1081. }
  1082. };
  1083. auto sig = [](LmdbDocumentStore& store, const smartbotic::database::Query& q) {
  1084. auto r = store.scan("c", q);
  1085. std::string out = "total=" + std::to_string(r.total_matched) +
  1086. " more=" + std::to_string(r.has_more ? 1 : 0) + " [";
  1087. for (const auto& d : r.documents) { out += d.id; out += ","; }
  1088. return out + "]";
  1089. };
  1090. auto Q = [](const char* field, bool desc, uint32_t limit, uint32_t offset) {
  1091. smartbotic::database::Query q;
  1092. q.sort = smartbotic::database::Sort{field, desc};
  1093. q.limit = limit;
  1094. q.offset = offset;
  1095. return q;
  1096. };
  1097. // ---- every row carries the sort field: the ordered plan applies ----
  1098. {
  1099. TmpEnv t("ord-total");
  1100. LmdbDocumentStore store(t.env);
  1101. // Deliberate duplicate keys (i/3) so tie-breaking by id is exercised in
  1102. // both directions.
  1103. make(store, 300, [](int i) {
  1104. return nlohmann::json{{"seq", i}, {"dup", i / 3}};
  1105. });
  1106. const std::vector<std::pair<const char*, smartbotic::database::Query>> cases = {
  1107. {"desc, first page", Q("seq", true, 10, 0)},
  1108. {"asc, first page", Q("seq", false, 10, 0)},
  1109. {"desc, deep offset", Q("seq", true, 10, 250)},
  1110. {"asc, deep offset", Q("seq", false, 7, 33)},
  1111. {"desc, last partial page", Q("seq", true, 10, 295)},
  1112. {"offset past the end", Q("seq", true, 10, 999)},
  1113. {"limit=0", Q("seq", true, 0, 0)},
  1114. {"whole collection", Q("seq", false, 1000, 0)},
  1115. {"desc with tied keys", Q("dup", true, 12, 0)},
  1116. {"asc with tied keys", Q("dup", false, 12, 0)},
  1117. {"tied keys, deep offset", Q("dup", true, 5, 40)},
  1118. };
  1119. uint64_t ordered_used = 0;
  1120. for (const auto& [name, q] : cases) {
  1121. store.set_indexed_fields("c", {});
  1122. const std::string without = sig(store, q);
  1123. store.set_indexed_fields("c", {"seq", "dup"});
  1124. store.build_index("c", "seq");
  1125. store.build_index("c", "dup");
  1126. store.reset_index_plan_stats();
  1127. const std::string with = sig(store, q);
  1128. ordered_used += store.index_plan_stats().ordered_scans;
  1129. if (with != without) {
  1130. std::cerr << " scan: " << without.substr(0, 180) << "\n"
  1131. << " ordered: " << with.substr(0, 180) << "\n";
  1132. }
  1133. const std::string msg = std::string("ordered plan == scan: ") + name;
  1134. check(with == without, msg.c_str());
  1135. }
  1136. // EVERY case above must have used the ordered plan, not just one. The
  1137. // first version of this test asserted only a single query and so passed
  1138. // while the plan was silently never taken (the collection's row count
  1139. // included the identity sentinel, so entries never equalled rows).
  1140. check(ordered_used == cases.size(),
  1141. ("the index SUPPLIED THE ORDER in all " + std::to_string(cases.size()) +
  1142. " cases (got " + std::to_string(ordered_used) + ") - otherwise the "
  1143. "comparisons above are scan against scan").c_str());
  1144. }
  1145. // ---- THE trap: one row lacks the sort field ----
  1146. {
  1147. TmpEnv t("ord-missing");
  1148. LmdbDocumentStore store(t.env);
  1149. make(store, 50, [](int i) {
  1150. nlohmann::json j{{"other", i}};
  1151. if (i != 7) j["seq"] = i; // r007 has no seq
  1152. return j;
  1153. });
  1154. store.set_indexed_fields("c", {});
  1155. const std::string without = sig(store, Q("seq", true, 5, 0));
  1156. store.set_indexed_fields("c", {"seq"});
  1157. store.build_index("c", "seq");
  1158. const std::string with = sig(store, Q("seq", true, 5, 0));
  1159. check(with == without,
  1160. "one row missing the sort field: DESCENDING puts it FIRST, and the "
  1161. "index has no posting for it - the ordered plan must decline");
  1162. store.reset_index_plan_stats();
  1163. (void)sig(store, Q("seq", true, 5, 0));
  1164. check(store.index_plan_stats().ordered_scans == 0,
  1165. "and it does decline - entries != rows is the guard");
  1166. }
  1167. // ---- an array-valued sort field must also decline ----
  1168. {
  1169. TmpEnv t("ord-array");
  1170. LmdbDocumentStore store(t.env);
  1171. make(store, 40, [](int i) {
  1172. return nlohmann::json{{"seq", nlohmann::json::array({i})}};
  1173. });
  1174. store.set_indexed_fields("c", {});
  1175. const std::string without = sig(store, Q("seq", true, 5, 0));
  1176. store.set_indexed_fields("c", {"seq"});
  1177. store.build_index("c", "seq");
  1178. const std::string with = sig(store, Q("seq", true, 5, 0));
  1179. check(with == without,
  1180. "a one-element-array sort field agrees - it satisfies entries==rows, "
  1181. "so the fetched-page array check is what catches it");
  1182. }
  1183. }
  1184. // v2.9.2 — IN as a union of posting lists, EXISTS as all postings.
  1185. void test_index_in_and_exists() {
  1186. using Op = smartbotic::database::FilterOp;
  1187. TmpEnv t("in-exists");
  1188. LmdbDocumentStore store(t.env);
  1189. for (int i = 0; i < 400; ++i) {
  1190. Document d;
  1191. d.id = "r" + std::to_string(1000 + i);
  1192. d.collection = "c";
  1193. nlohmann::json j{{"grp", "g" + std::to_string(i % 40)}};
  1194. // `rare` exists on 8 of 400 rows, so EXISTS=true is highly selective.
  1195. if (i % 50 == 0) j["rare"] = i;
  1196. d.set_data(j);
  1197. store.put("c", d.id, d);
  1198. }
  1199. auto F = [](const char* f, Op op, const nlohmann::json& v) {
  1200. smartbotic::database::Filter x;
  1201. x.field = f; x.op = op; x.value = v;
  1202. return x;
  1203. };
  1204. auto sig = [&](const std::vector<smartbotic::database::Filter>& fs) {
  1205. smartbotic::database::Query q;
  1206. q.filters = fs;
  1207. q.limit = 1000;
  1208. auto r = store.scan("c", q);
  1209. std::vector<std::string> ids;
  1210. for (const auto& d : r.documents) ids.push_back(d.id);
  1211. std::sort(ids.begin(), ids.end());
  1212. std::string out = "total=" + std::to_string(r.total_matched) + " [";
  1213. for (const auto& i : ids) { out += i; out += ","; }
  1214. return out + "]";
  1215. };
  1216. const std::vector<std::pair<const char*, std::vector<smartbotic::database::Filter>>> cases = {
  1217. {"IN over three values", {F("grp", Op::IN, nlohmann::json::array({"g1","g2","g3"}))}},
  1218. {"IN with a missing value", {F("grp", Op::IN, nlohmann::json::array({"g1","nope"}))}},
  1219. {"IN over one value", {F("grp", Op::IN, nlohmann::json::array({"g5"}))}},
  1220. {"IN matching nothing", {F("grp", Op::IN, nlohmann::json::array({"x","y"}))}},
  1221. {"IN too broad to help", {F("grp", Op::IN, nlohmann::json::array(
  1222. {"g0","g1","g2","g3","g4","g5","g6","g7","g8","g9","g10","g11"}))}},
  1223. {"EXISTS true on a sparse field", {F("rare", Op::EXISTS, true)}},
  1224. {"EXISTS false on a sparse field", {F("rare", Op::EXISTS, false)}},
  1225. {"EXISTS true on a dense field", {F("grp", Op::EXISTS, true)}},
  1226. };
  1227. for (const auto& [name, filters] : cases) {
  1228. store.set_indexed_fields("c", {});
  1229. const std::string without = sig(filters);
  1230. store.set_indexed_fields("c", {"grp", "rare"});
  1231. store.build_index("c", "grp");
  1232. store.build_index("c", "rare");
  1233. const std::string with = sig(filters);
  1234. if (with != without) {
  1235. std::cerr << " scan: " << without.substr(0, 160) << "\n"
  1236. << " indexed: " << with.substr(0, 160) << "\n";
  1237. }
  1238. const std::string msg = std::string("indexed == scan: ") + name;
  1239. check(with == without, msg.c_str());
  1240. }
  1241. store.set_indexed_fields("c", {"grp", "rare"});
  1242. {
  1243. store.reset_index_plan_stats();
  1244. (void)sig({F("grp", Op::IN, nlohmann::json::array({"g1","g2","g3"}))});
  1245. check(store.index_plan_stats().union_scans == 1,
  1246. "IN over a narrow set UNIONS posting lists - 30 of 400 rows");
  1247. }
  1248. {
  1249. store.reset_index_plan_stats();
  1250. (void)sig({F("grp", Op::IN, nlohmann::json::array(
  1251. {"g0","g1","g2","g3","g4","g5","g6","g7","g8","g9","g10","g11"}))});
  1252. auto st = store.index_plan_stats();
  1253. check(st.union_scans == 0 && st.full_scans == 1,
  1254. "a broad IN is rejected on the summed counts, before any list is read");
  1255. }
  1256. {
  1257. store.reset_index_plan_stats();
  1258. (void)sig({F("rare", Op::EXISTS, true)});
  1259. check(store.index_plan_stats().exists_scans == 1,
  1260. "EXISTS=true on a sparse field is served from all its postings");
  1261. }
  1262. {
  1263. store.reset_index_plan_stats();
  1264. (void)sig({F("rare", Op::EXISTS, false)});
  1265. auto st = store.index_plan_stats();
  1266. check(st.exists_scans == 0 && st.full_scans == 1,
  1267. "EXISTS=false cannot be served - rows WITHOUT a posting are not "
  1268. "enumerable from the index, so it must scan");
  1269. }
  1270. {
  1271. store.reset_index_plan_stats();
  1272. (void)sig({F("grp", Op::EXISTS, true)});
  1273. auto st = store.index_plan_stats();
  1274. check(st.exists_scans == 0 && st.full_scans == 1,
  1275. "EXISTS=true on a field every row has is not selective, so it scans");
  1276. }
  1277. }
  1278. } // namespace
  1279. int main() {
  1280. std::cout << "=== test_subdb_identity ===\n";
  1281. test_sentinel_roundtrip();
  1282. test_misbound_handle_is_refused();
  1283. test_unstamped_subdb_is_permitted();
  1284. test_identity_key_predicate();
  1285. test_sentinel_invisible_through_store();
  1286. test_vector_subdb_sentinel();
  1287. test_scan_limit_zero_reports_total();
  1288. test_scan_fast_path_matches_general_path();
  1289. test_aborted_write_does_not_poison_the_collection();
  1290. test_filtered_scan_operator_matrix();
  1291. test_concurrent_reads_do_not_rebind_cached_handles();
  1292. test_existing_collection_readable_without_writing_first();
  1293. test_index_tracks_documents_through_every_write();
  1294. test_build_index_over_existing_rows();
  1295. test_index_numeric_equality_matches_scan();
  1296. test_indexed_and_unindexed_plans_agree();
  1297. test_empty_id_reads_as_absent();
  1298. test_range_contains_and_intersection();
  1299. test_index_supplies_ordering();
  1300. test_index_in_and_exists();
  1301. std::cout << "passed: " << g_pass << ", failed: " << g_fail << "\n";
  1302. return g_fail == 0 ? 0 : 1;
  1303. }