|
|
@@ -53,6 +53,7 @@
|
|
|
#include "storage/lmdb_dbi.hpp"
|
|
|
#include "storage/lmdb_env.hpp"
|
|
|
#include "storage/lmdb_txn.hpp"
|
|
|
+#include "storage/secondary_index.hpp"
|
|
|
#include "storage/subdb_identity.hpp"
|
|
|
|
|
|
namespace smartbotic::db::storage {
|
|
|
@@ -307,7 +308,8 @@ LmdbDocumentStore::LmdbDocumentStore(LmdbEnv& env) : env_(env) {
|
|
|
}
|
|
|
|
|
|
unsigned int LmdbDocumentStore::open_for_write(WriteTxn& wtxn,
|
|
|
- std::string_view collection) {
|
|
|
+ std::string_view collection,
|
|
|
+ unsigned int extra_flags) {
|
|
|
std::string key(collection);
|
|
|
{
|
|
|
std::lock_guard<std::mutex> lock(cache_mutex_);
|
|
|
@@ -326,7 +328,8 @@ unsigned int LmdbDocumentStore::open_for_write(WriteTxn& wtxn,
|
|
|
}
|
|
|
MDB_dbi raw_dbi = 0;
|
|
|
std::string name = to_cstr(collection);
|
|
|
- mdb_check(mdb_dbi_open(wtxn.raw(), name.c_str(), MDB_CREATE, &raw_dbi),
|
|
|
+ mdb_check(mdb_dbi_open(wtxn.raw(), name.c_str(), MDB_CREATE | extra_flags,
|
|
|
+ &raw_dbi),
|
|
|
"dbi_open (collection create)");
|
|
|
// Stamp on every fresh open, not only on creation, so sub-dbs that
|
|
|
// pre-date v2.4.4 acquire a sentinel on their next write without a
|
|
|
@@ -444,7 +447,10 @@ size_t LmdbDocumentStore::prime_dbi_cache() {
|
|
|
|
|
|
for (const auto& n : names) {
|
|
|
MDB_dbi dbi = 0;
|
|
|
- int rc = mdb_dbi_open(wtxn.raw(), n.c_str(), 0, &dbi);
|
|
|
+ // Index sub-dbs are MDB_DUPSORT; pass the flag on reopen so the handle
|
|
|
+ // agrees with how the sub-db was created.
|
|
|
+ const unsigned int flags = is_index_subdb(n) ? MDB_DUPSORT : 0u;
|
|
|
+ int rc = mdb_dbi_open(wtxn.raw(), n.c_str(), flags, &dbi);
|
|
|
if (rc == MDB_NOTFOUND) continue; // vanished under us; ignore
|
|
|
if (rc != MDB_SUCCESS) throw_mdb(rc, "dbi_open (prime)");
|
|
|
opened[n] = dbi;
|
|
|
@@ -468,10 +474,276 @@ void LmdbDocumentStore::put(std::string_view collection,
|
|
|
WriteTxn wtxn(env_);
|
|
|
unsigned int dbi = open_for_write(wtxn, collection);
|
|
|
MDB_val k = to_val(id);
|
|
|
+
|
|
|
+ // v2.9.0 — index maintenance runs in THIS transaction, so a write that
|
|
|
+ // throws after this point rolls the index back with the document. An
|
|
|
+ // asynchronously-maintained index would let a query read entries for a row
|
|
|
+ // that was never stored, and return silently wrong rows rather than an error.
|
|
|
+ std::vector<std::pair<std::string, unsigned int>> index_dbis;
|
|
|
+ if (!indexed_fields(collection).empty()) {
|
|
|
+ MDB_val old{0, nullptr};
|
|
|
+ const int rc = mdb_get(wtxn.raw(), dbi, &k, &old);
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) throw_mdb(rc, "get (pre-index)");
|
|
|
+ const std::string_view old_payload =
|
|
|
+ rc == MDB_SUCCESS ? to_sv(old) : std::string_view{};
|
|
|
+ maintainIndexes(wtxn, collection, id, old_payload, &doc, index_dbis);
|
|
|
+ }
|
|
|
+
|
|
|
MDB_val v = to_val(payload);
|
|
|
mdb_check(mdb_put(wtxn.raw(), dbi, &k, &v, 0), "put");
|
|
|
wtxn.commit();
|
|
|
cacheCommittedDbi(collection, dbi);
|
|
|
+ for (const auto& [sub, d] : index_dbis) cacheCommittedDbi(sub, d);
|
|
|
+}
|
|
|
+
|
|
|
+
|
|
|
+// -------------------------------------------------------------------------
|
|
|
+// v2.9.0 secondary index maintenance
|
|
|
+// -------------------------------------------------------------------------
|
|
|
+
|
|
|
+void LmdbDocumentStore::set_indexed_fields(std::string_view collection,
|
|
|
+ std::vector<std::string> fields) {
|
|
|
+ std::lock_guard<std::mutex> lock(index_mutex_);
|
|
|
+ if (fields.empty()) {
|
|
|
+ indexed_fields_.erase(std::string(collection));
|
|
|
+ } else {
|
|
|
+ indexed_fields_[std::string(collection)] = std::move(fields);
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+std::vector<std::string>
|
|
|
+LmdbDocumentStore::indexed_fields(std::string_view collection) {
|
|
|
+ std::lock_guard<std::mutex> lock(index_mutex_);
|
|
|
+ auto it = indexed_fields_.find(std::string(collection));
|
|
|
+ if (it == indexed_fields_.end()) return {};
|
|
|
+ return it->second;
|
|
|
+}
|
|
|
+
|
|
|
+void LmdbDocumentStore::maintainIndexes(
|
|
|
+ WriteTxn& wtxn,
|
|
|
+ std::string_view collection,
|
|
|
+ std::string_view id,
|
|
|
+ std::string_view old_payload,
|
|
|
+ const smartbotic::database::Document* new_doc,
|
|
|
+ std::vector<std::pair<std::string, unsigned int>>& to_cache) {
|
|
|
+
|
|
|
+ const auto fields = indexed_fields(collection);
|
|
|
+ if (fields.empty()) return; // unindexed collections pay nothing
|
|
|
+
|
|
|
+ // Old keys come from the STORED bytes, parsed with yyjson and resolved one
|
|
|
+ // field at a time - never materialised into a Document. On a 505 MB
|
|
|
+ // collection a full decode_document costs ~1ms per write; resolving just the
|
|
|
+ // indexed fields off the yyjson tree is a fraction of that.
|
|
|
+ std::unordered_map<std::string, std::optional<std::string>> old_keys;
|
|
|
+ if (!old_payload.empty()) {
|
|
|
+ yyjson_doc* d = yyjson_read(old_payload.data(), old_payload.size(), 0);
|
|
|
+ if (d) {
|
|
|
+ yyjson_val* root = yyjson_doc_get_root(d);
|
|
|
+ for (const auto& f : fields) {
|
|
|
+ auto v = resolve_from_yyjson(root, f);
|
|
|
+ old_keys[f] = v ? encode_index_key(*v) : std::nullopt;
|
|
|
+ }
|
|
|
+ yyjson_doc_free(d);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ for (const auto& f : fields) {
|
|
|
+ std::optional<std::string> new_key;
|
|
|
+ if (new_doc != nullptr) {
|
|
|
+ // resolveFilterValue is the SAME resolver the scan uses, so an
|
|
|
+ // indexed lookup and a scan cannot disagree about which field a
|
|
|
+ // filter names.
|
|
|
+ auto v = filter_eval::resolveFilterValue(*new_doc, f);
|
|
|
+ if (v) new_key = encode_index_key(*v);
|
|
|
+ }
|
|
|
+ std::optional<std::string> old_key;
|
|
|
+ if (auto it = old_keys.find(f); it != old_keys.end()) old_key = it->second;
|
|
|
+
|
|
|
+ // An unchanged value needs no index write at all - the common case for
|
|
|
+ // an update that touches other fields.
|
|
|
+ if (old_key == new_key) continue;
|
|
|
+
|
|
|
+ const std::string sub = index_subdb_name(collection, f);
|
|
|
+ const unsigned int dbi = open_for_write(wtxn, sub, MDB_DUPSORT);
|
|
|
+ to_cache.emplace_back(sub, dbi);
|
|
|
+
|
|
|
+ if (old_key) {
|
|
|
+ MDB_val k = to_val(*old_key);
|
|
|
+ MDB_val v = to_val(id);
|
|
|
+ // DUPSORT: passing the data removes just THIS id from the value's
|
|
|
+ // posting list, not every id under that value.
|
|
|
+ const int rc = mdb_del(wtxn.raw(), dbi, &k, &v);
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) {
|
|
|
+ throw_mdb(rc, "index del");
|
|
|
+ }
|
|
|
+ }
|
|
|
+ if (new_key) {
|
|
|
+ MDB_val k = to_val(*new_key);
|
|
|
+ MDB_val v = to_val(id);
|
|
|
+ // MDB_NODUPDATA makes a repeat put a no-op rather than an error, so
|
|
|
+ // re-indexing the same row twice is safe (build_index re-runs).
|
|
|
+ const int rc = mdb_put(wtxn.raw(), dbi, &k, &v, MDB_NODUPDATA);
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_KEYEXIST) {
|
|
|
+ throw_mdb(rc, "index put");
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+std::optional<uint64_t>
|
|
|
+LmdbDocumentStore::index_count_eq(std::string_view collection,
|
|
|
+ const std::string& field,
|
|
|
+ const nlohmann::json& value) {
|
|
|
+ ReadTxn rtxn(env_);
|
|
|
+ return countIndexEqTxn(rtxn, collection, field, value);
|
|
|
+}
|
|
|
+
|
|
|
+std::optional<uint64_t>
|
|
|
+LmdbDocumentStore::countIndexEqTxn(ReadTxn& rtxn,
|
|
|
+ std::string_view collection,
|
|
|
+ const std::string& field,
|
|
|
+ const nlohmann::json& value) {
|
|
|
+ const auto key = encode_index_key(value);
|
|
|
+ if (!key) return std::nullopt; // unindexable value; caller must scan
|
|
|
+
|
|
|
+ auto dbi_opt = try_open_for_read(rtxn, index_subdb_name(collection, field));
|
|
|
+ if (!dbi_opt) return std::nullopt; // no such index
|
|
|
+
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
+ mdb_check(mdb_cursor_open(rtxn.raw(), *dbi_opt, &cur), "cursor_open (index count)");
|
|
|
+ struct G { MDB_cursor* c; ~G() { if (c) mdb_cursor_close(c); } } g{cur};
|
|
|
+
|
|
|
+ MDB_val k = to_val(*key);
|
|
|
+ MDB_val v{0, nullptr};
|
|
|
+ const int rc = mdb_cursor_get(cur, &k, &v, MDB_SET);
|
|
|
+ if (rc == MDB_NOTFOUND) return uint64_t{0};
|
|
|
+ if (rc != MDB_SUCCESS) throw_mdb(rc, "cursor_get (index count)");
|
|
|
+
|
|
|
+ // mdb_cursor_count reports the size of this key's duplicate set without
|
|
|
+ // reading the items - which is what makes the selectivity guard cheap.
|
|
|
+ size_t n = 0;
|
|
|
+ mdb_check(mdb_cursor_count(cur, &n), "cursor_count (index)");
|
|
|
+ return static_cast<uint64_t>(n);
|
|
|
+}
|
|
|
+
|
|
|
+std::optional<std::vector<std::string>>
|
|
|
+LmdbDocumentStore::index_lookup_eq(std::string_view collection,
|
|
|
+ const std::string& field,
|
|
|
+ const nlohmann::json& value) {
|
|
|
+ ReadTxn rtxn(env_);
|
|
|
+ return lookupIndexEqTxn(rtxn, collection, field, value);
|
|
|
+}
|
|
|
+
|
|
|
+std::optional<std::vector<std::string>>
|
|
|
+LmdbDocumentStore::lookupIndexEqTxn(ReadTxn& rtxn,
|
|
|
+ std::string_view collection,
|
|
|
+ const std::string& field,
|
|
|
+ const nlohmann::json& value) {
|
|
|
+ const auto key = encode_index_key(value);
|
|
|
+ if (!key) return std::nullopt;
|
|
|
+
|
|
|
+ auto dbi_opt = try_open_for_read(rtxn, index_subdb_name(collection, field));
|
|
|
+ if (!dbi_opt) return std::nullopt;
|
|
|
+
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
+ mdb_check(mdb_cursor_open(rtxn.raw(), *dbi_opt, &cur), "cursor_open (index)");
|
|
|
+ struct G { MDB_cursor* c; ~G() { if (c) mdb_cursor_close(c); } } g{cur};
|
|
|
+
|
|
|
+ std::vector<std::string> ids;
|
|
|
+ MDB_val k = to_val(*key);
|
|
|
+ MDB_val v{0, nullptr};
|
|
|
+ int rc = mdb_cursor_get(cur, &k, &v, MDB_SET);
|
|
|
+ if (rc == MDB_NOTFOUND) return ids; // indexed, genuinely no matches
|
|
|
+ if (rc != MDB_SUCCESS) throw_mdb(rc, "cursor_get (index)");
|
|
|
+
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST_DUP);
|
|
|
+ while (rc == MDB_SUCCESS) {
|
|
|
+ ids.emplace_back(to_sv(v));
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_DUP);
|
|
|
+ }
|
|
|
+ if (rc != MDB_NOTFOUND) throw_mdb(rc, "cursor next_dup (index)");
|
|
|
+ return ids;
|
|
|
+}
|
|
|
+
|
|
|
+LmdbDocumentStore::IndexPlanStats LmdbDocumentStore::index_plan_stats() const {
|
|
|
+ IndexPlanStats st;
|
|
|
+ st.indexed_scans = indexed_scans_.load(std::memory_order_relaxed);
|
|
|
+ st.full_scans = full_scans_.load(std::memory_order_relaxed);
|
|
|
+ st.declined_unselective =
|
|
|
+ declined_unselective_.load(std::memory_order_relaxed);
|
|
|
+ return st;
|
|
|
+}
|
|
|
+
|
|
|
+void LmdbDocumentStore::reset_index_plan_stats() {
|
|
|
+ indexed_scans_.store(0, std::memory_order_relaxed);
|
|
|
+ full_scans_.store(0, std::memory_order_relaxed);
|
|
|
+ declined_unselective_.store(0, std::memory_order_relaxed);
|
|
|
+}
|
|
|
+
|
|
|
+uint64_t LmdbDocumentStore::build_index(std::string_view collection,
|
|
|
+ const std::string& field) {
|
|
|
+ // Walk the collection once and index every row. Idempotent thanks to
|
|
|
+ // MDB_NODUPDATA, so a re-run over an existing index is harmless.
|
|
|
+ std::vector<std::pair<std::string, std::string>> rows; // id -> payload
|
|
|
+ {
|
|
|
+ ReadTxn rtxn(env_);
|
|
|
+ auto dbi_opt = try_open_for_read(rtxn, collection);
|
|
|
+ if (!dbi_opt) return 0;
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
+ mdb_check(mdb_cursor_open(rtxn.raw(), *dbi_opt, &cur), "cursor_open (build)");
|
|
|
+ struct G { MDB_cursor* c; ~G() { if (c) mdb_cursor_close(c); } } g{cur};
|
|
|
+ MDB_val k{0, nullptr};
|
|
|
+ MDB_val v{0, nullptr};
|
|
|
+ int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
+ while (rc == MDB_SUCCESS) {
|
|
|
+ if (!is_identity_key(to_sv(k))) {
|
|
|
+ rows.emplace_back(std::string(to_sv(k)), std::string(to_sv(v)));
|
|
|
+ }
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ const std::string sub = index_subdb_name(collection, field);
|
|
|
+ uint64_t indexed = 0;
|
|
|
+ {
|
|
|
+ WriteTxn wtxn(env_);
|
|
|
+ const unsigned int dbi = open_for_write(wtxn, sub, MDB_DUPSORT);
|
|
|
+ for (const auto& [id, payload] : rows) {
|
|
|
+ yyjson_doc* d = yyjson_read(payload.data(), payload.size(), 0);
|
|
|
+ if (!d) continue;
|
|
|
+ auto val = resolve_from_yyjson(yyjson_doc_get_root(d), field);
|
|
|
+ yyjson_doc_free(d);
|
|
|
+ if (!val) continue; // field absent on this row
|
|
|
+ const auto key = encode_index_key(*val);
|
|
|
+ if (!key) continue; // unindexable value
|
|
|
+ MDB_val kk = to_val(*key);
|
|
|
+ MDB_val vv = to_val(id);
|
|
|
+ const int rc = mdb_put(wtxn.raw(), dbi, &kk, &vv, MDB_NODUPDATA);
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_KEYEXIST) throw_mdb(rc, "index build put");
|
|
|
+ ++indexed;
|
|
|
+ }
|
|
|
+ wtxn.commit();
|
|
|
+ cacheCommittedDbi(sub, dbi);
|
|
|
+ }
|
|
|
+ return indexed;
|
|
|
+}
|
|
|
+
|
|
|
+bool LmdbDocumentStore::drop_index(std::string_view collection,
|
|
|
+ const std::string& field) {
|
|
|
+ const std::string sub = index_subdb_name(collection, field);
|
|
|
+ {
|
|
|
+ ReadTxn rtxn(env_);
|
|
|
+ if (!try_open_for_read(rtxn, sub)) return false;
|
|
|
+ }
|
|
|
+ WriteTxn wtxn(env_);
|
|
|
+ const unsigned int dbi = open_for_write(wtxn, sub, MDB_DUPSORT);
|
|
|
+ mdb_check(mdb_drop(wtxn.raw(), dbi, 1), "mdb_drop (index)");
|
|
|
+ wtxn.commit();
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(cache_mutex_);
|
|
|
+ dbi_cache_.erase(sub);
|
|
|
+ }
|
|
|
+ return true;
|
|
|
}
|
|
|
|
|
|
std::optional<smartbotic::database::Document>
|
|
|
@@ -502,6 +774,20 @@ bool LmdbDocumentStore::del(std::string_view collection, std::string_view id) {
|
|
|
WriteTxn wtxn(env_);
|
|
|
unsigned int dbi = open_for_write(wtxn, collection);
|
|
|
MDB_val k = to_val(id);
|
|
|
+
|
|
|
+ // Remove index entries before the row goes, while its stored bytes are
|
|
|
+ // still readable - they are the only record of which index keys it owns.
|
|
|
+ std::vector<std::pair<std::string, unsigned int>> index_dbis;
|
|
|
+ if (!indexed_fields(collection).empty()) {
|
|
|
+ MDB_val old{0, nullptr};
|
|
|
+ const int grc = mdb_get(wtxn.raw(), dbi, &k, &old);
|
|
|
+ if (grc == MDB_SUCCESS) {
|
|
|
+ maintainIndexes(wtxn, collection, id, to_sv(old), nullptr, index_dbis);
|
|
|
+ } else if (grc != MDB_NOTFOUND) {
|
|
|
+ throw_mdb(grc, "get (pre-index del)");
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
int rc = mdb_del(wtxn.raw(), dbi, &k, nullptr);
|
|
|
if (rc == MDB_NOTFOUND) {
|
|
|
wtxn.commit();
|
|
|
@@ -511,6 +797,7 @@ bool LmdbDocumentStore::del(std::string_view collection, std::string_view id) {
|
|
|
if (rc != MDB_SUCCESS) throw_mdb(rc, "del");
|
|
|
wtxn.commit();
|
|
|
cacheCommittedDbi(collection, dbi);
|
|
|
+ for (const auto& [sub, d] : index_dbis) cacheCommittedDbi(sub, d);
|
|
|
return true;
|
|
|
}
|
|
|
|
|
|
@@ -531,6 +818,70 @@ uint64_t LmdbDocumentStore::count(std::string_view collection) {
|
|
|
return entries;
|
|
|
}
|
|
|
|
|
|
+
|
|
|
+// v2.9.0 — pick an index for this query, or decline.
|
|
|
+//
|
|
|
+// Returns the candidate document ids when an index can serve one of the EQ
|
|
|
+// predicates cheaply, nullopt to mean "walk the whole collection".
|
|
|
+//
|
|
|
+// THE GUARD. An index is not automatically a win. Reading it costs one random
|
|
|
+// mdb_get + full decode_document per candidate, while the scan costs one yyjson
|
|
|
+// parse per row - and those are not the same price. Measured over the same 10,124
|
|
|
+// rows / 193 MB: full decode 3168ms against 369ms for the parse alone, a ratio of
|
|
|
+// ~8.6x. So the index wins only while
|
|
|
+//
|
|
|
+// candidates * 8.6 < total_rows
|
|
|
+//
|
|
|
+// i.e. below roughly 1/8.6 = 11.6% selectivity. kIndexSelectivityDivisor = 10
|
|
|
+// sits just inside that, measured rather than guessed.
|
|
|
+//
|
|
|
+// This matters on real data: on the live `executions` collection
|
|
|
+// status="completed" matches 6,676 of 10,101 rows (66%), so an index on `status`
|
|
|
+// would make that query SLOWER. Declaring an index must not be able to
|
|
|
+// pessimise a query, so the plan is re-decided per value, not per field.
|
|
|
+std::optional<std::vector<std::string>>
|
|
|
+LmdbDocumentStore::planIndexCandidates(ReadTxn& rtxn,
|
|
|
+ std::string_view collection,
|
|
|
+ const smartbotic::database::Query& query,
|
|
|
+ unsigned int coll_dbi) {
|
|
|
+ const auto fields = indexed_fields(collection);
|
|
|
+ if (fields.empty()) return std::nullopt;
|
|
|
+
|
|
|
+ MDB_stat st{};
|
|
|
+ if (mdb_stat(rtxn.raw(), coll_dbi, &st) != MDB_SUCCESS) return std::nullopt;
|
|
|
+ const uint64_t total = st.ms_entries;
|
|
|
+ if (total == 0) return std::nullopt;
|
|
|
+ const uint64_t budget = total / kIndexSelectivityDivisor;
|
|
|
+
|
|
|
+ const std::string* best_field = nullptr;
|
|
|
+ const nlohmann::json* best_value = nullptr;
|
|
|
+ uint64_t best_count = 0;
|
|
|
+ bool found = false;
|
|
|
+
|
|
|
+ for (const auto& f : query.filters) {
|
|
|
+ if (f.op != smartbotic::database::FilterOp::EQ) continue;
|
|
|
+ if (std::find(fields.begin(), fields.end(), f.field) == fields.end()) continue;
|
|
|
+ auto n = countIndexEqTxn(rtxn, collection, f.field, f.value);
|
|
|
+ if (!n) continue; // no index, or unindexable value
|
|
|
+ if (!found || *n < best_count) {
|
|
|
+ found = true;
|
|
|
+ best_count = *n;
|
|
|
+ best_field = &f.field;
|
|
|
+ best_value = &f.value;
|
|
|
+ }
|
|
|
+ }
|
|
|
+ if (!found) return std::nullopt;
|
|
|
+
|
|
|
+ // Zero candidates is a legitimate, extremely selective answer - the query
|
|
|
+ // matches nothing and we can say so without reading a single row.
|
|
|
+ if (best_count > budget) {
|
|
|
+ declined_unselective_.fetch_add(1, std::memory_order_relaxed);
|
|
|
+ return std::nullopt;
|
|
|
+ }
|
|
|
+
|
|
|
+ return lookupIndexEqTxn(rtxn, collection, *best_field, *best_value);
|
|
|
+}
|
|
|
+
|
|
|
ScanResult LmdbDocumentStore::scan(std::string_view collection,
|
|
|
const smartbotic::database::Query& query) {
|
|
|
ScanResult result;
|
|
|
@@ -647,38 +998,73 @@ ScanResult LmdbDocumentStore::scan(std::string_view collection,
|
|
|
};
|
|
|
std::vector<Hit> hits;
|
|
|
|
|
|
- int rc2 = mdb_cursor_get(cursor, &k, &v, MDB_FIRST);
|
|
|
- while (rc2 == MDB_SUCCESS) {
|
|
|
- if (!is_identity_key(to_sv(k))) {
|
|
|
- const auto bytes = to_sv(v);
|
|
|
- yyjson_doc* ydoc = yyjson_read(bytes.data(), bytes.size(), 0);
|
|
|
- if (ydoc) {
|
|
|
- yyjson_val* root = yyjson_doc_get_root(ydoc);
|
|
|
- const bool ok =
|
|
|
- smartbotic::db::storage::filter_eval::matchesFiltersResolved(
|
|
|
- [root](const std::string& f) {
|
|
|
- return resolve_from_yyjson(root, f);
|
|
|
- },
|
|
|
- query.filters);
|
|
|
- if (ok) {
|
|
|
- Hit h;
|
|
|
- h.id = std::string(to_sv(k));
|
|
|
- if (sorting) {
|
|
|
- auto kv = resolve_from_yyjson(root, query.sort->field);
|
|
|
- if (kv) { h.key = *kv; h.hasKey = true; }
|
|
|
- }
|
|
|
- hits.push_back(std::move(h));
|
|
|
- }
|
|
|
- yyjson_doc_free(ydoc);
|
|
|
- }
|
|
|
+ // One row body, two row sources. The predicate evaluation, sort-key
|
|
|
+ // extraction, sorting and pagination below are shared verbatim between
|
|
|
+ // the indexed and unindexed plans - an index changes only WHICH rows are
|
|
|
+ // visited, never how a row is judged. Two copies of the judging would
|
|
|
+ // drift, and the divergence would surface as wrong rows rather than an
|
|
|
+ // error.
|
|
|
+ const auto consider = [&](std::string_view row_id, std::string_view bytes) {
|
|
|
+ yyjson_doc* ydoc = yyjson_read(bytes.data(), bytes.size(), 0);
|
|
|
+ if (!ydoc) {
|
|
|
// A row that will not parse cannot be matched; skipping it is the
|
|
|
// same outcome the old path reached by throwing on decode, minus
|
|
|
// failing the whole query for one bad row.
|
|
|
+ return;
|
|
|
}
|
|
|
- rc2 = mdb_cursor_get(cursor, &k, &v, MDB_NEXT);
|
|
|
+ yyjson_val* root = yyjson_doc_get_root(ydoc);
|
|
|
+ const bool ok =
|
|
|
+ smartbotic::db::storage::filter_eval::matchesFiltersResolved(
|
|
|
+ [root](const std::string& f) {
|
|
|
+ return resolve_from_yyjson(root, f);
|
|
|
+ },
|
|
|
+ query.filters);
|
|
|
+ if (ok) {
|
|
|
+ Hit h;
|
|
|
+ h.id = std::string(row_id);
|
|
|
+ if (sorting) {
|
|
|
+ auto kv = resolve_from_yyjson(root, query.sort->field);
|
|
|
+ if (kv) { h.key = *kv; h.hasKey = true; }
|
|
|
+ }
|
|
|
+ hits.push_back(std::move(h));
|
|
|
+ }
|
|
|
+ yyjson_doc_free(ydoc);
|
|
|
+ };
|
|
|
+
|
|
|
+ auto candidates = planIndexCandidates(rtxn, collection, query, *dbi_opt);
|
|
|
+ if (candidates) {
|
|
|
+ indexed_scans_.fetch_add(1, std::memory_order_relaxed);
|
|
|
+ } else {
|
|
|
+ full_scans_.fetch_add(1, std::memory_order_relaxed);
|
|
|
}
|
|
|
- if (rc2 != MDB_NOTFOUND && rc2 != MDB_SUCCESS) {
|
|
|
- throw_mdb(rc2, "cursor_get");
|
|
|
+
|
|
|
+ if (candidates) {
|
|
|
+ // Indexed plan: visit only the rows the index named. Every filter is
|
|
|
+ // still applied to each one, so the index is allowed to be a
|
|
|
+ // superset - it narrows work, it does not decide the answer.
|
|
|
+ for (const auto& cid : *candidates) {
|
|
|
+ MDB_val ck = to_val(cid);
|
|
|
+ MDB_val cv{0, nullptr};
|
|
|
+ if (mdb_get(rtxn.raw(), *dbi_opt, &ck, &cv) != MDB_SUCCESS) {
|
|
|
+ // A posting with no row means the index is ahead of the
|
|
|
+ // collection, which the same-transaction maintenance is meant
|
|
|
+ // to prevent. Skipping is the safe reading: report what
|
|
|
+ // exists, never a row that does not.
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ consider(cid, to_sv(cv));
|
|
|
+ }
|
|
|
+ } else {
|
|
|
+ int rc2 = mdb_cursor_get(cursor, &k, &v, MDB_FIRST);
|
|
|
+ while (rc2 == MDB_SUCCESS) {
|
|
|
+ if (!is_identity_key(to_sv(k))) {
|
|
|
+ consider(to_sv(k), to_sv(v));
|
|
|
+ }
|
|
|
+ rc2 = mdb_cursor_get(cursor, &k, &v, MDB_NEXT);
|
|
|
+ }
|
|
|
+ if (rc2 != MDB_NOTFOUND && rc2 != MDB_SUCCESS) {
|
|
|
+ throw_mdb(rc2, "cursor_get");
|
|
|
+ }
|
|
|
}
|
|
|
|
|
|
if (sorting) {
|