|
|
@@ -534,58 +534,68 @@ void LmdbDocumentStore::maintainIndexes(
|
|
|
// 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;
|
|
|
+ // A field contributes a SET of keys, not one: an array value contributes one
|
|
|
+ // key per element so CONTAINS can be served. So the update is a set
|
|
|
+ // difference - remove the keys the row no longer owns, add the ones it gained,
|
|
|
+ // leave the rest untouched.
|
|
|
+ std::unordered_map<std::string, std::vector<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;
|
|
|
+ if (v) old_keys[f] = encode_index_keys(*v);
|
|
|
}
|
|
|
yyjson_doc_free(d);
|
|
|
}
|
|
|
}
|
|
|
|
|
|
for (const auto& f : fields) {
|
|
|
- std::optional<std::string> new_key;
|
|
|
+ std::vector<std::string> new_k;
|
|
|
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);
|
|
|
+ if (auto v = filter_eval::resolveFilterValue(*new_doc, f)) {
|
|
|
+ new_k = encode_index_keys(*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;
|
|
|
+ std::vector<std::string> old_k;
|
|
|
+ if (auto it = old_keys.find(f); it != old_keys.end()) old_k = it->second;
|
|
|
+
|
|
|
+ // Unchanged needs no index write at all - the common case for an update
|
|
|
+ // that touches other fields. Both sides are sorted and deduplicated by
|
|
|
+ // encode_index_keys, so this comparison is exact.
|
|
|
+ if (old_k == new_k) continue;
|
|
|
+
|
|
|
+ std::vector<std::string> to_remove;
|
|
|
+ std::vector<std::string> to_add;
|
|
|
+ std::set_difference(old_k.begin(), old_k.end(), new_k.begin(), new_k.end(),
|
|
|
+ std::back_inserter(to_remove));
|
|
|
+ std::set_difference(new_k.begin(), new_k.end(), old_k.begin(), old_k.end(),
|
|
|
+ std::back_inserter(to_add));
|
|
|
+ if (to_remove.empty() && to_add.empty()) 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);
|
|
|
+ for (const auto& key : to_remove) {
|
|
|
+ MDB_val k = to_val(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 (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) throw_mdb(rc, "index del");
|
|
|
}
|
|
|
- if (new_key) {
|
|
|
- MDB_val k = to_val(*new_key);
|
|
|
+ for (const auto& key : to_add) {
|
|
|
+ MDB_val k = to_val(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");
|
|
|
- }
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_KEYEXIST) throw_mdb(rc, "index put");
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
@@ -671,6 +681,8 @@ LmdbDocumentStore::IndexPlanStats LmdbDocumentStore::index_plan_stats() const {
|
|
|
st.full_scans = full_scans_.load(std::memory_order_relaxed);
|
|
|
st.declined_unselective =
|
|
|
declined_unselective_.load(std::memory_order_relaxed);
|
|
|
+ st.intersected_scans = intersected_scans_.load(std::memory_order_relaxed);
|
|
|
+ st.range_scans = range_scans_.load(std::memory_order_relaxed);
|
|
|
return st;
|
|
|
}
|
|
|
|
|
|
@@ -678,6 +690,8 @@ 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);
|
|
|
+ intersected_scans_.store(0, std::memory_order_relaxed);
|
|
|
+ range_scans_.store(0, std::memory_order_relaxed);
|
|
|
}
|
|
|
|
|
|
std::optional<LmdbDocumentStore::IndexStats>
|
|
|
@@ -715,6 +729,86 @@ LmdbDocumentStore::index_stats(std::string_view collection,
|
|
|
return out;
|
|
|
}
|
|
|
|
|
|
+
|
|
|
+std::optional<std::vector<std::string>>
|
|
|
+LmdbDocumentStore::lookupIndexRangeTxn(ReadTxn& rtxn,
|
|
|
+ std::string_view collection,
|
|
|
+ const std::string& field,
|
|
|
+ smartbotic::database::FilterOp op,
|
|
|
+ const nlohmann::json& value,
|
|
|
+ uint64_t budget) {
|
|
|
+ using Op = smartbotic::database::FilterOp;
|
|
|
+ if (op != Op::GT && op != Op::GTE && op != Op::LT && op != Op::LTE) {
|
|
|
+ return std::nullopt;
|
|
|
+ }
|
|
|
+ const auto bound = encode_index_key(value);
|
|
|
+ if (!bound) 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 range)");
|
|
|
+ struct G { MDB_cursor* c; ~G() { if (c) mdb_cursor_close(c); } } g{cur};
|
|
|
+
|
|
|
+ MDB_val k{0, nullptr};
|
|
|
+ MDB_val v{0, nullptr};
|
|
|
+
|
|
|
+ // Type homogeneity. Comparing the first and last key's tag is enough: keys are
|
|
|
+ // ordered, so one tag at both ends means one tag throughout. The identity
|
|
|
+ // sentinel sorts before every tag (it starts with NUL), so skip it.
|
|
|
+ int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
+ while (rc == MDB_SUCCESS && is_identity_key(to_sv(k))) {
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_NODUP);
|
|
|
+ }
|
|
|
+ if (rc == MDB_NOTFOUND) return std::vector<std::string>{}; // empty index
|
|
|
+ if (rc != MDB_SUCCESS) throw_mdb(rc, "cursor first (index range)");
|
|
|
+ const unsigned char first_tag = index_key_tag(to_sv(k));
|
|
|
+
|
|
|
+ MDB_val lk{0, nullptr};
|
|
|
+ MDB_val lv{0, nullptr};
|
|
|
+ rc = mdb_cursor_get(cur, &lk, &lv, MDB_LAST);
|
|
|
+ if (rc != MDB_SUCCESS) throw_mdb(rc, "cursor last (index range)");
|
|
|
+ if (index_key_tag(to_sv(lk)) != first_tag) return std::nullopt; // mixed types
|
|
|
+ if (first_tag != index_key_tag(*bound)) return std::nullopt; // wrong type
|
|
|
+
|
|
|
+ const std::string_view bsv = *bound;
|
|
|
+ std::vector<std::string> ids;
|
|
|
+
|
|
|
+ const bool ascending = (op == Op::GT || op == Op::GTE);
|
|
|
+ if (ascending) {
|
|
|
+ MDB_val sk = to_val(bsv);
|
|
|
+ rc = mdb_cursor_get(cur, &sk, &v, MDB_SET_RANGE); // first key >= bound
|
|
|
+ // For GT, step past the bound's own duplicates.
|
|
|
+ while (rc == MDB_SUCCESS && op == Op::GT && to_sv(sk) == bsv) {
|
|
|
+ rc = mdb_cursor_get(cur, &sk, &v, MDB_NEXT_NODUP);
|
|
|
+ }
|
|
|
+ k = sk;
|
|
|
+ } else {
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
+ }
|
|
|
+
|
|
|
+ while (rc == MDB_SUCCESS) {
|
|
|
+ const auto key = to_sv(k);
|
|
|
+ if (!is_identity_key(key)) {
|
|
|
+ if (index_key_tag(key) != first_tag) break; // left the type
|
|
|
+ if (!ascending) {
|
|
|
+ const int cmp = key.compare(bsv);
|
|
|
+ if (op == Op::LT ? cmp >= 0 : cmp > 0) break;
|
|
|
+ }
|
|
|
+ ids.emplace_back(to_sv(v));
|
|
|
+ // Bail out rather than truncate. A partial list would drop matching
|
|
|
+ // rows; declining sends the query to the scan, which is merely slower.
|
|
|
+ if (ids.size() > budget) return std::nullopt;
|
|
|
+ }
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT);
|
|
|
+ }
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) {
|
|
|
+ throw_mdb(rc, "cursor next (index range)");
|
|
|
+ }
|
|
|
+ return ids;
|
|
|
+}
|
|
|
+
|
|
|
uint64_t LmdbDocumentStore::build_index(std::string_view collection,
|
|
|
const std::string& field) {
|
|
|
// Walk the collection once and index every row. Idempotent thanks to
|
|
|
@@ -749,12 +843,16 @@ uint64_t LmdbDocumentStore::build_index(std::string_view collection,
|
|
|
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");
|
|
|
+ const auto keys = encode_index_keys(*val);
|
|
|
+ if (keys.empty()) continue; // unindexable value
|
|
|
+ for (const auto& key : keys) {
|
|
|
+ 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();
|
|
|
@@ -889,6 +987,7 @@ LmdbDocumentStore::planIndexCandidates(ReadTxn& rtxn,
|
|
|
std::string_view collection,
|
|
|
const smartbotic::database::Query& query,
|
|
|
unsigned int coll_dbi) {
|
|
|
+ using Op = smartbotic::database::FilterOp;
|
|
|
const auto fields = indexed_fields(collection);
|
|
|
if (fields.empty()) return std::nullopt;
|
|
|
|
|
|
@@ -898,33 +997,84 @@ LmdbDocumentStore::planIndexCandidates(ReadTxn& rtxn,
|
|
|
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;
|
|
|
+ const auto is_indexed = [&](const std::string& f) {
|
|
|
+ return std::find(fields.begin(), fields.end(), f) != fields.end();
|
|
|
+ };
|
|
|
|
|
|
+ // --- exact predicates: EQ, and CONTAINS (which asks about an array ELEMENT,
|
|
|
+ // and elements are exactly what an array contributes to the index) ---
|
|
|
+ struct Exact {
|
|
|
+ const smartbotic::database::Filter* f;
|
|
|
+ uint64_t count;
|
|
|
+ };
|
|
|
+ std::vector<Exact> exacts;
|
|
|
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 (f.op != Op::EQ && f.op != Op::CONTAINS) continue;
|
|
|
+ if (!is_indexed(f.field)) continue;
|
|
|
+ if (auto n = countIndexEqTxn(rtxn, collection, f.field, f.value)) {
|
|
|
+ exacts.push_back({&f, *n});
|
|
|
}
|
|
|
}
|
|
|
- if (!found) return std::nullopt;
|
|
|
+ std::sort(exacts.begin(), exacts.end(),
|
|
|
+ [](const Exact& a, const Exact& b) { return a.count < b.count; });
|
|
|
+
|
|
|
+ // Cheapest plan: one selective exact predicate. Zero candidates is a
|
|
|
+ // legitimate, maximally selective answer - the query matches nothing and we
|
|
|
+ // never read a row.
|
|
|
+ if (!exacts.empty() && exacts.front().count <= budget) {
|
|
|
+ const auto* f = exacts.front().f;
|
|
|
+ return lookupIndexEqTxn(rtxn, collection, f->field, f->value);
|
|
|
+ }
|
|
|
|
|
|
- // 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;
|
|
|
+ // --- intersection: two predicates each too broad alone, selective together.
|
|
|
+ //
|
|
|
+ // Worth trying because postings are ids, not documents: reading two lists
|
|
|
+ // costs nothing like the decode_document per candidate that the budget is
|
|
|
+ // protecting against. Only the INTERSECTION has to fit the budget.
|
|
|
+ //
|
|
|
+ // Correctness: each list is a superset of the rows matching its own
|
|
|
+ // predicate, so the intersection is a superset of the rows matching both -
|
|
|
+ // and every candidate is still checked against every filter afterwards.
|
|
|
+ if (exacts.size() >= 2) {
|
|
|
+ const uint64_t postings_cap = total; // reading ids is cheap; decoding is not
|
|
|
+ if (exacts[0].count <= postings_cap && exacts[1].count <= postings_cap) {
|
|
|
+ auto a = lookupIndexEqTxn(rtxn, collection, exacts[0].f->field,
|
|
|
+ exacts[0].f->value);
|
|
|
+ auto b = lookupIndexEqTxn(rtxn, collection, exacts[1].f->field,
|
|
|
+ exacts[1].f->value);
|
|
|
+ if (a && b) {
|
|
|
+ std::sort(a->begin(), a->end());
|
|
|
+ std::sort(b->begin(), b->end());
|
|
|
+ std::vector<std::string> both;
|
|
|
+ std::set_intersection(a->begin(), a->end(), b->begin(), b->end(),
|
|
|
+ std::back_inserter(both));
|
|
|
+ if (both.size() <= budget) {
|
|
|
+ intersected_scans_.fetch_add(1, std::memory_order_relaxed);
|
|
|
+ return both;
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ // --- ranges. There is no cheap count for a range, so the walk itself is the
|
|
|
+ // probe: it collects ids and gives up the moment it passes the budget,
|
|
|
+ // returning no-plan rather than a truncated list.
|
|
|
+ for (const auto& f : query.filters) {
|
|
|
+ if (f.op != Op::GT && f.op != Op::GTE && f.op != Op::LT && f.op != Op::LTE) {
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ if (!is_indexed(f.field)) continue;
|
|
|
+ if (auto ids = lookupIndexRangeTxn(rtxn, collection, f.field, f.op,
|
|
|
+ f.value, budget)) {
|
|
|
+ range_scans_.fetch_add(1, std::memory_order_relaxed);
|
|
|
+ return ids;
|
|
|
+ }
|
|
|
}
|
|
|
|
|
|
- return lookupIndexEqTxn(rtxn, collection, *best_field, *best_value);
|
|
|
+ if (!exacts.empty()) {
|
|
|
+ declined_unselective_.fetch_add(1, std::memory_order_relaxed);
|
|
|
+ }
|
|
|
+ return std::nullopt;
|
|
|
}
|
|
|
|
|
|
ScanResult LmdbDocumentStore::scan(std::string_view collection,
|