|
@@ -37,6 +37,7 @@
|
|
|
#include <cctype>
|
|
#include <cctype>
|
|
|
#include <cstring>
|
|
#include <cstring>
|
|
|
#include <functional>
|
|
#include <functional>
|
|
|
|
|
+#include <limits>
|
|
|
#include <regex>
|
|
#include <regex>
|
|
|
#include <sstream>
|
|
#include <sstream>
|
|
|
#include <stdexcept>
|
|
#include <stdexcept>
|
|
@@ -529,6 +530,7 @@ void LmdbDocumentStore::maintainIndexes(
|
|
|
|
|
|
|
|
const auto fields = indexed_fields(collection);
|
|
const auto fields = indexed_fields(collection);
|
|
|
if (fields.empty()) return; // unindexed collections pay nothing
|
|
if (fields.empty()) return; // unindexed collections pay nothing
|
|
|
|
|
+ const auto uniques = unique_fields(collection);
|
|
|
|
|
|
|
|
// Old keys come from the STORED bytes, parsed with yyjson and resolved one
|
|
// 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
|
|
// field at a time - never materialised into a Document. On a 505 MB
|
|
@@ -553,12 +555,14 @@ void LmdbDocumentStore::maintainIndexes(
|
|
|
|
|
|
|
|
for (const auto& f : fields) {
|
|
for (const auto& f : fields) {
|
|
|
std::vector<std::string> new_k;
|
|
std::vector<std::string> new_k;
|
|
|
|
|
+ bool new_is_array = false;
|
|
|
if (new_doc != nullptr) {
|
|
if (new_doc != nullptr) {
|
|
|
// resolveFilterValue is the SAME resolver the scan uses, so an
|
|
// resolveFilterValue is the SAME resolver the scan uses, so an
|
|
|
// indexed lookup and a scan cannot disagree about which field a
|
|
// indexed lookup and a scan cannot disagree about which field a
|
|
|
// filter names.
|
|
// filter names.
|
|
|
if (auto v = filter_eval::resolveFilterValue(*new_doc, f)) {
|
|
if (auto v = filter_eval::resolveFilterValue(*new_doc, f)) {
|
|
|
new_k = encode_index_keys(*v);
|
|
new_k = encode_index_keys(*v);
|
|
|
|
|
+ new_is_array = v->is_array();
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
std::vector<std::string> old_k;
|
|
std::vector<std::string> old_k;
|
|
@@ -581,6 +585,11 @@ void LmdbDocumentStore::maintainIndexes(
|
|
|
const unsigned int dbi = open_for_write(wtxn, sub, MDB_DUPSORT);
|
|
const unsigned int dbi = open_for_write(wtxn, sub, MDB_DUPSORT);
|
|
|
to_cache.emplace_back(sub, dbi);
|
|
to_cache.emplace_back(sub, dbi);
|
|
|
|
|
|
|
|
|
|
+ // An array value means one row owns several postings, which breaks the
|
|
|
|
|
+ // "duplicate count == match count" identity a count depends on. Record it
|
|
|
|
|
+ // once, in the same transaction, so a count can check cheaply.
|
|
|
|
|
+ if (new_is_array) markIndexMultiValued(wtxn, dbi);
|
|
|
|
|
+
|
|
|
for (const auto& key : to_remove) {
|
|
for (const auto& key : to_remove) {
|
|
|
MDB_val k = to_val(key);
|
|
MDB_val k = to_val(key);
|
|
|
MDB_val v = to_val(id);
|
|
MDB_val v = to_val(id);
|
|
@@ -589,6 +598,35 @@ void LmdbDocumentStore::maintainIndexes(
|
|
|
const int rc = mdb_del(wtxn.raw(), dbi, &k, &v);
|
|
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");
|
|
|
}
|
|
}
|
|
|
|
|
+ // UNIQUE enforcement, before any posting is written. Checked here rather
|
|
|
|
|
+ // than at a handler because this is the only place that knows the value's
|
|
|
|
|
+ // encoded key, and because it must share the document's transaction: a
|
|
|
|
|
+ // rejection must leave neither the row nor an index entry behind.
|
|
|
|
|
+ if (std::find(uniques.begin(), uniques.end(), f) != uniques.end()) {
|
|
|
|
|
+ for (const auto& key : to_add) {
|
|
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
|
|
+ mdb_check(mdb_cursor_open(wtxn.raw(), dbi, &cur),
|
|
|
|
|
+ "cursor_open (unique check)");
|
|
|
|
|
+ struct G { MDB_cursor* c; ~G() { if (c) mdb_cursor_close(c); } } g{cur};
|
|
|
|
|
+ MDB_val ck = to_val(key);
|
|
|
|
|
+ MDB_val cv{0, nullptr};
|
|
|
|
|
+ int crc = mdb_cursor_get(cur, &ck, &cv, MDB_SET);
|
|
|
|
|
+ while (crc == MDB_SUCCESS) {
|
|
|
|
|
+ const auto holder = to_sv(cv);
|
|
|
|
|
+ // Its own posting is not a conflict - that is this row being
|
|
|
|
|
+ // rewritten with the value it already had.
|
|
|
|
|
+ if (holder != id) {
|
|
|
|
|
+ throw UniqueViolation(std::string(collection), f,
|
|
|
|
|
+ std::string(holder));
|
|
|
|
|
+ }
|
|
|
|
|
+ crc = mdb_cursor_get(cur, &ck, &cv, MDB_NEXT_DUP);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (crc != MDB_SUCCESS && crc != MDB_NOTFOUND) {
|
|
|
|
|
+ throw_mdb(crc, "cursor (unique check)");
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
for (const auto& key : to_add) {
|
|
for (const auto& key : to_add) {
|
|
|
MDB_val k = to_val(key);
|
|
MDB_val k = to_val(key);
|
|
|
MDB_val v = to_val(id);
|
|
MDB_val v = to_val(id);
|
|
@@ -686,6 +724,8 @@ LmdbDocumentStore::IndexPlanStats LmdbDocumentStore::index_plan_stats() const {
|
|
|
st.ordered_scans = ordered_scans_.load(std::memory_order_relaxed);
|
|
st.ordered_scans = ordered_scans_.load(std::memory_order_relaxed);
|
|
|
st.union_scans = union_scans_.load(std::memory_order_relaxed);
|
|
st.union_scans = union_scans_.load(std::memory_order_relaxed);
|
|
|
st.exists_scans = exists_scans_.load(std::memory_order_relaxed);
|
|
st.exists_scans = exists_scans_.load(std::memory_order_relaxed);
|
|
|
|
|
+ st.counted_scans = counted_scans_.load(std::memory_order_relaxed);
|
|
|
|
|
+ st.range_ordered_scans = range_ordered_scans_.load(std::memory_order_relaxed);
|
|
|
return st;
|
|
return st;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -698,6 +738,70 @@ void LmdbDocumentStore::reset_index_plan_stats() {
|
|
|
ordered_scans_.store(0, std::memory_order_relaxed);
|
|
ordered_scans_.store(0, std::memory_order_relaxed);
|
|
|
union_scans_.store(0, std::memory_order_relaxed);
|
|
union_scans_.store(0, std::memory_order_relaxed);
|
|
|
exists_scans_.store(0, std::memory_order_relaxed);
|
|
exists_scans_.store(0, std::memory_order_relaxed);
|
|
|
|
|
+ counted_scans_.store(0, std::memory_order_relaxed);
|
|
|
|
|
+ range_ordered_scans_.store(0, std::memory_order_relaxed);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+std::vector<LmdbDocumentStore::IndexValue>
|
|
|
|
|
+LmdbDocumentStore::index_values(std::string_view collection,
|
|
|
|
|
+ const std::string& field,
|
|
|
|
|
+ size_t limit,
|
|
|
|
|
+ bool ascending) {
|
|
|
|
|
+ std::vector<IndexValue> out;
|
|
|
|
|
+ if (limit == 0) return out;
|
|
|
|
|
+
|
|
|
|
|
+ ReadTxn rtxn(env_);
|
|
|
|
|
+ auto idbi = try_open_for_read(rtxn, index_subdb_name(collection, field));
|
|
|
|
|
+ if (!idbi) return out;
|
|
|
|
|
+ auto cdbi = try_open_for_read(rtxn, collection);
|
|
|
|
|
+ if (!cdbi) return out;
|
|
|
|
|
+
|
|
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
|
|
+ mdb_check(mdb_cursor_open(rtxn.raw(), *idbi, &cur), "cursor_open (index values)");
|
|
|
|
|
+ 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, ascending ? MDB_FIRST : MDB_LAST);
|
|
|
|
|
+ while (rc == MDB_SUCCESS && out.size() < limit) {
|
|
|
|
|
+ if (!is_index_meta_key(to_sv(k))) {
|
|
|
|
|
+ size_t n = 0;
|
|
|
|
|
+ mdb_check(mdb_cursor_count(cur, &n), "cursor_count (index values)");
|
|
|
|
|
+
|
|
|
|
|
+ // Read the value from a row that holds it. Walking backwards lands on
|
|
|
|
|
+ // a key's LAST duplicate, which is just as good - any holder has the
|
|
|
|
|
+ // same value, that being what the key means.
|
|
|
|
|
+ MDB_val hk = to_val(to_sv(v));
|
|
|
|
|
+ MDB_val hv{0, nullptr};
|
|
|
|
|
+ if (mdb_get(rtxn.raw(), *cdbi, &hk, &hv) == MDB_SUCCESS) {
|
|
|
|
|
+ const auto bytes = to_sv(hv);
|
|
|
|
|
+ yyjson_doc* d = yyjson_read(bytes.data(), bytes.size(), 0);
|
|
|
|
|
+ if (d) {
|
|
|
|
|
+ auto val = resolve_from_yyjson(yyjson_doc_get_root(d), field);
|
|
|
|
|
+ yyjson_doc_free(d);
|
|
|
|
|
+ if (val) {
|
|
|
|
|
+ IndexValue iv;
|
|
|
|
|
+ // An array field contributes one posting per element, so
|
|
|
|
|
+ // the row's value is the whole array while this key names
|
|
|
|
|
+ // one element. Reporting the array would be wrong, so a
|
|
|
|
|
+ // multivalued index reports nothing rather than something
|
|
|
|
|
+ // misleading.
|
|
|
|
|
+ if (!val->is_array()) {
|
|
|
|
|
+ iv.value = std::move(*val);
|
|
|
|
|
+ iv.count = static_cast<uint64_t>(n);
|
|
|
|
|
+ out.push_back(std::move(iv));
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, ascending ? MDB_NEXT_NODUP : MDB_PREV_NODUP);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) {
|
|
|
|
|
+ throw_mdb(rc, "cursor step (index values)");
|
|
|
|
|
+ }
|
|
|
|
|
+ return out;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
std::optional<LmdbDocumentStore::IndexStats>
|
|
std::optional<LmdbDocumentStore::IndexStats>
|
|
@@ -715,9 +819,6 @@ LmdbDocumentStore::index_stats(std::string_view collection,
|
|
|
// Reporting it would overstate every index by exactly one, the same
|
|
// Reporting it would overstate every index by exactly one, the same
|
|
|
// off-by-one count() already subtracts for document sub-dbs.
|
|
// off-by-one count() already subtracts for document sub-dbs.
|
|
|
out.entries = st.ms_entries;
|
|
out.entries = st.ms_entries;
|
|
|
- if (!read_subdb_identity(rtxn, *dbi_opt).empty() && out.entries > 0) {
|
|
|
|
|
- --out.entries;
|
|
|
|
|
- }
|
|
|
|
|
|
|
|
|
|
// Distinct keys need a walk: MDB_NEXT_NODUP skips a key's remaining
|
|
// Distinct keys need a walk: MDB_NEXT_NODUP skips a key's remaining
|
|
|
// duplicates, so this costs one step per distinct value, not per posting.
|
|
// duplicates, so this costs one step per distinct value, not per posting.
|
|
@@ -726,12 +827,21 @@ LmdbDocumentStore::index_stats(std::string_view collection,
|
|
|
struct G { MDB_cursor* c; ~G() { if (c) mdb_cursor_close(c); } } g{cur};
|
|
struct G { MDB_cursor* c; ~G() { if (c) mdb_cursor_close(c); } } g{cur};
|
|
|
MDB_val k{0, nullptr};
|
|
MDB_val k{0, nullptr};
|
|
|
MDB_val v{0, nullptr};
|
|
MDB_val v{0, nullptr};
|
|
|
|
|
+ // Subtract the reserved keys from `entries` by counting them as we walk,
|
|
|
|
|
+ // rather than probing for each: they are the only NUL-prefixed keys and sort
|
|
|
|
|
+ // before every encoded value, so this sees them first.
|
|
|
|
|
+ uint64_t meta = 0;
|
|
|
int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
while (rc == MDB_SUCCESS) {
|
|
while (rc == MDB_SUCCESS) {
|
|
|
- if (!is_identity_key(to_sv(k))) ++out.distinct_values;
|
|
|
|
|
|
|
+ if (is_index_meta_key(to_sv(k))) {
|
|
|
|
|
+ ++meta;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ ++out.distinct_values;
|
|
|
|
|
+ }
|
|
|
rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_NODUP);
|
|
rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_NODUP);
|
|
|
}
|
|
}
|
|
|
if (rc != MDB_NOTFOUND) throw_mdb(rc, "cursor next_nodup (index stat)");
|
|
if (rc != MDB_NOTFOUND) throw_mdb(rc, "cursor next_nodup (index stat)");
|
|
|
|
|
+ out.entries = out.entries >= meta ? out.entries - meta : 0;
|
|
|
return out;
|
|
return out;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -764,7 +874,7 @@ LmdbDocumentStore::lookupIndexRangeTxn(ReadTxn& rtxn,
|
|
|
// ordered, so one tag at both ends means one tag throughout. The identity
|
|
// 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.
|
|
// sentinel sorts before every tag (it starts with NUL), so skip it.
|
|
|
int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
- while (rc == MDB_SUCCESS && is_identity_key(to_sv(k))) {
|
|
|
|
|
|
|
+ while (rc == MDB_SUCCESS && is_index_meta_key(to_sv(k))) {
|
|
|
rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_NODUP);
|
|
rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_NODUP);
|
|
|
}
|
|
}
|
|
|
if (rc == MDB_NOTFOUND) return std::vector<std::string>{}; // empty index
|
|
if (rc == MDB_NOTFOUND) return std::vector<std::string>{}; // empty index
|
|
@@ -796,7 +906,7 @@ LmdbDocumentStore::lookupIndexRangeTxn(ReadTxn& rtxn,
|
|
|
|
|
|
|
|
while (rc == MDB_SUCCESS) {
|
|
while (rc == MDB_SUCCESS) {
|
|
|
const auto key = to_sv(k);
|
|
const auto key = to_sv(k);
|
|
|
- if (!is_identity_key(key)) {
|
|
|
|
|
|
|
+ if (!is_index_meta_key(key)) {
|
|
|
if (index_key_tag(key) != first_tag) break; // left the type
|
|
if (index_key_tag(key) != first_tag) break; // left the type
|
|
|
if (!ascending) {
|
|
if (!ascending) {
|
|
|
const int cmp = key.compare(bsv);
|
|
const int cmp = key.compare(bsv);
|
|
@@ -851,6 +961,7 @@ uint64_t LmdbDocumentStore::build_index(std::string_view collection,
|
|
|
if (!val) continue; // field absent on this row
|
|
if (!val) continue; // field absent on this row
|
|
|
const auto keys = encode_index_keys(*val);
|
|
const auto keys = encode_index_keys(*val);
|
|
|
if (keys.empty()) continue; // unindexable value
|
|
if (keys.empty()) continue; // unindexable value
|
|
|
|
|
+ if (val->is_array()) markIndexMultiValued(wtxn, dbi);
|
|
|
for (const auto& key : keys) {
|
|
for (const auto& key : keys) {
|
|
|
MDB_val kk = to_val(key);
|
|
MDB_val kk = to_val(key);
|
|
|
MDB_val vv = to_val(id);
|
|
MDB_val vv = to_val(id);
|
|
@@ -1001,6 +1112,245 @@ uint64_t LmdbDocumentStore::rowCountTxn(ReadTxn& rtxn, unsigned int dbi) {
|
|
|
return n;
|
|
return n;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+void LmdbDocumentStore::set_unique_fields(std::string_view collection,
|
|
|
|
|
+ std::vector<std::string> fields) {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(index_mutex_);
|
|
|
|
|
+ if (fields.empty()) {
|
|
|
|
|
+ unique_fields_.erase(std::string(collection));
|
|
|
|
|
+ } else {
|
|
|
|
|
+ unique_fields_[std::string(collection)] = std::move(fields);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+std::vector<std::string>
|
|
|
|
|
+LmdbDocumentStore::unique_fields(std::string_view collection) {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(index_mutex_);
|
|
|
|
|
+ auto it = unique_fields_.find(std::string(collection));
|
|
|
|
|
+ if (it == unique_fields_.end()) return {};
|
|
|
|
|
+ return it->second;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+std::vector<LmdbDocumentStore::DuplicateValue>
|
|
|
|
|
+LmdbDocumentStore::find_duplicate_values(std::string_view collection,
|
|
|
|
|
+ const std::string& field,
|
|
|
|
|
+ size_t limit) {
|
|
|
|
|
+ std::vector<DuplicateValue> out;
|
|
|
|
|
+ ReadTxn rtxn(env_);
|
|
|
|
|
+ auto dbi_opt = try_open_for_read(rtxn, index_subdb_name(collection, field));
|
|
|
|
|
+ if (!dbi_opt) return out;
|
|
|
|
|
+
|
|
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
|
|
+ mdb_check(mdb_cursor_open(rtxn.raw(), *dbi_opt, &cur), "cursor_open (dupes)");
|
|
|
|
|
+ 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 && out.size() < limit) {
|
|
|
|
|
+ if (!is_index_meta_key(to_sv(k))) {
|
|
|
|
|
+ size_t n = 0;
|
|
|
|
|
+ mdb_check(mdb_cursor_count(cur, &n), "cursor_count (dupes)");
|
|
|
|
|
+ if (n > 1) {
|
|
|
|
|
+ DuplicateValue d;
|
|
|
|
|
+ d.count = static_cast<uint64_t>(n);
|
|
|
|
|
+ d.sample_id = std::string(to_sv(v));
|
|
|
|
|
+ out.push_back(std::move(d));
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_NODUP);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) throw_mdb(rc, "cursor (dupes)");
|
|
|
|
|
+ return out;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void LmdbDocumentStore::markIndexMultiValued(WriteTxn& wtxn, unsigned int dbi) {
|
|
|
|
|
+ MDB_val k{kIndexMultiValuedKey.size(),
|
|
|
|
|
+ const_cast<char*>(kIndexMultiValuedKey.data())};
|
|
|
|
|
+ MDB_val v{1, const_cast<char*>("1")};
|
|
|
|
|
+ // MDB_NODUPDATA so re-marking is a no-op rather than an error.
|
|
|
|
|
+ const int rc = mdb_put(wtxn.raw(), dbi, &k, &v, MDB_NODUPDATA);
|
|
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_KEYEXIST) throw_mdb(rc, "mark multivalued");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+bool LmdbDocumentStore::indexIsMultiValuedTxn(ReadTxn& rtxn,
|
|
|
|
|
+ std::string_view collection,
|
|
|
|
|
+ const std::string& field) {
|
|
|
|
|
+ auto dbi_opt = try_open_for_read(rtxn, index_subdb_name(collection, field));
|
|
|
|
|
+ if (!dbi_opt) return true; // unknown: assume the worst
|
|
|
|
|
+ MDB_val k{kIndexMultiValuedKey.size(),
|
|
|
|
|
+ const_cast<char*>(kIndexMultiValuedKey.data())};
|
|
|
|
|
+ MDB_val v{0, nullptr};
|
|
|
|
|
+ return mdb_get(rtxn.raw(), *dbi_opt, &k, &v) == MDB_SUCCESS;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+bool LmdbDocumentStore::rangeOrderedFromIndex(
|
|
|
|
|
+ ReadTxn& rtxn,
|
|
|
|
|
+ std::string_view collection,
|
|
|
|
|
+ const smartbotic::database::Query& query,
|
|
|
|
|
+ unsigned int coll_dbi,
|
|
|
|
|
+ ScanResult& result) {
|
|
|
|
|
+
|
|
|
|
|
+ using Op = smartbotic::database::FilterOp;
|
|
|
|
|
+ if (query.filters.size() != 1) return false;
|
|
|
|
|
+ if (!query.sort || query.sort->field.empty()) return false;
|
|
|
|
|
+
|
|
|
|
|
+ const auto& f = query.filters.front();
|
|
|
|
|
+ if (f.op != Op::GT && f.op != Op::GTE && f.op != Op::LT && f.op != Op::LTE) {
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (f.field != query.sort->field) return false; // the walk's order IS the sort
|
|
|
|
|
+
|
|
|
|
|
+ const auto fields = indexed_fields(collection);
|
|
|
|
|
+ if (std::find(fields.begin(), fields.end(), f.field) == fields.end()) return false;
|
|
|
|
|
+ // One posting per row, or the count is not the match count.
|
|
|
|
|
+ if (indexIsMultiValuedTxn(rtxn, collection, f.field)) return false;
|
|
|
|
|
+
|
|
|
|
|
+ // No budget: the walk collects ids only, and the total needs all of them.
|
|
|
|
|
+ // Passing UINT64_MAX asks for the whole range rather than declining on size.
|
|
|
|
|
+ auto ids = lookupIndexRangeTxn(rtxn, collection, f.field, f.op, f.value,
|
|
|
|
|
+ std::numeric_limits<uint64_t>::max());
|
|
|
|
|
+ if (!ids) return false; // no index, unindexable bound, or mixed types
|
|
|
|
|
+
|
|
|
|
|
+ // lookupIndexRangeTxn walks ascending. Reversing gives descending keys AND
|
|
|
|
|
+ // descending ids within a key, which is exactly how sort_documents breaks a
|
|
|
|
|
+ // tie when descending.
|
|
|
|
|
+ if (query.sort->descending) std::reverse(ids->begin(), ids->end());
|
|
|
|
|
+
|
|
|
|
|
+ const uint64_t total = ids->size();
|
|
|
|
|
+ const uint64_t start = std::min<uint64_t>(query.offset, total);
|
|
|
|
|
+ const uint64_t end = std::min<uint64_t>(
|
|
|
|
|
+ static_cast<uint64_t>(query.offset) + query.limit, total);
|
|
|
|
|
+
|
|
|
|
|
+ result.documents.clear();
|
|
|
|
|
+ result.documents.reserve(end > start ? end - start : 0);
|
|
|
|
|
+ for (uint64_t i = start; i < end; ++i) {
|
|
|
|
|
+ MDB_val hk = to_val((*ids)[i]);
|
|
|
|
|
+ MDB_val hv{0, nullptr};
|
|
|
|
|
+ if (mdb_get(rtxn.raw(), coll_dbi, &hk, &hv) != MDB_SUCCESS) continue;
|
|
|
|
|
+ result.documents.push_back(decode_document(to_sv(hv)));
|
|
|
|
|
+ }
|
|
|
|
|
+ result.total_matched = total;
|
|
|
|
|
+ result.has_more = end < total;
|
|
|
|
|
+
|
|
|
|
|
+ if (!query.projection.empty()) {
|
|
|
|
|
+ for (auto& d : result.documents) {
|
|
|
|
|
+ nlohmann::json data = d.data();
|
|
|
|
|
+ apply_projection_inplace(data, query.projection);
|
|
|
|
|
+ d.set_data(data);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ range_ordered_scans_.fetch_add(1, std::memory_order_relaxed);
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+bool LmdbDocumentStore::countAndPageFromIndex(
|
|
|
|
|
+ ReadTxn& rtxn,
|
|
|
|
|
+ std::string_view collection,
|
|
|
|
|
+ const smartbotic::database::Query& query,
|
|
|
|
|
+ unsigned int coll_dbi,
|
|
|
|
|
+ ScanResult& result) {
|
|
|
|
|
+
|
|
|
|
|
+ using Op = smartbotic::database::FilterOp;
|
|
|
|
|
+ // Exactly one predicate: with two, the total is the size of an intersection
|
|
|
|
|
+ // that no single index knows.
|
|
|
|
|
+ if (query.filters.size() != 1) return false;
|
|
|
|
|
+ if (query.sort && !query.sort->field.empty()) return false; // order differs
|
|
|
|
|
+
|
|
|
|
|
+ const auto& f = query.filters.front();
|
|
|
|
|
+ if (f.op != Op::EQ && f.op != Op::IN) return false;
|
|
|
|
|
+
|
|
|
|
|
+ const auto fields = indexed_fields(collection);
|
|
|
|
|
+ if (std::find(fields.begin(), fields.end(), f.field) == fields.end()) return false;
|
|
|
|
|
+
|
|
|
|
|
+ // The count is only the match count while every row owns at most one posting.
|
|
|
|
|
+ // Once an array has been indexed, postings under a key can name rows whose
|
|
|
|
|
+ // value merely CONTAINS it, and EQ compares the whole value.
|
|
|
|
|
+ if (indexIsMultiValuedTxn(rtxn, collection, f.field)) return false;
|
|
|
|
|
+
|
|
|
|
|
+ const uint64_t want = static_cast<uint64_t>(query.offset) + query.limit;
|
|
|
|
|
+ std::vector<std::string> page_ids;
|
|
|
|
|
+ uint64_t total = 0;
|
|
|
|
|
+
|
|
|
|
|
+ if (f.op == Op::EQ) {
|
|
|
|
|
+ auto n = countIndexEqTxn(rtxn, collection, f.field, f.value);
|
|
|
|
|
+ if (!n) return false; // no index, or unindexable value
|
|
|
|
|
+ total = *n;
|
|
|
|
|
+ // Walk only as far as the page needs. The ids under one key are ascending,
|
|
|
|
|
+ // which is the order a collection cursor would have produced.
|
|
|
|
|
+ if (total > 0 && want > 0) {
|
|
|
|
|
+ auto dbi_opt = try_open_for_read(rtxn,
|
|
|
|
|
+ index_subdb_name(collection, f.field));
|
|
|
|
|
+ if (!dbi_opt) return false;
|
|
|
|
|
+ const auto key = encode_index_key(f.value);
|
|
|
|
|
+ if (!key) return false;
|
|
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
|
|
+ mdb_check(mdb_cursor_open(rtxn.raw(), *dbi_opt, &cur),
|
|
|
|
|
+ "cursor_open (index page)");
|
|
|
|
|
+ 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};
|
|
|
|
|
+ int rc = mdb_cursor_get(cur, &k, &v, MDB_SET);
|
|
|
|
|
+ if (rc == MDB_SUCCESS) rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST_DUP);
|
|
|
|
|
+ uint64_t seen = 0;
|
|
|
|
|
+ while (rc == MDB_SUCCESS && page_ids.size() < query.limit) {
|
|
|
|
|
+ if (seen >= query.offset) page_ids.emplace_back(to_sv(v));
|
|
|
|
|
+ ++seen;
|
|
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT_DUP);
|
|
|
|
|
+ }
|
|
|
|
|
+ if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) {
|
|
|
|
|
+ throw_mdb(rc, "cursor next_dup (index page)");
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ } else {
|
|
|
|
|
+ // IN. A non-multivalued field gives each row one value, so a row appears
|
|
|
|
|
+ // under exactly one of the requested keys - the per-value counts are
|
|
|
|
|
+ // disjoint and their sum is the exact total.
|
|
|
|
|
+ if (!f.value.is_array()) return false;
|
|
|
|
|
+ std::vector<std::string> all;
|
|
|
|
|
+ for (const auto& el : f.value) {
|
|
|
|
|
+ auto n = countIndexEqTxn(rtxn, collection, f.field, el);
|
|
|
|
|
+ if (!n) return false;
|
|
|
|
|
+ total += *n;
|
|
|
|
|
+ if (want == 0) continue;
|
|
|
|
|
+ auto part = lookupIndexEqTxn(rtxn, collection, f.field, el);
|
|
|
|
|
+ if (!part) return false;
|
|
|
|
|
+ all.insert(all.end(), part->begin(), part->end());
|
|
|
|
|
+ }
|
|
|
|
|
+ // Merging is needed because the page must be in id order across values.
|
|
|
|
|
+ // Only ids are read, never documents, so this stays cheap.
|
|
|
|
|
+ std::sort(all.begin(), all.end());
|
|
|
|
|
+ const uint64_t start = std::min<uint64_t>(query.offset, all.size());
|
|
|
|
|
+ const uint64_t end = std::min<uint64_t>(want, all.size());
|
|
|
|
|
+ page_ids.assign(all.begin() + start, all.begin() + end);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Materialise ONLY the page.
|
|
|
|
|
+ result.documents.clear();
|
|
|
|
|
+ result.documents.reserve(page_ids.size());
|
|
|
|
|
+ for (const auto& id : page_ids) {
|
|
|
|
|
+ MDB_val hk = to_val(id);
|
|
|
|
|
+ MDB_val hv{0, nullptr};
|
|
|
|
|
+ if (mdb_get(rtxn.raw(), coll_dbi, &hk, &hv) != MDB_SUCCESS) continue;
|
|
|
|
|
+ result.documents.push_back(decode_document(to_sv(hv)));
|
|
|
|
|
+ }
|
|
|
|
|
+ result.total_matched = total;
|
|
|
|
|
+ result.has_more = std::min<uint64_t>(want, total) < total;
|
|
|
|
|
+
|
|
|
|
|
+ if (!query.projection.empty()) {
|
|
|
|
|
+ for (auto& d : result.documents) {
|
|
|
|
|
+ nlohmann::json data = d.data();
|
|
|
|
|
+ apply_projection_inplace(data, query.projection);
|
|
|
|
|
+ d.set_data(data);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ counted_scans_.fetch_add(1, std::memory_order_relaxed);
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
bool LmdbDocumentStore::orderedHitsFromIndex(
|
|
bool LmdbDocumentStore::orderedHitsFromIndex(
|
|
|
ReadTxn& rtxn,
|
|
ReadTxn& rtxn,
|
|
|
std::string_view collection,
|
|
std::string_view collection,
|
|
@@ -1026,11 +1376,8 @@ bool LmdbDocumentStore::orderedHitsFromIndex(
|
|
|
// Total order or nothing. entries != rows means some row has no posting (the
|
|
// Total order or nothing. entries != rows means some row has no posting (the
|
|
|
// field is absent, or its value is unindexable) or several (an array), and in
|
|
// field is absent, or its value is unindexable) or several (an array), and in
|
|
|
// either case the index cannot state where that row sorts.
|
|
// either case the index cannot state where that row sorts.
|
|
|
- MDB_stat ist{};
|
|
|
|
|
- if (mdb_stat(rtxn.raw(), *dbi_opt, &ist) != MDB_SUCCESS) return false;
|
|
|
|
|
- uint64_t entries = ist.ms_entries;
|
|
|
|
|
- if (!read_subdb_identity(rtxn, *dbi_opt).empty() && entries > 0) --entries;
|
|
|
|
|
- if (entries != rows) return false;
|
|
|
|
|
|
|
+ const auto istats = index_stats(collection, query.sort->field);
|
|
|
|
|
+ if (!istats || istats->entries != rows) return false;
|
|
|
|
|
|
|
|
MDB_cursor* cur = nullptr;
|
|
MDB_cursor* cur = nullptr;
|
|
|
mdb_check(mdb_cursor_open(rtxn.raw(), *dbi_opt, &cur), "cursor_open (index order)");
|
|
mdb_check(mdb_cursor_open(rtxn.raw(), *dbi_opt, &cur), "cursor_open (index order)");
|
|
@@ -1049,7 +1396,7 @@ bool LmdbDocumentStore::orderedHitsFromIndex(
|
|
|
int rc = mdb_cursor_get(cur, &k, &v, desc ? MDB_LAST : MDB_FIRST);
|
|
int rc = mdb_cursor_get(cur, &k, &v, desc ? MDB_LAST : MDB_FIRST);
|
|
|
std::vector<std::string> ids;
|
|
std::vector<std::string> ids;
|
|
|
while (rc == MDB_SUCCESS && ids.size() < want) {
|
|
while (rc == MDB_SUCCESS && ids.size() < want) {
|
|
|
- if (!is_identity_key(to_sv(k))) ids.emplace_back(to_sv(v));
|
|
|
|
|
|
|
+ if (!is_index_meta_key(to_sv(k))) ids.emplace_back(to_sv(v));
|
|
|
rc = mdb_cursor_get(cur, &k, &v, desc ? MDB_PREV : MDB_NEXT);
|
|
rc = mdb_cursor_get(cur, &k, &v, desc ? MDB_PREV : MDB_NEXT);
|
|
|
}
|
|
}
|
|
|
if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) {
|
|
if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) {
|
|
@@ -1102,7 +1449,7 @@ LmdbDocumentStore::lookupIndexAllTxn(ReadTxn& rtxn,
|
|
|
MDB_val v{0, nullptr};
|
|
MDB_val v{0, nullptr};
|
|
|
int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
while (rc == MDB_SUCCESS) {
|
|
while (rc == MDB_SUCCESS) {
|
|
|
- if (!is_identity_key(to_sv(k))) {
|
|
|
|
|
|
|
+ if (!is_index_meta_key(to_sv(k))) {
|
|
|
ids.emplace_back(to_sv(v));
|
|
ids.emplace_back(to_sv(v));
|
|
|
// Bail out rather than truncate, as everywhere else: a partial
|
|
// Bail out rather than truncate, as everywhere else: a partial
|
|
|
// candidate list drops matching rows.
|
|
// candidate list drops matching rows.
|
|
@@ -1401,6 +1748,18 @@ ScanResult LmdbDocumentStore::scan(std::string_view collection,
|
|
|
yyjson_doc_free(ydoc);
|
|
yyjson_doc_free(ydoc);
|
|
|
};
|
|
};
|
|
|
|
|
|
|
|
|
|
+ // v2.10.0 — count AND page from the index. Cheapest of all when it
|
|
|
|
|
+ // applies: nothing but the page is read, and no selectivity guard is
|
|
|
|
|
+ // needed because the total never requires visiting a match.
|
|
|
|
|
+ if (countAndPageFromIndex(rtxn, collection, query, *dbi_opt, result)) {
|
|
|
|
|
+ return result;
|
|
|
|
|
+ }
|
|
|
|
|
+ // A range on the field being sorted: the walk's order already IS the sort
|
|
|
|
|
+ // order, so no sorting pass and only the page is decoded.
|
|
|
|
|
+ if (rangeOrderedFromIndex(rtxn, collection, query, *dbi_opt, result)) {
|
|
|
|
|
+ return result;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// v2.9.2 — ORDERING from the index. Try this first: when it applies it
|
|
// v2.9.2 — ORDERING from the index. Try this first: when it applies it
|
|
|
// reads the page and nothing else, where every other plan still visits
|
|
// reads the page and nothing else, where every other plan still visits
|
|
|
// every matching row.
|
|
// every matching row.
|