|
|
@@ -0,0 +1,340 @@
|
|
|
+// v2.4.4 — offline sub-db placement audit and repair. See subdb_placement.hpp.
|
|
|
+
|
|
|
+#include "storage/subdb_placement.hpp"
|
|
|
+
|
|
|
+#include <lmdb.h>
|
|
|
+
|
|
|
+#include <algorithm>
|
|
|
+#include <cstring>
|
|
|
+#include <filesystem>
|
|
|
+#include <iostream>
|
|
|
+#include <map>
|
|
|
+#include <set>
|
|
|
+#include <stdexcept>
|
|
|
+
|
|
|
+#include <nlohmann/json.hpp>
|
|
|
+
|
|
|
+#include "storage/subdb_identity.hpp"
|
|
|
+
|
|
|
+namespace smartbotic::db::storage {
|
|
|
+
|
|
|
+namespace {
|
|
|
+
|
|
|
+namespace fs = std::filesystem;
|
|
|
+using nlohmann::json;
|
|
|
+
|
|
|
+// The sentinel key itself comes from storage/subdb_identity.hpp — constexpr,
|
|
|
+// so including it costs no link dependency. This file deliberately uses raw
|
|
|
+// mdb_* calls rather than the WriteTxn-based helpers there, because it also
|
|
|
+// builds into the CLI, which does not link the service storage stack.
|
|
|
+constexpr std::string_view kIdentityKey = kSubdbIdentityKey;
|
|
|
+
|
|
|
+std::string_view to_sv(const MDB_val& v) {
|
|
|
+ return std::string_view(static_cast<const char*>(v.mv_data), v.mv_size);
|
|
|
+}
|
|
|
+
|
|
|
+MDB_val to_val(std::string_view sv) {
|
|
|
+ return MDB_val{sv.size(), const_cast<char*>(sv.data())};
|
|
|
+}
|
|
|
+
|
|
|
+void ck(int rc, const char* where) {
|
|
|
+ if (rc != MDB_SUCCESS) {
|
|
|
+ throw std::runtime_error(std::string("LMDB ") + where + ": " + mdb_strerror(rc));
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+// A sub-db is "system" if it is internal bookkeeping rather than a user
|
|
|
+// collection: _meta, _vectors_*, _history_*, _files, _views, _orphans_*.
|
|
|
+bool is_system_subdb(std::string_view name) {
|
|
|
+ return !name.empty() && name.front() == '_';
|
|
|
+}
|
|
|
+
|
|
|
+// Infer the project from .../projects/<name>/env.
|
|
|
+std::string infer_project(const std::string& env_path) {
|
|
|
+ fs::path p(env_path);
|
|
|
+ if (p.filename() == "env" && p.has_parent_path()) {
|
|
|
+ return p.parent_path().filename().string();
|
|
|
+ }
|
|
|
+ return "default";
|
|
|
+}
|
|
|
+
|
|
|
+// Strip a leading "<project>:" if present.
|
|
|
+std::string dequalify(const std::string& name, const std::string& project) {
|
|
|
+ const std::string prefix = project + ":";
|
|
|
+ if (name.rfind(prefix, 0) == 0) return name.substr(prefix.size());
|
|
|
+ return name;
|
|
|
+}
|
|
|
+
|
|
|
+// Does a row declaring `declared` belong in the sub-db physically named
|
|
|
+// `physical`? Both the qualified and bare spellings are accepted, because the
|
|
|
+// envs in the field contain a mix of the two (the `default` project has
|
|
|
+// sub-dbs named both "conversations" and "default:mypure_test_col").
|
|
|
+bool declaration_matches(const std::string& declared,
|
|
|
+ const std::string& physical,
|
|
|
+ const std::string& project) {
|
|
|
+ if (declared == physical) return true;
|
|
|
+ if (dequalify(declared, project) == physical) return true;
|
|
|
+ if (declared == dequalify(physical, project)) return true;
|
|
|
+ return false;
|
|
|
+}
|
|
|
+
|
|
|
+std::vector<std::string> list_subdbs(MDB_txn* txn) {
|
|
|
+ std::vector<std::string> names;
|
|
|
+ MDB_dbi root = 0;
|
|
|
+ int rc = mdb_dbi_open(txn, nullptr, 0, &root);
|
|
|
+ if (rc == MDB_NOTFOUND) return names;
|
|
|
+ ck(rc, "dbi_open (root)");
|
|
|
+
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
+ ck(mdb_cursor_open(txn, root, &cur), "cursor_open (root)");
|
|
|
+ MDB_val k{}, v{};
|
|
|
+ int crc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
+ while (crc == MDB_SUCCESS) {
|
|
|
+ names.emplace_back(to_sv(k));
|
|
|
+ crc = mdb_cursor_get(cur, &k, &v, MDB_NEXT);
|
|
|
+ }
|
|
|
+ mdb_cursor_close(cur);
|
|
|
+ std::sort(names.begin(), names.end());
|
|
|
+ return names;
|
|
|
+}
|
|
|
+
|
|
|
+// Pull the "collection" field out of a stored document without building the
|
|
|
+// whole tree twice. Returns empty on parse failure or missing field.
|
|
|
+std::string declared_collection_of(std::string_view payload) {
|
|
|
+ try {
|
|
|
+ json j = json::parse(payload);
|
|
|
+ if (!j.is_object()) return {};
|
|
|
+ auto it = j.find("collection");
|
|
|
+ if (it == j.end() || !it->is_string()) return {};
|
|
|
+ return it->get<std::string>();
|
|
|
+ } catch (const std::exception&) {
|
|
|
+ return {};
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+struct EnvHandle {
|
|
|
+ MDB_env* env = nullptr;
|
|
|
+ explicit EnvHandle(const std::string& path, bool readonly) {
|
|
|
+ ck(mdb_env_create(&env), "env_create");
|
|
|
+ ck(mdb_env_set_maxdbs(env, 512), "set_maxdbs");
|
|
|
+ // Match the service's map size ceiling; repair may add sub-dbs.
|
|
|
+ ck(mdb_env_set_mapsize(env, 2ULL << 30), "set_mapsize");
|
|
|
+ unsigned int flags = readonly ? MDB_RDONLY : 0;
|
|
|
+ int rc = mdb_env_open(env, path.c_str(), flags, 0664);
|
|
|
+ if (rc != MDB_SUCCESS) {
|
|
|
+ mdb_env_close(env);
|
|
|
+ env = nullptr;
|
|
|
+ ck(rc, "env_open");
|
|
|
+ }
|
|
|
+ }
|
|
|
+ ~EnvHandle() { if (env) mdb_env_close(env); }
|
|
|
+ EnvHandle(const EnvHandle&) = delete;
|
|
|
+ EnvHandle& operator=(const EnvHandle&) = delete;
|
|
|
+};
|
|
|
+
|
|
|
+} // namespace
|
|
|
+
|
|
|
+AuditReport audit(const std::string& env_path, const std::string& project_in) {
|
|
|
+ AuditReport rep;
|
|
|
+ rep.project = project_in.empty() ? infer_project(env_path) : project_in;
|
|
|
+
|
|
|
+ EnvHandle eh(env_path, /*readonly=*/true);
|
|
|
+ MDB_txn* txn = nullptr;
|
|
|
+ ck(mdb_txn_begin(eh.env, nullptr, MDB_RDONLY, &txn), "txn_begin (audit)");
|
|
|
+
|
|
|
+ const auto all = list_subdbs(txn);
|
|
|
+
|
|
|
+ // First pass: which ids exist in which sub-db, so we can tell a move from
|
|
|
+ // a quarantine without a second env open.
|
|
|
+ std::map<std::string, std::set<std::string>> ids_by_subdb;
|
|
|
+
|
|
|
+ for (const auto& name : all) {
|
|
|
+ if (is_system_subdb(name)) continue;
|
|
|
+ rep.subdbs.push_back(name);
|
|
|
+
|
|
|
+ MDB_dbi dbi = 0;
|
|
|
+ if (mdb_dbi_open(txn, name.c_str(), 0, &dbi) != MDB_SUCCESS) continue;
|
|
|
+
|
|
|
+ // Identity sentinel present?
|
|
|
+ {
|
|
|
+ MDB_val k = to_val(kIdentityKey);
|
|
|
+ MDB_val v{};
|
|
|
+ if (mdb_get(txn, dbi, &k, &v) != MDB_SUCCESS) {
|
|
|
+ rep.unstamped.push_back(name);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ MDB_cursor* cur = nullptr;
|
|
|
+ ck(mdb_cursor_open(txn, dbi, &cur), "cursor_open (audit)");
|
|
|
+ MDB_val k{}, v{};
|
|
|
+ int rc = mdb_cursor_get(cur, &k, &v, MDB_FIRST);
|
|
|
+ while (rc == MDB_SUCCESS) {
|
|
|
+ std::string_view key = to_sv(k);
|
|
|
+ if (key != kIdentityKey) {
|
|
|
+ ids_by_subdb[name].insert(std::string(key));
|
|
|
+ ++rep.rows_scanned;
|
|
|
+ const std::string declared = declared_collection_of(to_sv(v));
|
|
|
+ if (declared.empty()) {
|
|
|
+ ++rep.rows_undecodable;
|
|
|
+ } else if (!declaration_matches(declared, name, rep.project)) {
|
|
|
+ MisplacedRow m;
|
|
|
+ m.physical_subdb = name;
|
|
|
+ m.declared_collection = declared;
|
|
|
+ m.id = std::string(key);
|
|
|
+ rep.misplaced.push_back(std::move(m));
|
|
|
+ }
|
|
|
+ }
|
|
|
+ rc = mdb_cursor_get(cur, &k, &v, MDB_NEXT);
|
|
|
+ }
|
|
|
+ mdb_cursor_close(cur);
|
|
|
+ }
|
|
|
+
|
|
|
+ // Resolve each misplaced row's destination. Prefer an existing sub-db
|
|
|
+ // spelled exactly as declared; else the de-qualified spelling, which is the
|
|
|
+ // dominant convention. Then note whether the home is already occupied.
|
|
|
+ const std::set<std::string> existing(all.begin(), all.end());
|
|
|
+ for (auto& m : rep.misplaced) {
|
|
|
+ const std::string bare = dequalify(m.declared_collection, rep.project);
|
|
|
+ if (existing.count(m.declared_collection)) {
|
|
|
+ m.home_subdb = m.declared_collection;
|
|
|
+ } else if (existing.count(bare)) {
|
|
|
+ m.home_subdb = bare;
|
|
|
+ } else {
|
|
|
+ m.home_subdb = bare; // will be created
|
|
|
+ }
|
|
|
+ auto it = ids_by_subdb.find(m.home_subdb);
|
|
|
+ m.home_occupied = (it != ids_by_subdb.end()) && it->second.count(m.id) > 0;
|
|
|
+ }
|
|
|
+
|
|
|
+ mdb_txn_abort(txn);
|
|
|
+ return rep;
|
|
|
+}
|
|
|
+
|
|
|
+RepairResult repair(const std::string& env_path,
|
|
|
+ const AuditReport& report,
|
|
|
+ bool stamp_identity) {
|
|
|
+ RepairResult res;
|
|
|
+ if (report.misplaced.empty() && !(stamp_identity && !report.unstamped.empty())) {
|
|
|
+ return res;
|
|
|
+ }
|
|
|
+
|
|
|
+ EnvHandle eh(env_path, /*readonly=*/false);
|
|
|
+ MDB_txn* txn = nullptr;
|
|
|
+ ck(mdb_txn_begin(eh.env, nullptr, 0, &txn), "txn_begin (repair)");
|
|
|
+
|
|
|
+ try {
|
|
|
+ // Cache handles for the duration of this single write txn. Safe here
|
|
|
+ // precisely because the txn commits at the end — the failure mode this
|
|
|
+ // whole tool exists to repair comes from caching across an ABORT.
|
|
|
+ std::map<std::string, MDB_dbi> handles;
|
|
|
+ auto open_db = [&](const std::string& name) -> MDB_dbi {
|
|
|
+ auto it = handles.find(name);
|
|
|
+ if (it != handles.end()) return it->second;
|
|
|
+ MDB_dbi dbi = 0;
|
|
|
+ ck(mdb_dbi_open(txn, name.c_str(), MDB_CREATE, &dbi), "dbi_open (repair)");
|
|
|
+ handles[name] = dbi;
|
|
|
+ return dbi;
|
|
|
+ };
|
|
|
+
|
|
|
+ for (const auto& m : report.misplaced) {
|
|
|
+ MDB_dbi src = open_db(m.physical_subdb);
|
|
|
+ MDB_val k = to_val(m.id);
|
|
|
+ MDB_val v{};
|
|
|
+ int rc = mdb_get(txn, src, &k, &v);
|
|
|
+ if (rc == MDB_NOTFOUND) {
|
|
|
+ res.errors.push_back("row vanished before repair: " +
|
|
|
+ m.physical_subdb + "/" + m.id);
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ ck(rc, "get (repair)");
|
|
|
+ // Copy out of the mmap before any write invalidates the page.
|
|
|
+ const std::string payload(static_cast<const char*>(v.mv_data), v.mv_size);
|
|
|
+
|
|
|
+ std::string dest;
|
|
|
+ if (m.home_occupied) {
|
|
|
+ // Never overwrite a document that is already correctly filed.
|
|
|
+ // Park the stray copy for manual comparison instead.
|
|
|
+ dest = "_orphans_" + m.physical_subdb;
|
|
|
+ ++res.quarantined;
|
|
|
+ } else {
|
|
|
+ dest = m.home_subdb;
|
|
|
+ ++res.moved;
|
|
|
+ }
|
|
|
+
|
|
|
+ MDB_dbi dst = open_db(dest);
|
|
|
+ MDB_val dk = to_val(m.id);
|
|
|
+ MDB_val dv = to_val(payload);
|
|
|
+ ck(mdb_put(txn, dst, &dk, &dv, 0), "put (repair)");
|
|
|
+
|
|
|
+ MDB_val delk = to_val(m.id);
|
|
|
+ int drc = mdb_del(txn, src, &delk, nullptr);
|
|
|
+ if (drc != MDB_SUCCESS && drc != MDB_NOTFOUND) {
|
|
|
+ ck(drc, "del (repair)");
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ // Identity sentinels — opt-in, because a pre-v2.4.4 server does not
|
|
|
+ // know to skip them and would choke on the extra key. See the header.
|
|
|
+ if (stamp_identity) {
|
|
|
+ std::set<std::string> to_stamp(report.unstamped.begin(),
|
|
|
+ report.unstamped.end());
|
|
|
+ for (const auto& m : report.misplaced) to_stamp.insert(m.home_subdb);
|
|
|
+ for (const auto& name : to_stamp) {
|
|
|
+ if (is_system_subdb(name)) continue;
|
|
|
+ MDB_dbi dbi = open_db(name);
|
|
|
+ MDB_val k = to_val(kIdentityKey);
|
|
|
+ MDB_val v = to_val(name);
|
|
|
+ ck(mdb_put(txn, dbi, &k, &v, 0), "put (identity stamp)");
|
|
|
+ ++res.stamped;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ ck(mdb_txn_commit(txn), "txn_commit (repair)");
|
|
|
+ } catch (...) {
|
|
|
+ mdb_txn_abort(txn);
|
|
|
+ throw;
|
|
|
+ }
|
|
|
+
|
|
|
+ return res;
|
|
|
+}
|
|
|
+
|
|
|
+void print_audit(const AuditReport& rep) {
|
|
|
+ std::cout << "project: " << rep.project << "\n"
|
|
|
+ << "sub-dbs: " << rep.subdbs.size() << "\n"
|
|
|
+ << "rows scanned: " << rep.rows_scanned << "\n";
|
|
|
+ if (rep.rows_undecodable) {
|
|
|
+ std::cout << "undecodable: " << rep.rows_undecodable
|
|
|
+ << " (no parseable `collection` field — not repairable here)\n";
|
|
|
+ }
|
|
|
+ std::cout << "unstamped: " << rep.unstamped.size()
|
|
|
+ << " (sub-dbs with no identity sentinel)\n";
|
|
|
+
|
|
|
+ if (rep.misplaced.empty()) {
|
|
|
+ std::cout << "\nNo misplaced rows. Placement is consistent.\n";
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ // Group for a readable summary rather than one line per row.
|
|
|
+ std::map<std::string, uint64_t> moves, quarantines;
|
|
|
+ for (const auto& m : rep.misplaced) {
|
|
|
+ const std::string key = m.physical_subdb + " -> " + m.home_subdb;
|
|
|
+ if (m.home_occupied) ++quarantines[key];
|
|
|
+ else ++moves[key];
|
|
|
+ }
|
|
|
+
|
|
|
+ std::cout << "\nMISPLACED ROWS: " << rep.misplaced.size() << "\n";
|
|
|
+ if (!moves.empty()) {
|
|
|
+ std::cout << "\n relocate to declared home:\n";
|
|
|
+ for (const auto& [k, n] : moves) {
|
|
|
+ std::cout << " " << n << "\t" << k << "\n";
|
|
|
+ }
|
|
|
+ }
|
|
|
+ if (!quarantines.empty()) {
|
|
|
+ std::cout << "\n home already occupied - park in _orphans_<subdb>:\n";
|
|
|
+ for (const auto& [k, n] : quarantines) {
|
|
|
+ std::cout << " " << n << "\t" << k << "\n";
|
|
|
+ }
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+} // namespace smartbotic::db::storage
|