|
@@ -16,6 +16,7 @@
|
|
|
#include <iostream>
|
|
#include <iostream>
|
|
|
#include <cstdio>
|
|
#include <cstdio>
|
|
|
#include <algorithm>
|
|
#include <algorithm>
|
|
|
|
|
+#include <thread>
|
|
|
#include <set>
|
|
#include <set>
|
|
|
#include <string>
|
|
#include <string>
|
|
|
#include <unistd.h>
|
|
#include <unistd.h>
|
|
@@ -543,6 +544,157 @@ void test_filtered_scan_operator_matrix() {
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+
|
|
|
|
|
+// v2.8.1 — concurrent reads must not rebind a cached handle.
|
|
|
|
|
+//
|
|
|
|
|
+// lmdb.h: "This function [mdb_dbi_open] must not be called from multiple
|
|
|
|
|
+// concurrent transactions in the same process. A transaction that uses this
|
|
|
|
|
+// function must finish (either commit or abort) before any other transaction in
|
|
|
|
|
+// the process may use this function."
|
|
|
|
|
+//
|
|
|
|
|
+// The old read path violated that on every read of a collection this process had
|
|
|
|
|
+// not yet written, from every gRPC thread at once. MDB_dbi is an index into the
|
|
|
|
|
+// env's shared handle table, so churning it silently REBOUND handles that
|
|
|
|
|
+// committed writes had already cached. Live on 2.8.0-3 that produced 41 refusals
|
|
|
|
|
+// in 45 minutes, all "caller asked for 'image_hashes' but the handle addresses
|
|
|
|
|
+// 'nsfw_images'" - and because each refusal bumped mirror drift, every read in
|
|
|
|
|
+// the process fell back to MemoryStore permanently.
|
|
|
|
|
+//
|
|
|
|
|
+// The fix primes all handles in one committed write txn at construction, so only
|
|
|
|
|
+// write txns ever call mdb_dbi_open and LMDB serialises those itself.
|
|
|
|
|
+//
|
|
|
|
|
+// HONEST LIMITATION: this test does NOT reproduce the race. Reverting the fix
|
|
|
|
|
+// leaves it green apart from the primed_count assertion - the torn slot-table
|
|
|
|
|
+// update needs an interleaving of concurrent mdb_dbi_open calls that 8 threads
|
|
|
|
|
+// over 10 collections does not reliably hit. What the test does pin is the
|
|
|
|
|
+// invariant the race violates (no read returns another collection's document, no
|
|
|
|
|
+// write is refused by the sentinel) plus the mechanism that removes the race:
|
|
|
|
|
+// handles are opened up front, so the read path has no mdb_dbi_open left to call.
|
|
|
|
|
+// The deterministic evidence is the documented contract in lmdb.h and the
|
|
|
|
|
+// production log.
|
|
|
|
|
+void test_concurrent_reads_do_not_rebind_cached_handles() {
|
|
|
|
|
+ const std::string path = make_tmpdir("dbi-concurrent");
|
|
|
|
|
+ const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
|
|
|
|
|
+ struct Cleanup {
|
|
|
|
|
+ const std::string& p;
|
|
|
|
|
+ ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
|
|
|
|
|
+ } cleanup{path};
|
|
|
|
|
+
|
|
|
|
|
+ // Enough collections that slot churn has somewhere to go.
|
|
|
|
|
+ const std::vector<std::string> colls = {
|
|
|
|
|
+ "alpha", "beta", "gamma", "delta", "epsilon",
|
|
|
|
|
+ "zeta", "eta", "theta", "iota", "kappa"};
|
|
|
|
|
+
|
|
|
|
|
+ {
|
|
|
|
|
+ LmdbEnv env(opts);
|
|
|
|
|
+ LmdbDocumentStore writer(env);
|
|
|
|
|
+ for (const auto& c : colls) {
|
|
|
|
|
+ Document d;
|
|
|
|
|
+ d.id = "seed";
|
|
|
|
|
+ d.collection = c;
|
|
|
|
|
+ d.set_data(nlohmann::json{{"who", c}});
|
|
|
|
|
+ writer.put(c, "seed", d);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Restart: fresh env, so the shared handle table starts empty and the store
|
|
|
|
|
+ // must prime it.
|
|
|
|
|
+ LmdbEnv env2(opts);
|
|
|
|
|
+ LmdbDocumentStore store(env2);
|
|
|
|
|
+ check(store.prime_error().empty(), "priming succeeded on reopen");
|
|
|
|
|
+ check(store.primed_count() >= colls.size(),
|
|
|
|
|
+ "priming opened a handle for every existing sub-db");
|
|
|
|
|
+
|
|
|
|
|
+ // Hammer reads from many threads. Every one of these used to call
|
|
|
|
|
+ // mdb_dbi_open inside its own read txn.
|
|
|
|
|
+ std::atomic<int> read_failures{0};
|
|
|
|
|
+ std::atomic<int> wrong_data{0};
|
|
|
|
|
+ {
|
|
|
|
|
+ std::vector<std::thread> threads;
|
|
|
|
|
+ for (int t = 0; t < 8; ++t) {
|
|
|
|
|
+ threads.emplace_back([&, t]() {
|
|
|
|
|
+ for (int i = 0; i < 40; ++i) {
|
|
|
|
|
+ const auto& c = colls[(t + i) % colls.size()];
|
|
|
|
|
+ try {
|
|
|
|
|
+ auto got = store.get(c, "seed");
|
|
|
|
|
+ if (!got) { ++read_failures; continue; }
|
|
|
|
|
+ // A rebound handle reads a STRANGER's sub-db, so the
|
|
|
|
|
+ // document that comes back belongs to another collection.
|
|
|
|
|
+ if (got->data().value("who", std::string{}) != c) ++wrong_data;
|
|
|
|
|
+ } catch (const std::exception&) {
|
|
|
|
|
+ ++read_failures;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+ for (auto& th : threads) th.join();
|
|
|
|
|
+ }
|
|
|
|
|
+ check(read_failures.load() == 0, "concurrent reads all succeeded");
|
|
|
|
|
+ check(wrong_data.load() == 0,
|
|
|
|
|
+ "no read returned another collection's document - a rebound handle "
|
|
|
|
|
+ "addresses whichever sub-db now occupies its slot");
|
|
|
|
|
+
|
|
|
|
|
+ // Now write to every collection through the cached handles. This is where
|
|
|
|
|
+ // the live failure surfaced: the sentinel refused the write.
|
|
|
|
|
+ int write_failures = 0;
|
|
|
|
|
+ for (const auto& c : colls) {
|
|
|
|
|
+ try {
|
|
|
|
|
+ Document d;
|
|
|
|
|
+ d.id = "after";
|
|
|
|
|
+ d.collection = c;
|
|
|
|
|
+ d.set_data(nlohmann::json{{"who", c}});
|
|
|
|
|
+ store.put(c, "after", d);
|
|
|
|
|
+ } catch (const std::exception&) {
|
|
|
|
|
+ ++write_failures;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ check(write_failures == 0,
|
|
|
|
|
+ "writes through primed handles are not refused by the identity "
|
|
|
|
|
+ "sentinel - the refusal is what production saw");
|
|
|
|
|
+
|
|
|
|
|
+ for (const auto& c : colls) {
|
|
|
|
|
+ check(store.count(c) == 2, ("both documents readable in " + c).c_str());
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// A collection that exists on disk but has NOT been written by this process must
|
|
|
|
|
+// still be readable. The read path no longer opens handles on demand, so if
|
|
|
|
|
+// priming missed anything a read would report the collection as EMPTY - a
|
|
|
|
|
+// silent-wrong-data failure worse than the bug being fixed.
|
|
|
|
|
+void test_existing_collection_readable_without_writing_first() {
|
|
|
|
|
+ const std::string path = make_tmpdir("dbi-prime-read");
|
|
|
|
|
+ const LmdbEnvOpts opts{path, 64ULL << 20, 256, 126, false};
|
|
|
|
|
+ struct Cleanup {
|
|
|
|
|
+ const std::string& p;
|
|
|
|
|
+ ~Cleanup() { std::error_code ec; fs::remove_all(p, ec); }
|
|
|
|
|
+ } cleanup{path};
|
|
|
|
|
+
|
|
|
|
|
+ {
|
|
|
|
|
+ LmdbEnv env(opts);
|
|
|
|
|
+ LmdbDocumentStore writer(env);
|
|
|
|
|
+ Document d;
|
|
|
|
|
+ d.id = "only";
|
|
|
|
|
+ d.collection = "archive";
|
|
|
|
|
+ d.set_data(nlohmann::json{{"kept", true}});
|
|
|
|
|
+ writer.put("archive", "only", d);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ LmdbEnv env2(opts);
|
|
|
|
|
+ LmdbDocumentStore reader(env2);
|
|
|
|
|
+
|
|
|
|
|
+ // Read-only access, never a write on this collection in this process.
|
|
|
|
|
+ auto got = reader.get("archive", "only");
|
|
|
|
|
+ check(got.has_value(), "a never-written-here collection is still readable");
|
|
|
|
|
+ check(reader.count("archive") == 1, "count sees it");
|
|
|
|
|
+ smartbotic::database::Query q; q.limit = 10;
|
|
|
|
|
+ check(reader.scan("archive", q).documents.size() == 1, "scan sees it");
|
|
|
|
|
+
|
|
|
|
|
+ // And a collection that genuinely does not exist still reads as absent.
|
|
|
|
|
+ check(!reader.get("nosuch", "x").has_value(),
|
|
|
|
|
+ "a missing collection is still absent, not an error");
|
|
|
|
|
+ check(reader.count("nosuch") == 0, "and counts zero");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
} // namespace
|
|
} // namespace
|
|
|
|
|
|
|
|
int main() {
|
|
int main() {
|
|
@@ -557,6 +709,8 @@ int main() {
|
|
|
test_scan_fast_path_matches_general_path();
|
|
test_scan_fast_path_matches_general_path();
|
|
|
test_aborted_write_does_not_poison_the_collection();
|
|
test_aborted_write_does_not_poison_the_collection();
|
|
|
test_filtered_scan_operator_matrix();
|
|
test_filtered_scan_operator_matrix();
|
|
|
|
|
+ test_concurrent_reads_do_not_rebind_cached_handles();
|
|
|
|
|
+ test_existing_collection_readable_without_writing_first();
|
|
|
|
|
|
|
|
std::cout << "passed: " << g_pass << ", failed: " << g_fail << "\n";
|
|
std::cout << "passed: " << g_pass << ", failed: " << g_fail << "\n";
|
|
|
return g_fail == 0 ? 0 : 1;
|
|
return g_fail == 0 ? 0 : 1;
|