| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931 |
- #include "database_service.hpp"
- #include "json_parse.hpp"
- #include "project_addressing.hpp"
- #include "storage/document_store.hpp"
- #include "storage/document_store_lmdb.hpp"
- #include "storage/dual_write_mirror.hpp"
- #include "relations/relation_cascade.hpp"
- #include "relations/relation_enforcement.hpp"
- #include "auth/auth_interceptor.hpp"
- #include "auth/principal.hpp"
- #include "tls/cert_generator.hpp"
- #include <fstream>
- #include <sstream>
- #include <map>
- #include <set>
- #include "storage/lmdb_env.hpp"
- #include "storage/migrate_v1_to_v2.hpp"
- #include "storage/project_store.hpp"
- #include "storage/subdb_placement.hpp"
- #include <grpcpp/grpcpp.h>
- #include <grpcpp/resource_quota.h>
- #include <spdlog/spdlog.h>
- #include <fstream>
- #ifdef HAVE_SYSTEMD
- #include <systemd/sd-daemon.h>
- #endif
- #if defined(__GLIBC__) && !defined(__APPLE__)
- #include <malloc.h>
- #endif
- namespace smartbotic::database {
- DatabaseService::DatabaseService(Config config)
- : config_(std::move(config))
- {
- }
- DatabaseService::~DatabaseService() {
- stop();
- }
- std::string DatabaseService::readOnlyReason() const {
- std::lock_guard<std::mutex> lock(reason_mutex_);
- return read_only_reason_;
- }
- void DatabaseService::setReadOnly(bool value, const std::string& reason) {
- {
- std::lock_guard<std::mutex> lock(reason_mutex_);
- read_only_reason_ = value ? reason : "";
- }
- bool prev = read_only_.exchange(value, std::memory_order_acq_rel);
- if (prev != value) {
- if (value) {
- spdlog::error("Database entered READ-ONLY mode: {}", reason);
- } else {
- spdlog::info("Database unlocked -- writes accepted");
- }
- }
- }
- bool DatabaseService::initialize() {
- spdlog::info("Initializing database service (node: {})", config_.nodeId);
- try {
- // Create data directory if needed
- std::error_code ec;
- std::filesystem::create_directories(config_.dataDirectory, ec);
- if (ec) {
- spdlog::error("Failed to create data directory: {}", ec.message());
- return false;
- }
- setupComponents();
- // Initialize encryption
- if (config_.encryptionEnabled) {
- if (!encryption_->initialize()) {
- spdlog::error("Failed to initialize encryption");
- return false;
- }
- }
- // Recover from persistence
- recovery_outcome_ = persistence_->recover(*store_);
- #if defined(__GLIBC__) && !defined(__APPLE__)
- // v1.9.3 — release freelist pages accumulated during recovery. v1.9.1
- // added trim at end of loadSnapshot, but the WAL replay phase that
- // runs afterward (`persistence_->recover`'s `replayWal` step) also
- // burns through GBs of small Document JSON allocations whose pages
- // sit on the per-thread freelist with no subsequent allocation to
- // shake them loose. On Zoe (docs/incidents/2026-04-22-zoe-rss-exceeds-budget.md update 15:05)
- // this was 5.7 GB at 23 min uptime — fully reclaimable via
- // malloc_trim, just nothing called it. Paired with the periodic
- // every-5-minute trim in MemoryStore::logMemoryCheck, so any
- // post-boot bloat that escapes this trim gets cleaned up shortly.
- ::malloc_trim(0);
- #endif
- if (recovery_outcome_.isFailure()) {
- const std::string modeStr =
- recoveryModeToString(config_.persistenceConfig.recoveryMode);
- spdlog::error("");
- spdlog::error("+------------------------------------------------------------------+");
- spdlog::error("| RECOVERY REFUSED |");
- spdlog::error("| |");
- spdlog::error("| Mode: {}", modeStr);
- spdlog::error("| Expected snapshot: {}", recovery_outcome_.expectedSnapshot.string());
- spdlog::error("| Reason: {}", recovery_outcome_.failureReason);
- spdlog::error("| |");
- spdlog::error("| Snapshots available: {}", recovery_outcome_.snapshotsAvailable);
- spdlog::error("| Snapshots attempted: {}", recovery_outcome_.snapshotsAttempted);
- spdlog::error("| |");
- spdlog::error("| To escalate, restart with one of: |");
- spdlog::error("| --recovery-mode=snapshot_fallback Try older snapshots |");
- spdlog::error("| --recovery-mode=wal_only Replay WAL only (slow) |");
- spdlog::error("| --recovery-mode=best_effort Try all of the above |");
- spdlog::error("| --recovery-mode=force_empty Start empty (LAST RESORT) |");
- spdlog::error("| |");
- spdlog::error("| Or set in config.json: \"recovery\": {{ \"mode\": \"<mode>\" }} |");
- spdlog::error("| |");
- spdlog::error("| PRESERVE /var/lib/smartbotic-database/ BEFORE ESCALATING. |");
- spdlog::error("+------------------------------------------------------------------+");
- throw std::runtime_error("recovery failed; refusing to start");
- }
- // Auto-readonly mode on non-trivial recovery
- if (recovery_outcome_.isNonTrivial() && !force_readwrite_) {
- std::string reason;
- switch (recovery_outcome_.kind) {
- case RecoveryOutcome::Kind::SnapshotFellBack:
- reason = "fell back to snapshot " +
- recovery_outcome_.snapshotUsed.filename().string() +
- " because " +
- recovery_outcome_.expectedSnapshot.filename().string() +
- " failed: " + recovery_outcome_.failureReason;
- break;
- case RecoveryOutcome::Kind::WalOnlyReplay:
- reason = "WAL-only replay, no snapshot loaded (" +
- std::to_string(recovery_outcome_.walEntriesReplayed) +
- " entries)";
- break;
- case RecoveryOutcome::Kind::ForcedEmpty:
- reason = "forced empty by operator (--recovery-mode=force_empty)";
- break;
- default:
- reason = "non-trivial recovery";
- break;
- }
- reason += ". Run `smartbotic-db-cli unlock` to accept this state, "
- "or restart with --force-readwrite to bypass this check.";
- setReadOnly(true, reason);
- spdlog::error("");
- spdlog::error("+------------------------------------------------------------------+");
- spdlog::error("| [ERROR] Database booted in READ-ONLY mode after non-trivial |");
- spdlog::error("| recovery. |");
- spdlog::error("| |");
- spdlog::error("| Reason: {}", reason);
- spdlog::error("| |");
- spdlog::error("| Writes will be REJECTED until you acknowledge this state: |");
- spdlog::error("| smartbotic-db-cli unlock # live, no restart |");
- spdlog::error("| smartbotic-database --force-readwrite # on next restart |");
- spdlog::error("+------------------------------------------------------------------+");
- }
- // Set replication sequence after recovery (WAL sequence is now known)
- replication_->setSequence(persistence_->currentWalSequence());
- // v1.8.0 — start the persistence manager before loadFromStore /
- // migrations so any system-collection writes those phases produce
- // (view docs, _collection_meta entries, _migrations recordings)
- // reach the WAL like normal mutations. Pre-v1.8, persistence_ was
- // started later in start(), which silently dropped those writes
- // (running_=false → logInsert no-op). That meant migrated views
- // only "survived" because the runner re-created them every boot,
- // and a manually-created view that landed during the racy startup
- // window (before full readiness) could be lost.
- if (!persistence_->start()) {
- spdlog::error("Failed to start persistence manager before migrations");
- return false;
- }
- // Load view definitions from the _views system collection (which is now
- // populated by the persistence recovery above).
- view_manager_->loadFromStore();
- // Load per-collection configs from the _collection_meta system collection.
- // Must happen AFTER persistence recovery and BEFORE migrations so any
- // migration-created documents are stamped with the correct precision.
- config_manager_->loadFromStore();
- policy_manager_->loadFromStore();
- relation_manager_->loadFromStore();
- applyRelationDeclarations();
- // v2.11.0 close-out — teach the TTL sweeper about relations. Installed
- // HERE, after loadFromStore() and the arming pass, deliberately: the
- // hook consults the declaration cache and the reverse index, and an
- // expiry that fired before either was ready would have decided from an
- // empty cache. Until this line the sweeper behaves exactly as it did
- // pre-v2.11.0 (which is also what every unit fixture that never
- // installs the hook keeps doing).
- store_->setTtlExpiryRelationHook(
- [this](const std::string& qualifiedCollection, const std::string& id,
- bool firstAttempt) {
- return ttlExpiryRelationDecision(qualifiedCollection, id, firstAttempt);
- });
- // v2.9.0 — re-apply persisted index declarations to each project's LMDB
- // store. This is load-bearing, not bookkeeping: the declaration is what
- // makes the write path maintain an index, and the planner consults an
- // index purely on the declaration's word. If this step were skipped, a
- // restart would leave indexes recorded but unmaintained, and queries
- // would be served from a frozen index - stale rows returned as current,
- // with nothing logged.
- applyIndexDeclarations();
- // v2.11.0 final review (finding 4) — THE POST-REPLAY RE-MIRROR PASS
- // RUNS HERE, and the position is load-bearing.
- //
- // PersistenceManager::recover() (above, before any of the arming
- // steps) only collects the id list. The pass writes documents through
- // LmdbDocumentStore::put(), which reads indexed_fields(),
- // unique_fields() and relations() LIVE and maintains every index
- // inside the document's own write transaction. Run from inside
- // recover(), all three maps were still empty, so it repaired the
- // document and left every posting derived from it stale: an
- // index-served query then returns the row under its OLD value and
- // misses it under its new one - silently wrong rows, which is exactly
- // what putting maintainIndexes() inside the transaction exists to
- // prevent, and which bit any install with a v2.9 index declared
- // whether or not it uses relations. It also made the reverse-index
- // half of relation_cascade.hpp's convergence claim false, and made
- // UniqueViolation (and therefore the pass's own retry phase)
- // unreachable.
- //
- // Must be after applyRelationDeclarations() AND
- // applyIndexDeclarations(); backfillIntoDocStore() below already sits
- // at this point for the same reason. Never throws.
- persistence_->runPendingRemirror(*store_, recovery_outcome_);
- // Run migrations if enabled
- if (config_.migrations.enabled && !config_.migrations.directory.empty()) {
- if (!runMigrations()) {
- if (config_.migrations.failOnError) {
- spdlog::error("Failed to run migrations");
- return false;
- }
- spdlog::warn("Some migrations failed, continuing anyway");
- }
- }
- // v2.0 Stage 4 — synchronous backfill into doc_store_ before
- // serving. Without this, only post-boot writes land in LMDB and
- // the read-flip downstream would see an empty mirror.
- backfillIntoDocStore();
- // v2.4.4 — placement audit. Reports documents sitting in a sub-db
- // other than the one they declare, which is the signature of a
- // misbound MDB_dbi (see storage/subdb_identity.hpp). Read-only and
- // advisory: it never blocks startup, because a misplacement is a
- // data-location problem an operator repairs offline with
- // `smartbotic-db-cli reconcile-subdbs`, not a reason to refuse
- // service on the other 99% of the dataset.
- auditSubdbPlacement();
- // v2.6.0 — stamp `project: "default"` onto file records written before
- // files were namespaced. Idempotent, and advisory: a file-metadata
- // stamp is not a reason to refuse service, so failures warn only.
- if (files_) {
- try {
- const uint32_t stamped = files_->stampMissingProjects();
- if (stamped > 0) {
- spdlog::info("v2.6.0 file migration: stamped {} legacy file "
- "record(s) with project 'default'", stamped);
- } else {
- spdlog::debug("v2.6.0 file migration: no legacy file records");
- }
- } catch (const std::exception& e) {
- spdlog::warn("v2.6.0 file migration failed: {}", e.what());
- }
- }
- // v2.11.0 close-out — LAST ACT OF THE BOOT PATH, and the position is
- // load-bearing. Everything above that can bump the drift counter (the
- // post-replay re-mirror pass, backfillIntoDocStore(), a malformed
- // legacy collection key) is now behind the baseline, so
- // mirrorDriftSinceReady() reports only drift a LIVE write caused.
- // Destructive relation policies gate on that, not on the raw counter:
- // see MemoryStore::markMirrorDriftBaseline() for the trap this closes
- // (one unrepairable row disabling every cascade/set_null delete and
- // every CreateRelation for the process's life, with an error message
- // telling the operator to restart, which re-incurs it). Reads keep
- // gating on the RAW counter - a stale row genuinely means MemoryStore
- // is ahead, and that fallback must stay.
- store_->markMirrorDriftBaseline();
- if (store_->mirrorDriftBaseline() > 0) {
- spdlog::warn("mirror drift at READY: {} - accrued on the boot path (see the "
- "ERROR lines above for the exact rows). Reads for this process "
- "will be served from MemoryStore rather than LMDB. Destructive "
- "relation policies are NOT disabled by this, but the rows named "
- "above are still stale in LMDB until each is rewritten through "
- "an ordinary write; a restart re-attempts and, if the cause "
- "persists, re-incurs the same count.",
- store_->mirrorDriftBaseline());
- }
- spdlog::info("Database service initialized successfully");
- return true;
- } catch (const std::exception& e) {
- spdlog::error("Failed to initialize database service: {}", e.what());
- return false;
- }
- }
- void DatabaseService::migrateLegacyEnvToDefaultProject() {
- const auto legacy = config_.dataDirectory / "env";
- const auto target = config_.dataDirectory / "projects" /
- smartbotic::database::kDefaultProject / "env";
- std::error_code ec;
- const bool legacy_exists = std::filesystem::exists(legacy, ec);
- const bool target_exists = std::filesystem::exists(target, ec);
- if (!legacy_exists) return; // fresh install or already migrated
- if (target_exists) {
- spdlog::error("v2.3 storage: both legacy '{}' and new '{}' exist — "
- "refusing to start. Inspect manually; the safe move is "
- "to either delete the legacy dir (if you confirm it's a "
- "leftover) or stop and contact ops.",
- legacy.string(), target.string());
- throw std::runtime_error("v2.3 storage migration: ambiguous layout");
- }
- std::filesystem::create_directories(target.parent_path(), ec);
- if (ec) {
- throw std::runtime_error("v2.3 storage migration: cannot create '"
- + target.parent_path().string() + "': " + ec.message());
- }
- std::filesystem::rename(legacy, target, ec);
- if (ec) {
- throw std::runtime_error("v2.3 storage migration: rename '"
- + legacy.string() + "' -> '"
- + target.string() + "' failed: " + ec.message());
- }
- spdlog::info("v2.3 storage: migrated legacy env '{}' -> '{}' (default project)",
- legacy.string(), target.string());
- }
- smartbotic::db::storage::DocumentStore* DatabaseService::docStore() noexcept {
- return docStore(smartbotic::database::kDefaultProject);
- }
- smartbotic::db::storage::DocumentStore*
- DatabaseService::docStore(std::string_view project) noexcept {
- if (!projects_) return nullptr;
- return projects_->get(project);
- }
- std::vector<std::string> DatabaseService::listProjects() const {
- if (!projects_) return {};
- return projects_->listOnDisk();
- }
- bool DatabaseService::createProject(const std::string& name, std::string& error) {
- if (!projects_) {
- error = "project registry not initialized";
- return false;
- }
- return projects_->create(name, error);
- }
- bool DatabaseService::dropProject(const std::string& name, std::string& error) {
- if (!projects_) {
- error = "project registry not initialized";
- return false;
- }
- // v2.11.0 close-out — capture the relation declarations belonging to this
- // project BEFORE the drop. Relations are validated same-project at
- // creation (no LMDB transaction spans two envs), so `child` and `parent`
- // are always in the same project as `name` itself; listRelations(name)
- // therefore names every declaration the drop invalidates. Captured first
- // because listRelations needs the project to still be nameable and, more
- // importantly, because un-arming needs the store to still exist.
- std::vector<RelationInfo> doomed;
- if (relation_manager_) {
- try {
- doomed = relation_manager_->listRelations(name);
- } catch (const std::exception& e) {
- spdlog::warn("v2.11 relations: could not enumerate relations of project '{}' "
- "before dropping it ({}); its declarations may be left behind",
- name, e.what());
- }
- }
- if (!projects_->drop(name, error)) return false;
- // ---- The drop succeeded; the declarations now point at a namespace that
- // no longer exists. Sweep them out of the global `_relations` collection
- // and un-arm each affected child collection.
- //
- // Order (drop first, sweep second) is deliberate: the registry is the one
- // that refuses "default" and rejects unknown names, so sweeping first would
- // destroy declarations for a project whose drop then failed. The cost of
- // this order is that the per-project store is usually already gone by the
- // time we try to un-arm, which makes the un-arm a no-op - correct, since a
- // store that no longer exists cannot be maintaining a reverse index. It
- // still runs, because ProjectStoreRegistry may hold a live
- // LmdbDocumentStore for a re-created project of the same name.
- //
- // ⚠ Advisory, never fatal: the project IS dropped at this point, and a
- // leftover declaration is litter, not corruption (its child and parent
- // collections are both gone, so nothing can be enforced against them). A
- // failure here must not report the drop as failed.
- for (const auto& r : doomed) {
- try {
- std::string err;
- if (!relation_manager_->dropRelation(r.name, err)) {
- spdlog::warn("v2.11 relations: could not drop declaration '{}' left by "
- "dropped project '{}': {}", r.name, name, err);
- continue;
- }
- const auto rc = resolveCollection(r.child);
- auto* lmdb = dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(
- docStore(rc.project));
- if (lmdb != nullptr) {
- // Empty list = "this collection has no relations", which is what
- // set_relations replaces the previous list with. Every relation
- // of a dropped project goes, so clearing the child wholesale is
- // exactly right here (unlike DropRelation, which must re-arm the
- // survivors).
- lmdb->set_relations(rc.collection, {});
- }
- spdlog::info("v2.11 relations: dropped declaration '{}' with the project '{}'",
- r.name, name);
- } catch (const std::exception& e) {
- spdlog::warn("v2.11 relations: could not clean up declaration '{}' after "
- "dropping project '{}': {}", r.name, name, e.what());
- }
- }
- // ⚠ KNOWN, NOT FIXED HERE (v2.11.0 close-out): _views, _policies and
- // _collection_meta have the SAME gap - ProjectStoreRegistry::drop() removes
- // the project's env and nothing else, so a dropped project leaves its view
- // definitions, access policies and per-collection configs behind in those
- // global system collections too. Not a shared fix: each manager keys and
- // caches differently (ViewManager keys `<project>:<name>`, PolicyManager
- // `<project>:<principal>`, CollectionConfigManager `<project>:<collection>`
- // and canonicalises the cache key only), so each needs its own sweep, and
- // PolicyManager's has a security dimension this one does not (a re-created
- // project inheriting a stale policy). Scoped out deliberately; recorded so
- // it is not rediscovered as new.
- return true;
- }
- bool DatabaseService::runMigrations() {
- if (!config_.migrations.enabled) {
- return true;
- }
- MigrationRunner::Config migrationConfig;
- migrationConfig.directory = config_.migrations.directory;
- migrationConfig.autoApply = config_.migrations.autoApply;
- migrationConfig.failOnError = config_.migrations.failOnError;
- migrationRunner_ = std::make_unique<MigrationRunner>(
- *store_, *view_manager_, *relation_manager_, migrationConfig);
- const bool ok = migrationRunner_->runMigrations();
- // v2.11.0 close-out — a `create_relation` migration op only PERSISTS the
- // declaration. Arming the write path and building the reverse index over
- // rows that already exist is what applyRelationDeclarations() does, and it
- // has already run for this boot (initialize() calls it before migrations),
- // so without this a migration-declared relation would maintain no reverse
- // index and enforce nothing until the NEXT restart - the "declared but not
- // enforcing" state this feature's self-heal exists to make unreachable.
- // Idempotent and cheap: re-arming an already-armed child replaces an
- // identical list, and the bootstrap loop skips any relation whose index
- // sub-db already exists. Runs even when the migration run reported failure,
- // because a partially-applied run may still have created a relation.
- applyRelationDeclarations(/*afterMigrations=*/true);
- return ok;
- }
- void DatabaseService::applyIndexDeclarations() {
- size_t applied = 0;
- for (const auto& [qualified, cfg] : config_manager_->allConfigs()) {
- if (cfg.indexedFields.empty()) continue;
- try {
- const auto rc = resolveCollection(qualified);
- auto* ds = docStore(rc.project);
- auto* lmdb =
- dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(ds);
- if (lmdb == nullptr) continue;
- lmdb->set_indexed_fields(rc.collection, cfg.indexedFields);
- // v2.11.0 T11 — arm uniqueFields the same way. This is the fix
- // for the exact failure mode the type was built to avoid:
- // uniqueFields has been persisted since T8 with no RPC or proto
- // surface, which was safe (nothing consulted it); now that
- // set_unique_fields is reachable, a restart that skipped this
- // call would silently stop enforcing a constraint every write
- // handler still advertises as active.
- lmdb->set_unique_fields(rc.collection, cfg.uniqueFields);
- ++applied;
- spdlog::info("v2.9 index: {} field(s) active on {} ({} unique)",
- cfg.indexedFields.size(), qualified, cfg.uniqueFields.size());
- // Self-heal a declared index whose sub-db is absent. That happens
- // when the KEY FORMAT VERSION in the sub-db prefix changes (v2.9.1
- // unified the numeric encoding, so `_idx_` became `_idx2_`), and it
- // would otherwise leave the declaration pointing at nothing: writes
- // would maintain the new index correctly, but rows written BEFORE the
- // upgrade would have no postings, so an indexed query would silently
- // return only the newer ones. Rebuilding is the only safe reading.
- for (const auto& field : cfg.indexedFields) {
- if (lmdb->index_stats(rc.collection, field).has_value()) continue;
- const uint64_t rows = lmdb->build_index(rc.collection, field);
- spdlog::warn("v2.9 index: rebuilt {}#{} over {} row(s) - the index "
- "was declared but its sub-db was absent (key format "
- "change or a manual removal)",
- qualified, field, rows);
- }
- } catch (const std::exception& e) {
- // Advisory per collection: one unparseable name must not stop the
- // rest from being armed. Loud, because a missing declaration means
- // that collection's index silently stops being maintained.
- spdlog::error("v2.9 index: could not apply declarations for {}: {}",
- qualified, e.what());
- }
- }
- if (applied > 0) {
- spdlog::info("v2.9 index: applied declarations for {} collection(s)", applied);
- }
- }
- MemoryStore::TtlExpiryAction DatabaseService::ttlExpiryRelationDecision(
- const std::string& qualifiedCollection, const std::string& id, bool firstAttempt) {
- using Action = MemoryStore::TtlExpiryAction;
- // Thin binder over relations/relation_cascade.cpp's ttlExpiryDecision() -
- // the policy itself lives there so it is unit-testable against the same
- // fixtures the cascade tests use, without a DatabaseService. All this does
- // is resolve the per-project LMDB store and supply the replication/events
- // notifier.
- if (!relation_manager_ || !config_manager_ || !store_ || !persistence_) {
- return Action::Proceed;
- }
- if (qualifiedCollection.empty() || qualifiedCollection[0] == '_') return Action::Proceed;
- try {
- const auto rc = resolveCollection(qualifiedCollection);
- auto* lmdb = dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(
- docStore(rc.project));
- // No LMDB substrate means no reverse index to consult, so the
- // pre-v2.11.0 behaviour is all that is available.
- if (lmdb == nullptr) return Action::Proceed;
- // Replication + Subscribe events, driven per mutation - the same wiring
- // DatabaseGrpcImpl::performRelationCascade() uses. Without it a
- // TTL-driven cascade would be invisible to followers and subscribers,
- // and a follower that never saw the child deletions diverges
- // permanently (the v2.3.1 class of bug).
- auto notify = [this](const std::string& coll, const std::string& docId,
- const std::optional<Document>& doc, EventType eventType) {
- notifyReplicationAndEvents(coll, docId, doc, eventType);
- };
- // firstAttempt -> the refusal (if any) is logged at WARN; a retry of an
- // already-reported document logs at DEBUG. The sweeper owns that
- // decision because it is the only thing that knows the document has
- // been refused before.
- return smartbotic::database::ttlExpiryDecision(
- *relation_manager_, *lmdb, *persistence_, *store_, *config_manager_,
- qualifiedCollection, id, notify, firstAttempt);
- } catch (const std::exception& e) {
- // Never expire a parent whose children could not be evaluated.
- spdlog::error("TTL expiry of '{}/{}': could not resolve its storage ({}) - the "
- "document is left in place and will be retried",
- qualifiedCollection, id, e.what());
- return Action::Skip;
- }
- }
- void DatabaseService::applyRelationDeclarations(bool afterMigrations) {
- // Group by (project, bare child collection) — set_relations() is a
- // per-collection call on that project's LmdbDocumentStore and replaces
- // whatever was declared before, so every relation sharing a child
- // collection must land in one call.
- std::unordered_map<std::string,
- std::vector<smartbotic::db::storage::RelationRef>> byChild;
- // Track which project each grouping key belongs to alongside the bare
- // collection name, since the map key alone doesn't carry it.
- std::unordered_map<std::string, std::pair<std::string, std::string>> keyToProjectCollection;
- for (const auto& r : relation_manager_->listRelations()) {
- try {
- const auto rn = resolveCollection(r.name);
- const auto rc = resolveCollection(r.child);
- // v2.11.0 T13 — parent resolved to its bare name, same reasoning
- // as armRelationsForChild's mirror of this construction.
- const auto pc = resolveCollection(r.parent);
- const std::string key = rc.project + ":" + rc.collection;
- // v2.11.0 T13 round 2 — snapshot relationsEnforced at boot-time
- // arming, same reasoning as armRelationsForChild's mirror of
- // this construction. `key` is already the qualified child name
- // configFor expects.
- const bool enforced = config_manager_->configFor(key).relationsEnforced;
- byChild[key].push_back(smartbotic::db::storage::RelationRef{
- rn.collection, r.childField, pc.collection, r.validateOnWrite, enforced});
- keyToProjectCollection[key] = {rc.project, rc.collection};
- } catch (const std::exception& e) {
- // Advisory per relation: one unparseable declaration must not
- // stop the rest from being armed. Loud, because a skipped
- // relation silently maintains no reverse index for its child.
- spdlog::error("v2.11 relations: could not apply declaration '{}': {}",
- r.name, e.what());
- }
- }
- size_t applied = 0;
- for (const auto& [key, refs] : byChild) {
- const auto& [project, collection] = keyToProjectCollection[key];
- try {
- // ⚠ v2.11.0 close-out — getOrCreate, not get(), on the POST-MIGRATION
- // pass. docStore() is projects_->get(), which does NOT create an env,
- // and a `create_relation` migration in a project that has never been
- // written to has no env yet: the declaration persisted, this pass
- // skipped it, and the relation maintained no reverse index and
- // enforced NOTHING until the next restart - the exact
- // "declared but not enforcing" state the op was added to avoid, just
- // narrowed to new projects. A migration declaring a relation IS a
- // statement that the project exists, and the mirror resolver would
- // getOrCreate the same env on the first write anyway.
- //
- // The BOOT pass deliberately keeps get(): creating envs there would
- // resurrect the env of a project whose declarations are merely stale.
- auto* ds = (afterMigrations && projects_)
- ? projects_->getOrCreate(project)
- : docStore(project);
- auto* lmdb = dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(ds);
- if (lmdb == nullptr) continue;
- lmdb->set_relations(collection, refs);
- ++applied;
- spdlog::info("v2.11 relations: {} relation(s) active with child '{}:{}'",
- refs.size(), project, collection);
- // v2.11.0 T7 self-heal — mirrors applyIndexDeclarations()'s
- // rebuild of a declared-but-absent index sub-db, for the same
- // reason: a relation whose reverse index sub-db is missing is
- // exactly the "declared but not enforcing" bug this task exists
- // to close. Reachable paths that leave a relation in this state:
- // a snapshot/backup restored from before the sub-db existed, or
- // an operator manually dropping the `_relidx1_*` sub-db (there
- // is no dropRelation-only-the-index tool). CreateRelation's own
- // bootstrap scan (see database_grpc_impl.cpp) covers the normal
- // declare-time path; this covers everything else.
- //
- // ⚠ v2.11.0 close-out (review minor 2) — the SAME code path is the
- // NORMAL, expected one for a relation a migration just declared:
- // migrations run after this pass, so runMigrations() calls it again
- // and the brand-new relation legitimately has no sub-db yet. Logging
- // "declared but its sub-db was absent (restored snapshot or a manual
- // removal)" at WARN there describes a fault that did not happen.
- // `afterMigrations` distinguishes the two so the text and the level
- // match the event.
- for (const auto& ref : refs) {
- if (lmdb->relation_index_exists(ref.name)) continue;
- const uint64_t rows =
- lmdb->build_relation_index(ref.name, collection, ref.childField);
- if (afterMigrations) {
- spdlog::info("v2.11 relations: built the reverse index for '{}' over {} "
- "existing row(s) in '{}:{}' - a migration declared this "
- "relation on this boot",
- ref.name, rows, project, collection);
- } else {
- spdlog::warn("v2.11 relations: rebuilt '{}' over {} row(s) in '{}:{}' - the "
- "relation was declared but its reverse index sub-db was absent "
- "(restored snapshot predating it, or a manual removal)",
- ref.name, rows, project, collection);
- }
- }
- } catch (const std::exception& e) {
- spdlog::error("v2.11 relations: could not apply declarations for '{}:{}': {}",
- project, collection, e.what());
- }
- }
- if (applied > 0) {
- spdlog::info("v2.11 relations: applied declarations for {} child collection(s)", applied);
- }
- }
- void DatabaseService::auditSubdbPlacement() {
- if (!projects_) return;
- for (const auto& name : projects_->listOpen()) {
- auto h = projects_->getHandle(name);
- if (!h.env) continue;
- try {
- // audit_env(), NOT audit(). The path-taking overload opens a
- // second MDB_env on the same file; closing it would drop the
- // POSIX record locks this process already holds for `h.env`
- // (POSIX locks are per-process, and closing any fd on a file
- // releases all of them), leaving every later mdb_txn_begin
- // failing EINVAL. Borrow the open handle instead.
- const auto rep = smartbotic::db::storage::audit_env(h.env->raw(), name);
- if (rep.misplaced.empty()) {
- spdlog::debug("placement audit: project '{}' consistent ({} rows)",
- name, rep.rows_scanned);
- continue;
- }
- // Summarise per (physical -> declared home) pair; one line per
- // affected document would be unreadable at scale.
- std::map<std::string, uint64_t> pairs;
- for (const auto& m : rep.misplaced) {
- ++pairs[m.physical_subdb + " -> " + m.home_subdb];
- }
- spdlog::error("placement audit: project '{}' has {} document(s) in the "
- "wrong sub-db. Reads served from LMDB will not find them "
- "under their own collection. Repair offline with: "
- "smartbotic-db-cli reconcile-subdbs --env {} --project {} --apply",
- name, rep.misplaced.size(), h.env->path(), name);
- for (const auto& [pair, count] : pairs) {
- spdlog::error("placement audit: {} document(s) {}", count, pair);
- }
- } catch (const std::exception& e) {
- // Advisory only — a failed audit must not take the service down.
- spdlog::warn("placement audit: project '{}' could not be audited: {}",
- name, e.what());
- }
- }
- }
- bool DatabaseService::backfillIntoDocStore() {
- if (!projects_) {
- // Registry open failed earlier; nothing to back-fill into.
- return true;
- }
- // v2.3 — MemoryStore stores collection keys in qualified form
- // ("project:collection"). For the default project (back-compat
- // path) the marker check operates on the default env. Skip when
- // already migrated.
- auto handle = projects_->getHandle(smartbotic::database::kDefaultProject);
- if (handle.env) {
- try {
- if (smartbotic::db::storage::migration_complete(*handle.env)) {
- spdlog::info("v2.3 backfill: default project already migrated, skipping");
- return true;
- }
- } catch (const std::exception& e) {
- spdlog::warn("v2.3 backfill: migration_complete probe failed: {} — "
- "treating default env as fresh and proceeding",
- e.what());
- }
- }
- const auto collections = store_->listCollections();
- uint64_t total_docs = 0;
- uint64_t total_failures = 0;
- auto t0 = std::chrono::steady_clock::now();
- for (const auto& qualified : collections) {
- if (qualified.empty() || qualified[0] == '_') continue;
- // Each qualified key parses to (project, collection); route the
- // mirror write to the right project's env.
- smartbotic::database::ProjectCollection pc;
- try {
- pc = smartbotic::database::parseProjectCollection(qualified);
- } catch (const std::exception&) {
- // Legacy unqualified key — treat as default project.
- pc = {std::string(smartbotic::database::kDefaultProject), qualified};
- }
- auto* ds = projects_->getOrCreate(pc.project);
- if (!ds) {
- ++total_failures;
- continue;
- }
- const auto docs = store_->getAllDocuments(qualified);
- for (const auto& doc : docs) {
- std::optional<Document> opt_doc(doc);
- const uint64_t drift_before = mirror_drift_count_.load(std::memory_order_relaxed);
- // v2.11.0 T13 round 2 — UniqueViolation/MissingParentReference are
- // rethrown UNCAUGHT by applyDualWriteMirror (deliberately - see
- // both exceptions' header comments), unlike a genuine mirror
- // fault which is swallowed and only shows up as a drift bump
- // below. Uncaught here would escape this loop, backfillIntoDocStore
- // and initialize() itself, refusing startup over one rejected row.
- // Practically unreachable today - backfill only runs migrating a
- // v1.x dataset, which predates both `_relations` and unique-index
- // declarations - but catching removes that reasoning burden for
- // the next reader, same as any other row here that fails.
- try {
- smartbotic::db::storage::applyDualWriteMirror(
- ds, mirror_healthy_, mirror_drift_count_,
- pc.collection, doc.id, opt_doc, EventType::INSERT);
- } catch (const std::exception& e) {
- spdlog::error("v2.3 backfill: mirror rejected {}/{}: {} - row stays "
- "in MemoryStore only, LMDB does not have it",
- qualified, doc.id, e.what());
- ++total_failures;
- }
- if (mirror_drift_count_.load(std::memory_order_relaxed) > drift_before) {
- ++total_failures;
- }
- ++total_docs;
- if (total_docs % 10000 == 0) {
- sd_notify(0, "EXTEND_TIMEOUT_USEC=600000000");
- const std::string status = "STATUS=v2.3 backfill: " +
- std::to_string(total_docs) + " docs mirrored";
- sd_notify(0, status.c_str());
- spdlog::info("v2.3 backfill: {} docs mirrored ({} failures so far)",
- total_docs, total_failures);
- }
- }
- }
- const auto elapsed = std::chrono::duration_cast<std::chrono::milliseconds>(
- std::chrono::steady_clock::now() - t0).count();
- if (total_failures == 0) {
- // Stamp schema_version=2 on every opened project env. Subsequent
- // boots short-circuit at migration_complete().
- for (const auto& name : projects_->listOpen()) {
- auto h = projects_->getHandle(name);
- if (!h.env) continue;
- try {
- smartbotic::db::storage::mark_migration_complete(*h.env);
- } catch (const std::exception& e) {
- spdlog::warn("v2.3 backfill: schema_version marker write failed for "
- "project '{}': {}", name, e.what());
- }
- }
- spdlog::info("v2.3 backfill: complete — {} docs across {} collections in {} ms",
- total_docs, collections.size(), elapsed);
- } else {
- spdlog::error("v2.3 backfill: completed with {} failures out of {} docs "
- "in {} ms — mirror_healthy_=false, reads must NOT flip to LMDB",
- total_failures, total_docs, elapsed);
- }
- return true;
- }
- void DatabaseService::start() {
- if (running_.exchange(true)) {
- return;
- }
- spdlog::info("Starting database service with {} listener(s)", config_.listeners.size());
- // Start components
- store_->start();
- persistence_->start();
- events_->start();
- files_->start();
- if (config_.replicationEnabled) {
- replication_->start();
- // Set initial local collections for discovery
- replication_->setLocalCollections(store_->listCollections());
- }
- // Start gRPC server
- startGrpcServer();
- spdlog::info("Database service started");
- }
- void DatabaseService::stop() {
- if (!running_.exchange(false)) {
- return;
- }
- spdlog::info("Stopping database service...");
- // Stop gRPC server first
- stopGrpcServer();
- // Stop components in reverse order
- if (replication_) {
- replication_->stop();
- }
- if (files_) {
- files_->stop();
- }
- if (events_) {
- events_->stop();
- }
- if (persistence_ && store_) {
- // v1.8.0 — take a final snapshot on graceful shutdown so the next
- // boot's recovery is "trivial" (snapshot-loaded) instead of "WAL-only
- // replay" which trips auto-readonly mode. This makes ordinary
- // restarts pass through cleanly without needing
- // `smartbotic-db-cli unlock`.
- //
- // v1.8.2 — extend systemd's stop watchdog before doing it. On Zoe
- // (1.85M docs / ~5 GB tracked) the serialize takes ~70 s and the
- // default TimeoutStopSec=30s SIGKILLed the process mid-write, peak
- // RSS hit 11 GB (live state + uncompressed serialize buffer), and
- // the new snapshot never made it to disk anyway. EXTEND_TIMEOUT_USEC
- // tells systemd we're working — same protocol the recovery path
- // uses during startup. Paired with TimeoutStopSec=300 in the unit
- // file as the static cap.
- #ifdef HAVE_SYSTEMD
- sd_notify(0, "EXTEND_TIMEOUT_USEC=600000000");
- sd_notify(0, "STATUS=Writing final snapshot");
- #endif
- try {
- persistence_->forceSnapshot(*store_);
- spdlog::info("Final snapshot taken on shutdown");
- } catch (const std::exception& e) {
- spdlog::warn("Final snapshot on shutdown failed: {} — recovery on next "
- "boot may be WAL-only and trip auto-readonly mode", e.what());
- }
- }
- if (persistence_) {
- persistence_->stop();
- }
- if (store_) {
- store_->stop();
- }
- spdlog::info("Database service stopped");
- }
- void DatabaseService::wait() {
- if (serverThread_.joinable()) {
- serverThread_.join();
- }
- }
- void DatabaseService::signalStop() {
- stopRequested_ = true;
- stop();
- }
- nlohmann::json DatabaseService::getStats() const {
- nlohmann::json stats;
- if (store_) {
- auto storeStats = store_->getStats();
- stats["documents"] = storeStats.totalDocuments;
- stats["collections"] = storeStats.totalCollections;
- stats["memory_bytes"] = storeStats.estimatedMemoryBytes;
- stats["inserts"] = storeStats.insertCount;
- stats["updates"] = storeStats.updateCount;
- stats["deletes"] = storeStats.deleteCount;
- stats["queries"] = storeStats.queryCount;
- }
- if (persistence_) {
- auto persistStats = persistence_->getStats();
- stats["wal_sequence"] = persistStats.walSequence;
- stats["wal_size_bytes"] = persistStats.walSizeBytes;
- stats["snapshot_count"] = persistStats.snapshotCount;
- stats["last_snapshot_sequence"] = persistStats.lastSnapshotSequence;
- }
- if (events_) {
- stats["subscriptions"] = events_->subscriptionCount();
- }
- return stats;
- }
- DatabaseService::Config DatabaseService::loadConfig(const std::filesystem::path& configPath) {
- std::ifstream file(configPath);
- if (!file) {
- throw std::runtime_error("Failed to open config file: " + configPath.string());
- }
- nlohmann::json json;
- file >> json;
- return parseConfig(json);
- }
- DatabaseService::Config DatabaseService::parseConfig(const nlohmann::json& json) {
- Config config;
- // Check for both "storage" and "database" keys for backward compatibility
- const nlohmann::json* dbConfig = nullptr;
- if (json.contains("database")) {
- dbConfig = &json["database"];
- } else if (json.contains("storage")) {
- dbConfig = &json["storage"];
- }
- if (dbConfig) {
- const auto& db = *dbConfig;
- config.nodeId = db.value("node_id", config.nodeId);
- config.bindAddress = db.value("bind_address", config.bindAddress);
- config.rpcPort = db.value("rpc_port", config.rpcPort);
- // v2.4 — parse the per-listener fleet. If "listeners" is set,
- // it overrides the legacy bind_address+rpc_port. Otherwise we
- // synthesise one listener below using those legacy fields.
- if (db.contains("listeners") && db["listeners"].is_array()) {
- for (const auto& l : db["listeners"]) {
- ListenerConfig lc;
- lc.bind = l.value("bind", lc.bind);
- lc.port = l.value("port", lc.port);
- if (l.contains("tls") && l["tls"].is_object()) {
- const auto& t = l["tls"];
- lc.tls.enabled = t.value("enabled", false);
- lc.tls.cert_path = t.value("cert_path", std::string{});
- lc.tls.key_path = t.value("key_path", std::string{});
- lc.tls.auto_self_signed_if_missing =
- t.value("auto_self_signed_if_missing", true);
- }
- if (l.contains("auth") && l["auth"].is_object()) {
- const auto& a = l["auth"];
- lc.auth.required = a.value("required", false);
- if (a.contains("keys") && a["keys"].is_array()) {
- for (const auto& k : a["keys"]) {
- // v2.7.0 — a key is either a bare token (v2.4-v2.6
- // shape) or {"name","key"}. The name becomes the
- // principal that policy is written against; a bare
- // token maps to the reserved `unnamed`, so old
- // configs keep authenticating but cannot be told
- // apart in a policy until they are named.
- if (k.is_string()) {
- lc.auth.keys.push_back(
- {smartbotic::database::auth::kPrincipalUnnamed,
- k.get<std::string>()});
- } else if (k.is_object()) {
- const std::string name = k.value("name", std::string{});
- const std::string key = k.value("key", std::string{});
- if (key.empty()) {
- spdlog::warn("auth: listener {}:{} has a key entry "
- "with no `key` value - skipped",
- lc.bind, lc.port);
- continue;
- }
- lc.auth.keys.push_back(
- {name.empty()
- ? std::string(smartbotic::database::auth::kPrincipalUnnamed)
- : name,
- key});
- }
- }
- }
- }
- config.listeners.push_back(std::move(lc));
- }
- }
- // Expand environment variables in data directory
- std::string dataDir = db.value("data_directory", "");
- if (dataDir.find("${HOME}") != std::string::npos) {
- const char* home = std::getenv("HOME");
- if (home) {
- size_t pos = dataDir.find("${HOME}");
- dataDir.replace(pos, 7, home);
- }
- }
- config.dataDirectory = dataDir;
- // Migrations settings
- if (db.contains("migrations")) {
- auto& migrations = db["migrations"];
- config.migrations.enabled = migrations.value("enabled", config.migrations.enabled);
- config.migrations.autoApply = migrations.value("auto_apply", config.migrations.autoApply);
- config.migrations.failOnError = migrations.value("fail_on_error", config.migrations.failOnError);
- std::string migDir = migrations.value("directory", "");
- if (migDir.find("${HOME}") != std::string::npos) {
- const char* home = std::getenv("HOME");
- if (home) {
- size_t pos = migDir.find("${HOME}");
- migDir.replace(pos, 7, home);
- }
- }
- config.migrations.directory = migDir;
- }
- // Memory eviction settings
- if (db.contains("memory")) {
- auto& memory = db["memory"];
- // Deprecation log. These knobs control MemoryStore eviction, which
- // is still active for the MemoryStore-side mirror but does NOT bound
- // RSS since v2.0 (LMDB mmap is the dominant RSS contributor). The
- // substrate-level equivalent is the LMDB env mapsize and the OS page
- // cache.
- //
- // They disappear when the write-handler migration deletes MemoryStore
- // (docs/ROADMAP.md, "Pending"). An earlier revision of this comment
- // promised a v2.1 rename to `buffer_pool_size_mb` "per the Phase C
- // plan"; that plan was superseded by LMDB before it was written, no
- // such knob exists, and the rename is not planned.
- for (const char* deprecated : {"max_memory_mb",
- "eviction_threshold_percent",
- "eviction_target_percent",
- "eviction_check_interval_ms",
- "eviction_chunk_size",
- "eviction_chunk_pause_ms",
- "max_eviction_passes_per_trigger",
- "hot_write_floor_ms",
- "memory_priority"}) {
- if (memory.contains(deprecated)) {
- spdlog::warn("v2.0 deprecation: storage.memory.{} is deprecated "
- "and will be removed in v2.1. v2.0 RSS is bounded "
- "by the LMDB env mapsize + OS page cache, not by "
- "MemoryStore eviction. Setting still applied to "
- "the MemoryStore mirror for back-compat.",
- deprecated);
- }
- }
- config.maxMemoryMb = memory.value("max_memory_mb", config.maxMemoryMb);
- config.evictionThresholdPercent = memory.value("eviction_threshold_percent", config.evictionThresholdPercent);
- config.evictionTargetPercent = memory.value("eviction_target_percent", config.evictionTargetPercent);
- config.evictionCheckIntervalMs = memory.value("eviction_check_interval_ms", config.evictionCheckIntervalMs);
- // NEW v1.7.0 eviction tuning (schema-only in T2; runtime use lands in T3/T4)
- config.evictionChunkSize = memory.value("eviction_chunk_size", config.evictionChunkSize);
- config.evictionChunkPauseMs = memory.value("eviction_chunk_pause_ms", config.evictionChunkPauseMs);
- config.maxEvictionPassesPerTrigger = memory.value("max_eviction_passes_per_trigger", config.maxEvictionPassesPerTrigger);
- config.hotWriteFloorMs = memory.value("hot_write_floor_ms", config.hotWriteFloorMs);
- config.memorySoftPercent = memory.value("memory_soft_percent", config.memorySoftPercent);
- config.memoryHardPercent = memory.value("memory_hard_percent", config.memoryHardPercent);
- config.memoryEmergencyPercent = memory.value("memory_emergency_percent", config.memoryEmergencyPercent);
- // v1.7.0 T10 — eviction burst event threshold
- config.evictionBurstThreshold = memory.value("eviction_burst_threshold", config.evictionBurstThreshold);
- // v2.4.3 — eviction drain cap (0 disables)
- config.evictionMaxEpisodePercent = memory.value("eviction_max_episode_percent", config.evictionMaxEpisodePercent);
- }
- // v2.11.0 close-out — TTL sweep budgets. Its own block rather than
- // `memory`, because these bound expiry work per sweep, not memory.
- if (db.contains("ttl")) {
- auto& ttl = db["ttl"];
- config.ttlMaxCandidatesPerSweep =
- ttl.value("max_candidates_per_sweep", config.ttlMaxCandidatesPerSweep);
- config.ttlMaxBlockedRetriesPerSweep =
- ttl.value("max_blocked_retries_per_sweep", config.ttlMaxBlockedRetriesPerSweep);
- config.ttlBlockedRetrySweeps =
- ttl.value("blocked_retry_sweeps", config.ttlBlockedRetrySweeps);
- config.ttlBlockedSummarySweeps =
- ttl.value("blocked_summary_sweeps", config.ttlBlockedSummarySweeps);
- }
- // Persistence settings
- if (db.contains("persistence")) {
- auto& persistence = db["persistence"];
- config.walSyncIntervalMs = persistence.value("wal_sync_interval_ms", config.walSyncIntervalMs);
- config.snapshotIntervalSec = persistence.value("snapshot_interval_sec", config.snapshotIntervalSec);
- config.compressionEnabled = persistence.value("compression", "lz4") != "none";
- // Snapshot durability (NEW in v1.6.1)
- if (persistence.contains("snapshots")) {
- const auto& snap = persistence["snapshots"];
- config.persistenceConfig.validateAfterWrite = snap.value("validate_after_write", config.persistenceConfig.validateAfterWrite);
- config.persistenceConfig.cleanupOnlyIfVerified = snap.value("cleanup_only_if_verified", config.persistenceConfig.cleanupOnlyIfVerified);
- }
- // Recovery (NEW in v1.6.1)
- if (persistence.contains("recovery")) {
- const auto& rec = persistence["recovery"];
- std::string modeStr = rec.value("mode", std::string("normal"));
- try {
- config.persistenceConfig.recoveryMode = recoveryModeFromString(modeStr);
- } catch (const std::exception& e) {
- spdlog::warn("Invalid recovery.mode '{}', defaulting to 'normal': {}",
- modeStr, e.what());
- config.persistenceConfig.recoveryMode = RecoveryMode::Normal;
- }
- config.persistenceConfig.autoEscalate = rec.value("auto_escalate", config.persistenceConfig.autoEscalate);
- config.persistenceConfig.allowEmptyOnFreshInstall = rec.value("allow_empty_on_fresh_install", config.persistenceConfig.allowEmptyOnFreshInstall);
- }
- }
- // File settings
- if (db.contains("files")) {
- auto& files = db["files"];
- config.maxFileSizeMb = files.value("max_file_size_mb", config.maxFileSizeMb);
- // v2.8.0 — this interval now drives the file EXPIRY sweep as well
- // as orphan cleanup, so accept the plainer name too. The original
- // key kept working; nothing read either of them before 2.8.0
- // because there was no sweeper thread at all.
- config.fileCleanupIntervalSec =
- files.value("cleanup_orphans_interval_sec", config.fileCleanupIntervalSec);
- config.fileCleanupIntervalSec =
- files.value("cleanup_interval_sec", config.fileCleanupIntervalSec);
- // v2.8.0 — default retention per file type. The analogue of a
- // collection's defaultTtlSeconds, so retention is an operator
- // setting rather than something every uploader must remember.
- //
- // "default_ttl_seconds": {
- // "acme:generated": 86400, // one project, one type
- // "generated": 604800, // any project
- // "acme:*": 2592000, // every type in one project
- // "*": 0 // everything (0 = keep)
- // }
- if (files.contains("default_ttl_seconds") &&
- files["default_ttl_seconds"].is_object()) {
- for (const auto& [k, v] : files["default_ttl_seconds"].items()) {
- if (v.is_number_unsigned()) {
- config.fileDefaultTtlByType[k] = v.get<uint32_t>();
- }
- }
- }
- if (files.contains("allowed_types")) {
- for (const auto& type : files["allowed_types"]) {
- config.allowedFileTypes.push_back(type.get<std::string>());
- }
- }
- }
- // Encryption settings
- if (db.contains("encryption")) {
- auto& encryption = db["encryption"];
- config.encryptionEnabled = encryption.value("enabled", config.encryptionEnabled);
- config.autoGenerateKey = encryption.value("auto_generate_key", config.autoGenerateKey);
- std::string keyFile = encryption.value("key_file", "");
- if (keyFile.find("${HOME}") != std::string::npos) {
- const char* home = std::getenv("HOME");
- if (home) {
- size_t pos = keyFile.find("${HOME}");
- keyFile.replace(pos, 7, home);
- }
- }
- config.keyFilePath = keyFile;
- }
- // Replication settings
- if (db.contains("replication")) {
- auto& replication = db["replication"];
- config.replicationEnabled = replication.value("enabled", config.replicationEnabled);
- config.conflictResolution = replication.value("conflict_resolution", config.conflictResolution);
- if (replication.contains("peers")) {
- for (const auto& peer : replication["peers"]) {
- config.peerAddresses.push_back(peer.get<std::string>());
- }
- }
- }
- // gRPC settings (v1.6.2 — concurrency cap + configurable message sizes)
- if (db.contains("grpc")) {
- const auto& grpc = db["grpc"];
- config.grpc.maxReceiveMessageSizeMb =
- grpc.value("max_receive_message_size_mb", config.grpc.maxReceiveMessageSizeMb);
- config.grpc.maxSendMessageSizeMb =
- grpc.value("max_send_message_size_mb", config.grpc.maxSendMessageSizeMb);
- config.grpc.resourceQuotaMemoryMb =
- grpc.value("resource_quota_memory_mb", config.grpc.resourceQuotaMemoryMb);
- config.grpc.maxConcurrentSubscribeStreams =
- grpc.value("max_concurrent_subscribe_streams", config.grpc.maxConcurrentSubscribeStreams);
- config.grpc.maxConcurrentFileStreams =
- grpc.value("max_concurrent_file_streams", config.grpc.maxConcurrentFileStreams);
- }
- }
- // v2.4 back-compat shim. If no listeners[] was provided, synthesize
- // one from the legacy bind_address + rpc_port fields. This is what
- // every existing v2.3 deployment will hit on first v2.4 boot —
- // their configs don't know about listeners[] yet, and they get
- // exactly the same listener they had before.
- if (config.listeners.empty()) {
- ListenerConfig legacy;
- legacy.bind = config.bindAddress;
- legacy.port = config.rpcPort;
- legacy.tls.enabled = false;
- legacy.auth.required = false;
- config.listeners.push_back(std::move(legacy));
- }
- return config;
- }
- void DatabaseService::setupComponents() {
- // v2.3 multi-project storage substrate.
- //
- // Each project lives at `<dataDir>/projects/<name>/env/`. The default
- // project always exists and is what v2.0-v2.2 clients (which don't
- // address projects) implicitly use.
- //
- // Pre-2.3 installs have their data at the legacy `<dataDir>/env/`
- // path. migrateLegacyEnvToDefaultProject() detects that layout and
- // atomically moves it under `projects/default/env/` before the
- // registry opens.
- migrateLegacyEnvToDefaultProject();
- try {
- constexpr size_t kMapSize = 4ULL << 30; // 4 GiB per project; LMDB grows lazily.
- constexpr uint32_t kMaxDbs = 1024;
- projects_ = std::make_unique<smartbotic::db::storage::ProjectStoreRegistry>(
- config_.dataDirectory / "projects", kMapSize, kMaxDbs);
- const auto opened = projects_->openExisting();
- spdlog::info("v2.3 storage: opened {} project env(s): [{}]",
- opened.size(),
- [&opened]() {
- std::string s;
- for (size_t i = 0; i < opened.size(); ++i) {
- if (i) s += ", ";
- s += opened[i];
- }
- return s;
- }());
- } catch (const std::exception& e) {
- spdlog::error("v2.3 storage: failed to open project registry at '{}': {} -- "
- "service cannot start without LMDB; bailing",
- (config_.dataDirectory / "projects").string(), e.what());
- projects_.reset();
- throw;
- }
- // Create memory store with eviction config
- MemoryStore::Config storeConfig;
- storeConfig.nodeId = config_.nodeId;
- storeConfig.maxMemoryBytes = config_.maxMemoryMb * 1024ULL * 1024ULL;
- storeConfig.evictionThresholdPercent = config_.evictionThresholdPercent;
- storeConfig.evictionTargetPercent = config_.evictionTargetPercent;
- storeConfig.evictionCheckIntervalMs = config_.evictionCheckIntervalMs;
- // NEW v1.7.0 eviction tuning (wired in T2; consumed in T3/T4).
- storeConfig.evictionChunkSize = config_.evictionChunkSize;
- storeConfig.evictionChunkPauseMs = config_.evictionChunkPauseMs;
- storeConfig.maxEvictionPassesPerTrigger = config_.maxEvictionPassesPerTrigger;
- storeConfig.hotWriteFloorMs = config_.hotWriteFloorMs;
- storeConfig.memorySoftPercent = config_.memorySoftPercent;
- storeConfig.memoryHardPercent = config_.memoryHardPercent;
- storeConfig.memoryEmergencyPercent = config_.memoryEmergencyPercent;
- storeConfig.evictionBurstThreshold = config_.evictionBurstThreshold;
- storeConfig.evictionMaxEpisodePercent = config_.evictionMaxEpisodePercent;
- // v2.11.0 close-out — TTL sweep budgets (storage.ttl).
- storeConfig.ttlMaxCandidatesPerSweep = config_.ttlMaxCandidatesPerSweep;
- storeConfig.ttlMaxBlockedRetriesPerSweep = config_.ttlMaxBlockedRetriesPerSweep;
- storeConfig.ttlBlockedRetrySweeps = config_.ttlBlockedRetrySweeps;
- storeConfig.ttlBlockedSummarySweeps = config_.ttlBlockedSummarySweeps;
- store_ = std::make_unique<MemoryStore>(storeConfig);
- // Create view manager (cache loaded in initialize() after persistence recovery)
- view_manager_ = std::make_unique<ViewManager>(*store_);
- // Create per-collection config manager and attach it to the store so the
- // write paths route document timestamp stamps through it.
- config_manager_ = std::make_unique<CollectionConfigManager>(*store_);
- // v2.8.0 — access policy. Constructed here; its cache is loaded after
- // recovery alongside the view and collection-config caches.
- policy_manager_ = std::make_unique<PolicyManager>(*store_);
- // v2.11.0 T6a — relation declarations. Constructed here; its cache is
- // loaded after recovery alongside the view/config/policy caches.
- relation_manager_ = std::make_unique<RelationManager>(*store_);
- store_->setConfigManager(config_manager_.get());
- // v1.9.0 — disk-resident version history. Replaces the in-heap
- // `CollectionData::versionHistory` deques that were the dominant
- // source of the Zoe untracked-RSS gap. One file per collection at
- // `<dataDir>/history/<collection>.hlog`. Lifecycle: constructed
- // here, wired into the store immediately so snapshot deserializer
- // can migrate pre-v1.9 in-memory history blocks straight to disk.
- HistoryStore::Config histConfig;
- histConfig.dataDir = config_.dataDirectory;
- history_store_ = std::make_unique<HistoryStore>(histConfig);
- store_->setHistoryStore(history_store_.get());
- // v2.3 — wire the LMDB mirror INTO MemoryStore with a project
- // resolver. Every write that flows through store_ commits to the
- // right per-project LMDB env BEFORE the per-collection lock releases.
- if (projects_) {
- auto* registry = projects_.get();
- store_->setDocumentStoreMirror(
- [registry](std::string_view project) -> smartbotic::db::storage::DocumentStore* {
- return registry->getOrCreate(project);
- },
- &mirror_healthy_,
- &mirror_drift_count_);
- }
- // Create persistence manager. Start from any persistenceConfig values the
- // caller (config loader / CLI parser) has already populated — including
- // recoveryMode — and overlay the top-level convenience fields.
- PersistenceManager::Config persistConfig = config_.persistenceConfig;
- persistConfig.dataDir = config_.dataDirectory;
- persistConfig.walSyncIntervalMs = config_.walSyncIntervalMs;
- persistConfig.snapshotIntervalSec = config_.snapshotIntervalSec;
- persistConfig.compressionEnabled = config_.compressionEnabled;
- persistConfig.nodeId = config_.nodeId;
- // Mirror the resolved recoveryMode back into our Config so downstream code
- // (logging, RPC handlers) can inspect config_.persistenceConfig.recoveryMode
- // without having to reach into the PersistenceManager.
- config_.persistenceConfig = persistConfig;
- persistence_ = std::make_unique<PersistenceManager>(persistConfig);
- // Connect store callbacks to persistence - uses setPersistCallback for WAL logging
- store_->setPersistCallback([this](const std::string& collection, const std::string& id,
- const std::optional<Document>& doc, EventType eventType) {
- // Log to WAL based on event type. v1.7.4: capture the assigned WAL
- // sequence and record it on the store so eviction stubs can carry
- // a real seq (instead of doc.version, which was useless as a
- // `fromSequence` hint and forced a full WAL scan on every fault).
- uint64_t walSeq = 0;
- switch (eventType) {
- case EventType::INSERT:
- if (doc) walSeq = persistence_->logInsert(collection, *doc);
- break;
- case EventType::UPDATE:
- if (doc) walSeq = persistence_->logUpdate(collection, *doc);
- break;
- case EventType::DELETE:
- persistence_->logDelete(collection, id);
- break;
- default:
- break;
- }
- if (walSeq > 0 && (eventType == EventType::INSERT || eventType == EventType::UPDATE)) {
- store_->recordWalSequence(collection, id, walSeq);
- }
- // v2.0 dual-write — moved INTO MemoryStore::mirrorWriteToDocStore so
- // that it runs UNDER the per-collection write lock. That closes the
- // race window where readers could see a doc in MemoryStore but not
- // yet in LMDB. The callback now does only WAL + replication + events.
- // v2.11.0 T12 review (C2) — replication + events extracted into
- // notifyReplicationAndEvents() so relations/relation_cascade.cpp can
- // drive them explicitly too. See that method's doc comment.
- notifyReplicationAndEvents(collection, id, doc, eventType);
- });
- // Set up document load callback for LRU eviction recovery.
- //
- // v1.7.4: walSequence is now a real WAL sequence (set by the persist
- // callback below at append time, not the per-doc version counter).
- // Pass `walSequence - 1` as the exclusive `fromSequence` so replay
- // returns entries with seq >= walSequence — the doc's last write is
- // exactly at walSequence, so the very first hit is the one we want.
- // For docs whose seq is unknown (0 — e.g. loaded from a snapshot
- // that predates v1.7.4), fall back to the full-WAL scan from 0.
- store_->setDocumentLoadCallback([this](const std::string& collection,
- const std::string& id,
- uint64_t walSequence) -> std::optional<Document> {
- uint64_t fromSequence = walSequence > 0 ? walSequence - 1 : 0;
- return persistence_->loadDocument(collection, id, fromSequence);
- });
- // Create encryption manager
- EncryptionManager::Config encryptConfig;
- encryptConfig.enabled = config_.encryptionEnabled;
- encryptConfig.keyFilePath = config_.keyFilePath;
- encryptConfig.autoGenerateKey = config_.autoGenerateKey;
- encryption_ = std::make_unique<EncryptionManager>(encryptConfig);
- // Create event manager
- EventManager::Config eventConfig;
- eventConfig.nodeId = config_.nodeId;
- events_ = std::make_unique<EventManager>(eventConfig);
- // Create file manager
- FileManager::Config fileConfig;
- fileConfig.filesDir = config_.dataDirectory / "files";
- fileConfig.maxFileSizeMb = config_.maxFileSizeMb;
- fileConfig.allowedMimeTypes = config_.allowedFileTypes;
- fileConfig.cleanupIntervalSec = config_.fileCleanupIntervalSec;
- fileConfig.defaultTtlSecondsByType = config_.fileDefaultTtlByType;
- files_ = std::make_unique<FileManager>(fileConfig);
- // Create replication manager
- ReplicationManager::Config replConfig;
- replConfig.nodeId = config_.nodeId;
- replConfig.peerAddresses = config_.peerAddresses;
- replConfig.conflictResolution = config_.conflictResolution;
- replication_ = std::make_unique<ReplicationManager>(replConfig);
- // Wire replication entry apply callback
- replication_->setEntryApplyCallback([this](const databasepb::ReplicationEntry& entry) {
- applyReplicatedEntry(entry);
- });
- // Create gRPC implementations
- storageImpl_ = std::make_unique<DatabaseGrpcImpl>(
- *this, *store_, *persistence_, *events_, *files_, *encryption_, *view_manager_, *config_manager_,
- *policy_manager_, *relation_manager_
- );
- // v2.7.0 — hand the impl the union of every listener's keys so it can
- // resolve a principal per call. A token is only accepted if some listener's
- // auth processor accepted it, so the union cannot widen authentication; it
- // only names what was already authenticated.
- {
- std::vector<smartbotic::database::auth::PrincipalResolver::NamedKey> all;
- for (const auto& l : config_.listeners) {
- for (const auto& k : l.auth.keys) all.push_back({k.name, k.key});
- }
- storageImpl_->setPrincipalKeys(std::move(all));
- }
- // v1.6.2 — wire the streaming-RPC concurrency limits from GrpcConfig.
- storageImpl_->setStreamLimits(
- config_.grpc.maxConcurrentSubscribeStreams,
- config_.grpc.maxConcurrentFileStreams);
- replicationImpl_ = std::make_unique<DatabaseReplicationGrpcImpl>(
- *store_, *replication_, *persistence_
- );
- }
- void DatabaseService::notifyReplicationAndEvents(const std::string& collection,
- const std::string& id,
- const std::optional<Document>& doc,
- EventType eventType) {
- // Queue for replication broadcast
- if (replication_ && config_.replicationEnabled) {
- databasepb::ReplicationEntry entry;
- entry.set_collection(collection);
- entry.set_document_id(id);
- // v1.8.0 — stamp our nodeId so peers can attribute the write to
- // its actual origin. Pre-v1.8 broadcasts left this empty, which
- // is now treated by receivers as "legacy / unknown origin".
- entry.set_node_id(config_.nodeId);
- entry.set_global_timestamp(
- std::chrono::duration_cast<std::chrono::milliseconds>(
- std::chrono::system_clock::now().time_since_epoch()
- ).count());
- switch (eventType) {
- case EventType::INSERT:
- entry.set_op(databasepb::OP_INSERT);
- // Send the full Document envelope (id, collection, data,
- // version, timestamps, encryption state). The follower's
- // applyReplicatedEntry uses Document::fromJson() which
- // expects this envelope — sending only doc->data silently
- // dropped every user field on the other side.
- if (doc) entry.set_data(doc->toJson().dump());
- break;
- case EventType::UPDATE:
- entry.set_op(databasepb::OP_UPDATE);
- if (doc) entry.set_data(doc->toJson().dump());
- break;
- case EventType::DELETE:
- entry.set_op(databasepb::OP_DELETE);
- break;
- default:
- break;
- }
- replication_->queueForReplication(entry);
- }
- // Publish event
- if (events_) {
- DatabaseEvent event;
- event.type = eventType;
- event.collection = collection;
- event.documentId = id;
- event.timestamp = std::chrono::duration_cast<std::chrono::milliseconds>(
- std::chrono::system_clock::now().time_since_epoch()
- ).count();
- event.nodeId = config_.nodeId;
- if (doc) {
- event.data = doc->data();
- }
- events_->publish(event);
- }
- }
- void DatabaseService::applyReplicatedEntry(const databasepb::ReplicationEntry& entry) {
- try {
- // v1.8.0 — also append the replicated entry to the local WAL, tagged
- // with the originating node's id. This is what makes follower-side
- // recovery (eviction → WAL fallback, restart → WAL replay) preserve
- // replicated documents. Echo amplification is prevented by:
- // - GetEntriesSince emitting walEntry.nodeId (not local node id)
- // - SyncProtocol skipping entries whose origin is the peer itself
- // both implemented in the same v1.8.0 cycle.
- const std::string origin = entry.node_id().empty()
- ? config_.nodeId // legacy peer (pre-1.8) — best-guess local
- : entry.node_id();
- switch (entry.op()) {
- case databasepb::OP_INSERT:
- case databasepb::OP_UPDATE:
- case databasepb::OP_UPSERT: {
- if (entry.data().empty()) {
- spdlog::warn("Replication entry has no data for op {}", static_cast<int>(entry.op()));
- return;
- }
- auto json = smartbotic::db::parse_to_nlohmann(entry.data());
- Document doc = Document::fromJson(json);
- doc.id = entry.document_id();
- doc.nodeId = entry.node_id();
- // loadDocument bypasses the persist-and-broadcast callback to avoid
- // re-replicating; we explicitly write to the WAL right after with
- // the origin-aware overload so eviction/WAL recovery still works.
- store_->loadDocument(entry.collection(), doc);
- // v2.3.1 — also mirror to LMDB. loadDocument bypasses the
- // MemoryStore persist callback that normally drives the
- // dual-write mirror; in v2.0-v2.2 the mirror was a side
- // channel and reads still hit MemoryStore, so the bypass
- // was harmless. In v2.3 reads are LMDB-first, so a
- // follower that only filled MemoryStore would return
- // empty Find/Get for every replicated doc. Route the
- // entry through the same project-aware mirror the write
- // handlers use.
- if (projects_) {
- auto pc = smartbotic::database::parseProjectCollection(
- entry.collection());
- if (auto* ds = projects_->getOrCreate(pc.project)) {
- std::optional<Document> opt_doc(doc);
- try {
- smartbotic::db::storage::applyDualWriteMirror(
- ds, mirror_healthy_, mirror_drift_count_,
- pc.collection, doc.id, opt_doc, EventType::INSERT);
- } catch (const smartbotic::db::storage::UniqueViolation& e) {
- // v2.11.0 T11 finding 1 — deliberate choice, not an
- // accident of the generic catch below. By this
- // point store_->loadDocument() above has ALREADY
- // put the row into MemoryStore (loadDocument
- // bypasses the persist callback specifically to
- // avoid re-replicating), so unlike every
- // client-facing write path there is no "reject the
- // whole operation" available: the origin node
- // already committed this write and every other
- // follower is expected to converge to it too.
- // Undoing loadDocument here would silently diverge
- // this follower's data from the rest of the
- // cluster over a constraint that may not even
- // exist on the origin (this follower could have
- // declared the unique field locally, after the
- // fact — replication carries no guarantee the
- // origin shares this node's config). So
- // MemoryStore keeps the row — it is the true
- // record of what replication decided — and the
- // mirror is marked unhealthy ON PURPOSE: this is
- // exactly the pre-T11 swallow-and-flip outcome,
- // chosen deliberately for this one caller so an
- // operator sees the drift (and LMDB-first reads
- // fall back to MemoryStore, which has the answer)
- // rather than nothing signalling a permanent gap
- // between the two stores.
- spdlog::error(
- "v2.11 replication: unique constraint violated applying "
- "{}/{}: {} - MemoryStore holds the row, LMDB does not; "
- "marking the mirror unhealthy so reads fall back instead "
- "of silently disagreeing with MemoryStore",
- entry.collection(), doc.id, e.what());
- mirror_healthy_.store(false, std::memory_order_release);
- mirror_drift_count_.fetch_add(1, std::memory_order_relaxed);
- } catch (const smartbotic::db::storage::MissingParentReference& e) {
- // v2.11.0 T13 — same deliberate swallow-and-flip as
- // the UniqueViolation catch just above, for the same
- // reason: loadDocument() has already committed this
- // row into MemoryStore, the origin already accepted
- // the write, and this follower may not even have the
- // same validate_on_write declaration (or the same
- // parent data) as the origin did at the time it
- // wrote. MemoryStore keeps the row; the mirror is
- // marked unhealthy on purpose so an operator sees
- // the drift and reads fall back to MemoryStore
- // rather than the two stores silently disagreeing.
- spdlog::error(
- "v2.11 replication: validate_on_write rejected applying "
- "{}/{}: {} - MemoryStore holds the row, LMDB does not; "
- "marking the mirror unhealthy so reads fall back instead "
- "of silently disagreeing with MemoryStore",
- entry.collection(), doc.id, e.what());
- mirror_healthy_.store(false, std::memory_order_release);
- mirror_drift_count_.fetch_add(1, std::memory_order_relaxed);
- }
- }
- }
- if (persistence_) {
- uint64_t walSeq = (entry.op() == databasepb::OP_INSERT)
- ? persistence_->logInsert(entry.collection(), doc, origin)
- : persistence_->logUpdate(entry.collection(), doc, origin);
- if (walSeq > 0) {
- store_->recordWalSequence(entry.collection(), doc.id, walSeq);
- }
- }
- spdlog::trace("Applied replicated {} to {}/{} from {}",
- entry.op() == databasepb::OP_INSERT ? "insert" :
- entry.op() == databasepb::OP_UPDATE ? "update" : "upsert",
- entry.collection(), entry.document_id(), entry.node_id());
- break;
- }
- case databasepb::OP_DELETE: {
- store_->remove(entry.collection(), entry.document_id());
- if (persistence_) {
- persistence_->logDelete(entry.collection(), entry.document_id(), origin);
- }
- spdlog::trace("Applied replicated delete to {}/{} from {}",
- entry.collection(), entry.document_id(), entry.node_id());
- break;
- }
- case databasepb::OP_CREATE_COLLECTION: {
- CollectionOptions options;
- store_->createCollection(entry.collection(), options);
- if (persistence_) {
- persistence_->logCreateCollection(entry.collection(), options, origin);
- }
- spdlog::trace("Applied replicated create collection {} from {}",
- entry.collection(), entry.node_id());
- break;
- }
- case databasepb::OP_DROP_COLLECTION: {
- store_->dropCollection(entry.collection());
- if (persistence_) {
- persistence_->logDropCollection(entry.collection(), origin);
- }
- spdlog::trace("Applied replicated drop collection {} from {}",
- entry.collection(), entry.node_id());
- break;
- }
- default:
- spdlog::warn("Unknown replication operation type: {}", static_cast<int>(entry.op()));
- break;
- }
- } catch (const std::exception& e) {
- spdlog::error("Failed to apply replicated entry: {}", e.what());
- }
- }
- void DatabaseService::startGrpcServer() {
- if (config_.listeners.empty()) {
- // Should never hit — parseConfig's back-compat shim always
- // synthesises at least one. Defend anyway.
- throw std::runtime_error("startGrpcServer: no listeners configured");
- }
- const int64_t mb = 1024 * 1024;
- spdlog::info("gRPC config: recv={}MB send={}MB quota={}MB subs={} files={}",
- config_.grpc.maxReceiveMessageSizeMb,
- config_.grpc.maxSendMessageSizeMb,
- config_.grpc.resourceQuotaMemoryMb,
- config_.grpc.maxConcurrentSubscribeStreams,
- config_.grpc.maxConcurrentFileStreams);
- grpcServers_.reserve(config_.listeners.size());
- for (const auto& listener : config_.listeners) {
- const std::string addr = listener.bind + ":" + std::to_string(listener.port);
- grpc::ServerBuilder builder;
- // v2.4 — choose credentials per listener. Plaintext for
- // tls.enabled=false; SslServerCredentials with the operator's
- // cert+key (or an auto-generated self-signed pair) for true.
- std::shared_ptr<grpc::ServerCredentials> creds;
- if (listener.tls.enabled) {
- std::filesystem::path cert_path = listener.tls.cert_path;
- std::filesystem::path key_path = listener.tls.key_path;
- // Resolve cert+key. If either is missing/unreadable AND
- // auto_self_signed_if_missing is set, generate a pair.
- const bool cert_ok = !cert_path.empty()
- && std::filesystem::exists(cert_path);
- const bool key_ok = !key_path.empty()
- && std::filesystem::exists(key_path);
- if (!cert_ok || !key_ok) {
- if (!listener.tls.auto_self_signed_if_missing) {
- throw std::runtime_error(
- "TLS listener " + addr +
- ": cert or key missing and auto_self_signed_if_missing=false");
- }
- const auto out_dir = config_.dataDirectory / "tls";
- auto gen = smartbotic::database::tls_util::generateSelfSignedCert(
- out_dir, listener.bind);
- cert_path = gen.cert_path;
- key_path = gen.key_path;
- spdlog::warn(
- "TLS listener {}: no operator cert at '{}' (or key at '{}'); "
- "auto-generated self-signed cert at '{}' (key '{}'). "
- "OK for dev. For production, drop a real cert + key at "
- "the configured cert_path/key_path.",
- addr,
- listener.tls.cert_path.string(),
- listener.tls.key_path.string(),
- cert_path.string(),
- key_path.string());
- }
- auto slurp = [](const std::filesystem::path& p) {
- std::ifstream f(p);
- if (!f) throw std::runtime_error("read " + p.string());
- std::stringstream ss; ss << f.rdbuf();
- return ss.str();
- };
- grpc::SslServerCredentialsOptions sslOpts;
- grpc::SslServerCredentialsOptions::PemKeyCertPair kp;
- kp.private_key = slurp(key_path);
- kp.cert_chain = slurp(cert_path);
- sslOpts.pem_key_cert_pairs.push_back(std::move(kp));
- creds = grpc::SslServerCredentials(sslOpts);
- } else {
- creds = grpc::InsecureServerCredentials();
- }
- // v2.4 Stage D — attach the bearer-token auth processor when
- // the listener requires it. The processor validates the
- // `authorization: Bearer <key>` metadata against the listener's
- // configured keys list. gRPC rejects with UNAUTHENTICATED
- // (no per-handler code needed).
- if (listener.auth.required) {
- if (listener.auth.keys.empty()) {
- throw std::runtime_error(
- "Listener " + addr + ": auth.required=true but auth.keys is empty");
- }
- // v2.7.0 — a key name is a principal, so it must be unambiguous.
- // Refuse at startup rather than resolving a policy against a name
- // that means two different callers.
- {
- std::set<std::string> seen;
- for (const auto& k : listener.auth.keys) {
- if (k.name != smartbotic::database::auth::kPrincipalUnnamed &&
- smartbotic::database::auth::isReservedPrincipal(k.name)) {
- throw std::runtime_error(
- "Listener " + addr + ": auth key name '" + k.name +
- "' is reserved. `anonymous` denotes an unauthenticated "
- "caller and `unnamed` denotes a legacy bare key; a real "
- "key must not be able to impersonate either in a policy.");
- }
- if (!seen.insert(k.name).second &&
- k.name != smartbotic::database::auth::kPrincipalUnnamed) {
- throw std::runtime_error(
- "Listener " + addr + ": duplicate auth key name '" +
- k.name + "'. A principal must identify one caller, "
- "otherwise a policy written against it is ambiguous.");
- }
- }
- }
- if (!listener.tls.enabled) {
- // gRPC ABORTS the process if an auth metadata processor is
- // attached to insecure credentials
- // (insecure_server_credentials.cc: "assertion failed: 0").
- // This used to be a WARN, which meant a plaintext+auth listener
- // looked merely inadvisable in the config and then killed the
- // server on start with an unexplained assert. Refuse it here
- // with a message that says what to do instead.
- throw std::runtime_error(
- "Listener " + addr + ": auth.required=true requires "
- "tls.enabled=true. gRPC does not support an auth metadata "
- "processor on insecure credentials and will abort the "
- "process. Enable TLS for this listener (set "
- "tls.auto_self_signed_if_missing for a local dev cert), or "
- "drop auth and rely on binding to loopback.");
- }
- std::vector<smartbotic::database::auth::BearerAuthProcessor::NamedKey> nk;
- nk.reserve(listener.auth.keys.size());
- for (const auto& k : listener.auth.keys) nk.push_back({k.name, k.key});
- creds->SetAuthMetadataProcessor(
- std::make_shared<smartbotic::database::auth::BearerAuthProcessor>(
- std::move(nk)));
- }
- builder.AddListeningPort(addr, creds);
- builder.RegisterService(storageImpl_.get());
- builder.RegisterService(replicationImpl_.get());
- builder.SetMaxReceiveMessageSize(
- static_cast<int>(config_.grpc.maxReceiveMessageSizeMb * mb));
- builder.SetMaxSendMessageSize(
- static_cast<int>(config_.grpc.maxSendMessageSizeMb * mb));
- grpc::ResourceQuota quota("smartbotic-db-" + addr);
- quota.Resize(config_.grpc.resourceQuotaMemoryMb * mb);
- builder.SetResourceQuota(quota);
- auto server = builder.BuildAndStart();
- if (!server) {
- throw std::runtime_error("Failed to start gRPC server on " + addr);
- }
- // v2.4 Stage B log line — listener identity printed so operators
- // know exactly which interface is exposed with which policy.
- // Stage C/D will extend with TLS + auth status.
- spdlog::info("gRPC listener up: {} (tls={}, auth_required={})",
- addr,
- listener.tls.enabled ? "true" : "false",
- listener.auth.required ? "true" : "false");
- grpcServers_.push_back(std::move(server));
- }
- }
- void DatabaseService::stopGrpcServer() {
- // ⚠ Shutdown() WITHOUT a deadline waits for every in-flight RPC to return,
- // and it does NOT cancel them. `Subscribe` never returns on its own: its
- // handler parks in `while (!context->IsCancelled())`, and that context is
- // cancelled only when the client goes away or a shutdown DEADLINE expires.
- // So a single connected subscriber - the normal state of any consumer using
- // events - hung this call forever.
- //
- // The cost is not just a slow stop. `stop()` calls this FIRST and takes the
- // final snapshot AFTER it, so the hang blocked the snapshot; systemd then
- // SIGKILLed at TimeoutStopSec=300, and the next boot recovered by WAL-only
- // replay, which trips auto read-only mode and needs an operator `unlock`.
- // A graceful restart therefore looked like a crash to the next boot.
- //
- // With a deadline, gRPC cancels whatever is still running once it expires,
- // the Subscribe loops observe IsCancelled() and unwind, and shutdown
- // proceeds. The deadline is ABSOLUTE and shared across listeners, so the
- // total wait is bounded by the grace period rather than by
- // grace × listener count.
- //
- // Five seconds: long enough for an ordinary unary RPC in flight to finish
- // (they are milliseconds), short enough that it is invisible next to the
- // final snapshot, which is the part of shutdown that legitimately takes
- // time (~70 s on the largest known dataset).
- constexpr auto kShutdownGrace = std::chrono::seconds(5);
- const auto deadline = std::chrono::system_clock::now() + kShutdownGrace;
- for (auto& server : grpcServers_) {
- if (server) server->Shutdown(deadline);
- }
- grpcServers_.clear();
- }
- } // namespace smartbotic::database
|