database_service.cpp 95 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931
  1. #include "database_service.hpp"
  2. #include "json_parse.hpp"
  3. #include "project_addressing.hpp"
  4. #include "storage/document_store.hpp"
  5. #include "storage/document_store_lmdb.hpp"
  6. #include "storage/dual_write_mirror.hpp"
  7. #include "relations/relation_cascade.hpp"
  8. #include "relations/relation_enforcement.hpp"
  9. #include "auth/auth_interceptor.hpp"
  10. #include "auth/principal.hpp"
  11. #include "tls/cert_generator.hpp"
  12. #include <fstream>
  13. #include <sstream>
  14. #include <map>
  15. #include <set>
  16. #include "storage/lmdb_env.hpp"
  17. #include "storage/migrate_v1_to_v2.hpp"
  18. #include "storage/project_store.hpp"
  19. #include "storage/subdb_placement.hpp"
  20. #include <grpcpp/grpcpp.h>
  21. #include <grpcpp/resource_quota.h>
  22. #include <spdlog/spdlog.h>
  23. #include <fstream>
  24. #ifdef HAVE_SYSTEMD
  25. #include <systemd/sd-daemon.h>
  26. #endif
  27. #if defined(__GLIBC__) && !defined(__APPLE__)
  28. #include <malloc.h>
  29. #endif
  30. namespace smartbotic::database {
  31. DatabaseService::DatabaseService(Config config)
  32. : config_(std::move(config))
  33. {
  34. }
  35. DatabaseService::~DatabaseService() {
  36. stop();
  37. }
  38. std::string DatabaseService::readOnlyReason() const {
  39. std::lock_guard<std::mutex> lock(reason_mutex_);
  40. return read_only_reason_;
  41. }
  42. void DatabaseService::setReadOnly(bool value, const std::string& reason) {
  43. {
  44. std::lock_guard<std::mutex> lock(reason_mutex_);
  45. read_only_reason_ = value ? reason : "";
  46. }
  47. bool prev = read_only_.exchange(value, std::memory_order_acq_rel);
  48. if (prev != value) {
  49. if (value) {
  50. spdlog::error("Database entered READ-ONLY mode: {}", reason);
  51. } else {
  52. spdlog::info("Database unlocked -- writes accepted");
  53. }
  54. }
  55. }
  56. bool DatabaseService::initialize() {
  57. spdlog::info("Initializing database service (node: {})", config_.nodeId);
  58. try {
  59. // Create data directory if needed
  60. std::error_code ec;
  61. std::filesystem::create_directories(config_.dataDirectory, ec);
  62. if (ec) {
  63. spdlog::error("Failed to create data directory: {}", ec.message());
  64. return false;
  65. }
  66. setupComponents();
  67. // Initialize encryption
  68. if (config_.encryptionEnabled) {
  69. if (!encryption_->initialize()) {
  70. spdlog::error("Failed to initialize encryption");
  71. return false;
  72. }
  73. }
  74. // Recover from persistence
  75. recovery_outcome_ = persistence_->recover(*store_);
  76. #if defined(__GLIBC__) && !defined(__APPLE__)
  77. // v1.9.3 — release freelist pages accumulated during recovery. v1.9.1
  78. // added trim at end of loadSnapshot, but the WAL replay phase that
  79. // runs afterward (`persistence_->recover`'s `replayWal` step) also
  80. // burns through GBs of small Document JSON allocations whose pages
  81. // sit on the per-thread freelist with no subsequent allocation to
  82. // shake them loose. On Zoe (docs/incidents/2026-04-22-zoe-rss-exceeds-budget.md update 15:05)
  83. // this was 5.7 GB at 23 min uptime — fully reclaimable via
  84. // malloc_trim, just nothing called it. Paired with the periodic
  85. // every-5-minute trim in MemoryStore::logMemoryCheck, so any
  86. // post-boot bloat that escapes this trim gets cleaned up shortly.
  87. ::malloc_trim(0);
  88. #endif
  89. if (recovery_outcome_.isFailure()) {
  90. const std::string modeStr =
  91. recoveryModeToString(config_.persistenceConfig.recoveryMode);
  92. spdlog::error("");
  93. spdlog::error("+------------------------------------------------------------------+");
  94. spdlog::error("| RECOVERY REFUSED |");
  95. spdlog::error("| |");
  96. spdlog::error("| Mode: {}", modeStr);
  97. spdlog::error("| Expected snapshot: {}", recovery_outcome_.expectedSnapshot.string());
  98. spdlog::error("| Reason: {}", recovery_outcome_.failureReason);
  99. spdlog::error("| |");
  100. spdlog::error("| Snapshots available: {}", recovery_outcome_.snapshotsAvailable);
  101. spdlog::error("| Snapshots attempted: {}", recovery_outcome_.snapshotsAttempted);
  102. spdlog::error("| |");
  103. spdlog::error("| To escalate, restart with one of: |");
  104. spdlog::error("| --recovery-mode=snapshot_fallback Try older snapshots |");
  105. spdlog::error("| --recovery-mode=wal_only Replay WAL only (slow) |");
  106. spdlog::error("| --recovery-mode=best_effort Try all of the above |");
  107. spdlog::error("| --recovery-mode=force_empty Start empty (LAST RESORT) |");
  108. spdlog::error("| |");
  109. spdlog::error("| Or set in config.json: \"recovery\": {{ \"mode\": \"<mode>\" }} |");
  110. spdlog::error("| |");
  111. spdlog::error("| PRESERVE /var/lib/smartbotic-database/ BEFORE ESCALATING. |");
  112. spdlog::error("+------------------------------------------------------------------+");
  113. throw std::runtime_error("recovery failed; refusing to start");
  114. }
  115. // Auto-readonly mode on non-trivial recovery
  116. if (recovery_outcome_.isNonTrivial() && !force_readwrite_) {
  117. std::string reason;
  118. switch (recovery_outcome_.kind) {
  119. case RecoveryOutcome::Kind::SnapshotFellBack:
  120. reason = "fell back to snapshot " +
  121. recovery_outcome_.snapshotUsed.filename().string() +
  122. " because " +
  123. recovery_outcome_.expectedSnapshot.filename().string() +
  124. " failed: " + recovery_outcome_.failureReason;
  125. break;
  126. case RecoveryOutcome::Kind::WalOnlyReplay:
  127. reason = "WAL-only replay, no snapshot loaded (" +
  128. std::to_string(recovery_outcome_.walEntriesReplayed) +
  129. " entries)";
  130. break;
  131. case RecoveryOutcome::Kind::ForcedEmpty:
  132. reason = "forced empty by operator (--recovery-mode=force_empty)";
  133. break;
  134. default:
  135. reason = "non-trivial recovery";
  136. break;
  137. }
  138. reason += ". Run `smartbotic-db-cli unlock` to accept this state, "
  139. "or restart with --force-readwrite to bypass this check.";
  140. setReadOnly(true, reason);
  141. spdlog::error("");
  142. spdlog::error("+------------------------------------------------------------------+");
  143. spdlog::error("| [ERROR] Database booted in READ-ONLY mode after non-trivial |");
  144. spdlog::error("| recovery. |");
  145. spdlog::error("| |");
  146. spdlog::error("| Reason: {}", reason);
  147. spdlog::error("| |");
  148. spdlog::error("| Writes will be REJECTED until you acknowledge this state: |");
  149. spdlog::error("| smartbotic-db-cli unlock # live, no restart |");
  150. spdlog::error("| smartbotic-database --force-readwrite # on next restart |");
  151. spdlog::error("+------------------------------------------------------------------+");
  152. }
  153. // Set replication sequence after recovery (WAL sequence is now known)
  154. replication_->setSequence(persistence_->currentWalSequence());
  155. // v1.8.0 — start the persistence manager before loadFromStore /
  156. // migrations so any system-collection writes those phases produce
  157. // (view docs, _collection_meta entries, _migrations recordings)
  158. // reach the WAL like normal mutations. Pre-v1.8, persistence_ was
  159. // started later in start(), which silently dropped those writes
  160. // (running_=false → logInsert no-op). That meant migrated views
  161. // only "survived" because the runner re-created them every boot,
  162. // and a manually-created view that landed during the racy startup
  163. // window (before full readiness) could be lost.
  164. if (!persistence_->start()) {
  165. spdlog::error("Failed to start persistence manager before migrations");
  166. return false;
  167. }
  168. // Load view definitions from the _views system collection (which is now
  169. // populated by the persistence recovery above).
  170. view_manager_->loadFromStore();
  171. // Load per-collection configs from the _collection_meta system collection.
  172. // Must happen AFTER persistence recovery and BEFORE migrations so any
  173. // migration-created documents are stamped with the correct precision.
  174. config_manager_->loadFromStore();
  175. policy_manager_->loadFromStore();
  176. relation_manager_->loadFromStore();
  177. applyRelationDeclarations();
  178. // v2.11.0 close-out — teach the TTL sweeper about relations. Installed
  179. // HERE, after loadFromStore() and the arming pass, deliberately: the
  180. // hook consults the declaration cache and the reverse index, and an
  181. // expiry that fired before either was ready would have decided from an
  182. // empty cache. Until this line the sweeper behaves exactly as it did
  183. // pre-v2.11.0 (which is also what every unit fixture that never
  184. // installs the hook keeps doing).
  185. store_->setTtlExpiryRelationHook(
  186. [this](const std::string& qualifiedCollection, const std::string& id,
  187. bool firstAttempt) {
  188. return ttlExpiryRelationDecision(qualifiedCollection, id, firstAttempt);
  189. });
  190. // v2.9.0 — re-apply persisted index declarations to each project's LMDB
  191. // store. This is load-bearing, not bookkeeping: the declaration is what
  192. // makes the write path maintain an index, and the planner consults an
  193. // index purely on the declaration's word. If this step were skipped, a
  194. // restart would leave indexes recorded but unmaintained, and queries
  195. // would be served from a frozen index - stale rows returned as current,
  196. // with nothing logged.
  197. applyIndexDeclarations();
  198. // v2.11.0 final review (finding 4) — THE POST-REPLAY RE-MIRROR PASS
  199. // RUNS HERE, and the position is load-bearing.
  200. //
  201. // PersistenceManager::recover() (above, before any of the arming
  202. // steps) only collects the id list. The pass writes documents through
  203. // LmdbDocumentStore::put(), which reads indexed_fields(),
  204. // unique_fields() and relations() LIVE and maintains every index
  205. // inside the document's own write transaction. Run from inside
  206. // recover(), all three maps were still empty, so it repaired the
  207. // document and left every posting derived from it stale: an
  208. // index-served query then returns the row under its OLD value and
  209. // misses it under its new one - silently wrong rows, which is exactly
  210. // what putting maintainIndexes() inside the transaction exists to
  211. // prevent, and which bit any install with a v2.9 index declared
  212. // whether or not it uses relations. It also made the reverse-index
  213. // half of relation_cascade.hpp's convergence claim false, and made
  214. // UniqueViolation (and therefore the pass's own retry phase)
  215. // unreachable.
  216. //
  217. // Must be after applyRelationDeclarations() AND
  218. // applyIndexDeclarations(); backfillIntoDocStore() below already sits
  219. // at this point for the same reason. Never throws.
  220. persistence_->runPendingRemirror(*store_, recovery_outcome_);
  221. // Run migrations if enabled
  222. if (config_.migrations.enabled && !config_.migrations.directory.empty()) {
  223. if (!runMigrations()) {
  224. if (config_.migrations.failOnError) {
  225. spdlog::error("Failed to run migrations");
  226. return false;
  227. }
  228. spdlog::warn("Some migrations failed, continuing anyway");
  229. }
  230. }
  231. // v2.0 Stage 4 — synchronous backfill into doc_store_ before
  232. // serving. Without this, only post-boot writes land in LMDB and
  233. // the read-flip downstream would see an empty mirror.
  234. backfillIntoDocStore();
  235. // v2.4.4 — placement audit. Reports documents sitting in a sub-db
  236. // other than the one they declare, which is the signature of a
  237. // misbound MDB_dbi (see storage/subdb_identity.hpp). Read-only and
  238. // advisory: it never blocks startup, because a misplacement is a
  239. // data-location problem an operator repairs offline with
  240. // `smartbotic-db-cli reconcile-subdbs`, not a reason to refuse
  241. // service on the other 99% of the dataset.
  242. auditSubdbPlacement();
  243. // v2.6.0 — stamp `project: "default"` onto file records written before
  244. // files were namespaced. Idempotent, and advisory: a file-metadata
  245. // stamp is not a reason to refuse service, so failures warn only.
  246. if (files_) {
  247. try {
  248. const uint32_t stamped = files_->stampMissingProjects();
  249. if (stamped > 0) {
  250. spdlog::info("v2.6.0 file migration: stamped {} legacy file "
  251. "record(s) with project 'default'", stamped);
  252. } else {
  253. spdlog::debug("v2.6.0 file migration: no legacy file records");
  254. }
  255. } catch (const std::exception& e) {
  256. spdlog::warn("v2.6.0 file migration failed: {}", e.what());
  257. }
  258. }
  259. // v2.11.0 close-out — LAST ACT OF THE BOOT PATH, and the position is
  260. // load-bearing. Everything above that can bump the drift counter (the
  261. // post-replay re-mirror pass, backfillIntoDocStore(), a malformed
  262. // legacy collection key) is now behind the baseline, so
  263. // mirrorDriftSinceReady() reports only drift a LIVE write caused.
  264. // Destructive relation policies gate on that, not on the raw counter:
  265. // see MemoryStore::markMirrorDriftBaseline() for the trap this closes
  266. // (one unrepairable row disabling every cascade/set_null delete and
  267. // every CreateRelation for the process's life, with an error message
  268. // telling the operator to restart, which re-incurs it). Reads keep
  269. // gating on the RAW counter - a stale row genuinely means MemoryStore
  270. // is ahead, and that fallback must stay.
  271. store_->markMirrorDriftBaseline();
  272. if (store_->mirrorDriftBaseline() > 0) {
  273. spdlog::warn("mirror drift at READY: {} - accrued on the boot path (see the "
  274. "ERROR lines above for the exact rows). Reads for this process "
  275. "will be served from MemoryStore rather than LMDB. Destructive "
  276. "relation policies are NOT disabled by this, but the rows named "
  277. "above are still stale in LMDB until each is rewritten through "
  278. "an ordinary write; a restart re-attempts and, if the cause "
  279. "persists, re-incurs the same count.",
  280. store_->mirrorDriftBaseline());
  281. }
  282. spdlog::info("Database service initialized successfully");
  283. return true;
  284. } catch (const std::exception& e) {
  285. spdlog::error("Failed to initialize database service: {}", e.what());
  286. return false;
  287. }
  288. }
  289. void DatabaseService::migrateLegacyEnvToDefaultProject() {
  290. const auto legacy = config_.dataDirectory / "env";
  291. const auto target = config_.dataDirectory / "projects" /
  292. smartbotic::database::kDefaultProject / "env";
  293. std::error_code ec;
  294. const bool legacy_exists = std::filesystem::exists(legacy, ec);
  295. const bool target_exists = std::filesystem::exists(target, ec);
  296. if (!legacy_exists) return; // fresh install or already migrated
  297. if (target_exists) {
  298. spdlog::error("v2.3 storage: both legacy '{}' and new '{}' exist — "
  299. "refusing to start. Inspect manually; the safe move is "
  300. "to either delete the legacy dir (if you confirm it's a "
  301. "leftover) or stop and contact ops.",
  302. legacy.string(), target.string());
  303. throw std::runtime_error("v2.3 storage migration: ambiguous layout");
  304. }
  305. std::filesystem::create_directories(target.parent_path(), ec);
  306. if (ec) {
  307. throw std::runtime_error("v2.3 storage migration: cannot create '"
  308. + target.parent_path().string() + "': " + ec.message());
  309. }
  310. std::filesystem::rename(legacy, target, ec);
  311. if (ec) {
  312. throw std::runtime_error("v2.3 storage migration: rename '"
  313. + legacy.string() + "' -> '"
  314. + target.string() + "' failed: " + ec.message());
  315. }
  316. spdlog::info("v2.3 storage: migrated legacy env '{}' -> '{}' (default project)",
  317. legacy.string(), target.string());
  318. }
  319. smartbotic::db::storage::DocumentStore* DatabaseService::docStore() noexcept {
  320. return docStore(smartbotic::database::kDefaultProject);
  321. }
  322. smartbotic::db::storage::DocumentStore*
  323. DatabaseService::docStore(std::string_view project) noexcept {
  324. if (!projects_) return nullptr;
  325. return projects_->get(project);
  326. }
  327. std::vector<std::string> DatabaseService::listProjects() const {
  328. if (!projects_) return {};
  329. return projects_->listOnDisk();
  330. }
  331. bool DatabaseService::createProject(const std::string& name, std::string& error) {
  332. if (!projects_) {
  333. error = "project registry not initialized";
  334. return false;
  335. }
  336. return projects_->create(name, error);
  337. }
  338. bool DatabaseService::dropProject(const std::string& name, std::string& error) {
  339. if (!projects_) {
  340. error = "project registry not initialized";
  341. return false;
  342. }
  343. // v2.11.0 close-out — capture the relation declarations belonging to this
  344. // project BEFORE the drop. Relations are validated same-project at
  345. // creation (no LMDB transaction spans two envs), so `child` and `parent`
  346. // are always in the same project as `name` itself; listRelations(name)
  347. // therefore names every declaration the drop invalidates. Captured first
  348. // because listRelations needs the project to still be nameable and, more
  349. // importantly, because un-arming needs the store to still exist.
  350. std::vector<RelationInfo> doomed;
  351. if (relation_manager_) {
  352. try {
  353. doomed = relation_manager_->listRelations(name);
  354. } catch (const std::exception& e) {
  355. spdlog::warn("v2.11 relations: could not enumerate relations of project '{}' "
  356. "before dropping it ({}); its declarations may be left behind",
  357. name, e.what());
  358. }
  359. }
  360. if (!projects_->drop(name, error)) return false;
  361. // ---- The drop succeeded; the declarations now point at a namespace that
  362. // no longer exists. Sweep them out of the global `_relations` collection
  363. // and un-arm each affected child collection.
  364. //
  365. // Order (drop first, sweep second) is deliberate: the registry is the one
  366. // that refuses "default" and rejects unknown names, so sweeping first would
  367. // destroy declarations for a project whose drop then failed. The cost of
  368. // this order is that the per-project store is usually already gone by the
  369. // time we try to un-arm, which makes the un-arm a no-op - correct, since a
  370. // store that no longer exists cannot be maintaining a reverse index. It
  371. // still runs, because ProjectStoreRegistry may hold a live
  372. // LmdbDocumentStore for a re-created project of the same name.
  373. //
  374. // ⚠ Advisory, never fatal: the project IS dropped at this point, and a
  375. // leftover declaration is litter, not corruption (its child and parent
  376. // collections are both gone, so nothing can be enforced against them). A
  377. // failure here must not report the drop as failed.
  378. for (const auto& r : doomed) {
  379. try {
  380. std::string err;
  381. if (!relation_manager_->dropRelation(r.name, err)) {
  382. spdlog::warn("v2.11 relations: could not drop declaration '{}' left by "
  383. "dropped project '{}': {}", r.name, name, err);
  384. continue;
  385. }
  386. const auto rc = resolveCollection(r.child);
  387. auto* lmdb = dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(
  388. docStore(rc.project));
  389. if (lmdb != nullptr) {
  390. // Empty list = "this collection has no relations", which is what
  391. // set_relations replaces the previous list with. Every relation
  392. // of a dropped project goes, so clearing the child wholesale is
  393. // exactly right here (unlike DropRelation, which must re-arm the
  394. // survivors).
  395. lmdb->set_relations(rc.collection, {});
  396. }
  397. spdlog::info("v2.11 relations: dropped declaration '{}' with the project '{}'",
  398. r.name, name);
  399. } catch (const std::exception& e) {
  400. spdlog::warn("v2.11 relations: could not clean up declaration '{}' after "
  401. "dropping project '{}': {}", r.name, name, e.what());
  402. }
  403. }
  404. // ⚠ KNOWN, NOT FIXED HERE (v2.11.0 close-out): _views, _policies and
  405. // _collection_meta have the SAME gap - ProjectStoreRegistry::drop() removes
  406. // the project's env and nothing else, so a dropped project leaves its view
  407. // definitions, access policies and per-collection configs behind in those
  408. // global system collections too. Not a shared fix: each manager keys and
  409. // caches differently (ViewManager keys `<project>:<name>`, PolicyManager
  410. // `<project>:<principal>`, CollectionConfigManager `<project>:<collection>`
  411. // and canonicalises the cache key only), so each needs its own sweep, and
  412. // PolicyManager's has a security dimension this one does not (a re-created
  413. // project inheriting a stale policy). Scoped out deliberately; recorded so
  414. // it is not rediscovered as new.
  415. return true;
  416. }
  417. bool DatabaseService::runMigrations() {
  418. if (!config_.migrations.enabled) {
  419. return true;
  420. }
  421. MigrationRunner::Config migrationConfig;
  422. migrationConfig.directory = config_.migrations.directory;
  423. migrationConfig.autoApply = config_.migrations.autoApply;
  424. migrationConfig.failOnError = config_.migrations.failOnError;
  425. migrationRunner_ = std::make_unique<MigrationRunner>(
  426. *store_, *view_manager_, *relation_manager_, migrationConfig);
  427. const bool ok = migrationRunner_->runMigrations();
  428. // v2.11.0 close-out — a `create_relation` migration op only PERSISTS the
  429. // declaration. Arming the write path and building the reverse index over
  430. // rows that already exist is what applyRelationDeclarations() does, and it
  431. // has already run for this boot (initialize() calls it before migrations),
  432. // so without this a migration-declared relation would maintain no reverse
  433. // index and enforce nothing until the NEXT restart - the "declared but not
  434. // enforcing" state this feature's self-heal exists to make unreachable.
  435. // Idempotent and cheap: re-arming an already-armed child replaces an
  436. // identical list, and the bootstrap loop skips any relation whose index
  437. // sub-db already exists. Runs even when the migration run reported failure,
  438. // because a partially-applied run may still have created a relation.
  439. applyRelationDeclarations(/*afterMigrations=*/true);
  440. return ok;
  441. }
  442. void DatabaseService::applyIndexDeclarations() {
  443. size_t applied = 0;
  444. for (const auto& [qualified, cfg] : config_manager_->allConfigs()) {
  445. if (cfg.indexedFields.empty()) continue;
  446. try {
  447. const auto rc = resolveCollection(qualified);
  448. auto* ds = docStore(rc.project);
  449. auto* lmdb =
  450. dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(ds);
  451. if (lmdb == nullptr) continue;
  452. lmdb->set_indexed_fields(rc.collection, cfg.indexedFields);
  453. // v2.11.0 T11 — arm uniqueFields the same way. This is the fix
  454. // for the exact failure mode the type was built to avoid:
  455. // uniqueFields has been persisted since T8 with no RPC or proto
  456. // surface, which was safe (nothing consulted it); now that
  457. // set_unique_fields is reachable, a restart that skipped this
  458. // call would silently stop enforcing a constraint every write
  459. // handler still advertises as active.
  460. lmdb->set_unique_fields(rc.collection, cfg.uniqueFields);
  461. ++applied;
  462. spdlog::info("v2.9 index: {} field(s) active on {} ({} unique)",
  463. cfg.indexedFields.size(), qualified, cfg.uniqueFields.size());
  464. // Self-heal a declared index whose sub-db is absent. That happens
  465. // when the KEY FORMAT VERSION in the sub-db prefix changes (v2.9.1
  466. // unified the numeric encoding, so `_idx_` became `_idx2_`), and it
  467. // would otherwise leave the declaration pointing at nothing: writes
  468. // would maintain the new index correctly, but rows written BEFORE the
  469. // upgrade would have no postings, so an indexed query would silently
  470. // return only the newer ones. Rebuilding is the only safe reading.
  471. for (const auto& field : cfg.indexedFields) {
  472. if (lmdb->index_stats(rc.collection, field).has_value()) continue;
  473. const uint64_t rows = lmdb->build_index(rc.collection, field);
  474. spdlog::warn("v2.9 index: rebuilt {}#{} over {} row(s) - the index "
  475. "was declared but its sub-db was absent (key format "
  476. "change or a manual removal)",
  477. qualified, field, rows);
  478. }
  479. } catch (const std::exception& e) {
  480. // Advisory per collection: one unparseable name must not stop the
  481. // rest from being armed. Loud, because a missing declaration means
  482. // that collection's index silently stops being maintained.
  483. spdlog::error("v2.9 index: could not apply declarations for {}: {}",
  484. qualified, e.what());
  485. }
  486. }
  487. if (applied > 0) {
  488. spdlog::info("v2.9 index: applied declarations for {} collection(s)", applied);
  489. }
  490. }
  491. MemoryStore::TtlExpiryAction DatabaseService::ttlExpiryRelationDecision(
  492. const std::string& qualifiedCollection, const std::string& id, bool firstAttempt) {
  493. using Action = MemoryStore::TtlExpiryAction;
  494. // Thin binder over relations/relation_cascade.cpp's ttlExpiryDecision() -
  495. // the policy itself lives there so it is unit-testable against the same
  496. // fixtures the cascade tests use, without a DatabaseService. All this does
  497. // is resolve the per-project LMDB store and supply the replication/events
  498. // notifier.
  499. if (!relation_manager_ || !config_manager_ || !store_ || !persistence_) {
  500. return Action::Proceed;
  501. }
  502. if (qualifiedCollection.empty() || qualifiedCollection[0] == '_') return Action::Proceed;
  503. try {
  504. const auto rc = resolveCollection(qualifiedCollection);
  505. auto* lmdb = dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(
  506. docStore(rc.project));
  507. // No LMDB substrate means no reverse index to consult, so the
  508. // pre-v2.11.0 behaviour is all that is available.
  509. if (lmdb == nullptr) return Action::Proceed;
  510. // Replication + Subscribe events, driven per mutation - the same wiring
  511. // DatabaseGrpcImpl::performRelationCascade() uses. Without it a
  512. // TTL-driven cascade would be invisible to followers and subscribers,
  513. // and a follower that never saw the child deletions diverges
  514. // permanently (the v2.3.1 class of bug).
  515. auto notify = [this](const std::string& coll, const std::string& docId,
  516. const std::optional<Document>& doc, EventType eventType) {
  517. notifyReplicationAndEvents(coll, docId, doc, eventType);
  518. };
  519. // firstAttempt -> the refusal (if any) is logged at WARN; a retry of an
  520. // already-reported document logs at DEBUG. The sweeper owns that
  521. // decision because it is the only thing that knows the document has
  522. // been refused before.
  523. return smartbotic::database::ttlExpiryDecision(
  524. *relation_manager_, *lmdb, *persistence_, *store_, *config_manager_,
  525. qualifiedCollection, id, notify, firstAttempt);
  526. } catch (const std::exception& e) {
  527. // Never expire a parent whose children could not be evaluated.
  528. spdlog::error("TTL expiry of '{}/{}': could not resolve its storage ({}) - the "
  529. "document is left in place and will be retried",
  530. qualifiedCollection, id, e.what());
  531. return Action::Skip;
  532. }
  533. }
  534. void DatabaseService::applyRelationDeclarations(bool afterMigrations) {
  535. // Group by (project, bare child collection) — set_relations() is a
  536. // per-collection call on that project's LmdbDocumentStore and replaces
  537. // whatever was declared before, so every relation sharing a child
  538. // collection must land in one call.
  539. std::unordered_map<std::string,
  540. std::vector<smartbotic::db::storage::RelationRef>> byChild;
  541. // Track which project each grouping key belongs to alongside the bare
  542. // collection name, since the map key alone doesn't carry it.
  543. std::unordered_map<std::string, std::pair<std::string, std::string>> keyToProjectCollection;
  544. for (const auto& r : relation_manager_->listRelations()) {
  545. try {
  546. const auto rn = resolveCollection(r.name);
  547. const auto rc = resolveCollection(r.child);
  548. // v2.11.0 T13 — parent resolved to its bare name, same reasoning
  549. // as armRelationsForChild's mirror of this construction.
  550. const auto pc = resolveCollection(r.parent);
  551. const std::string key = rc.project + ":" + rc.collection;
  552. // v2.11.0 T13 round 2 — snapshot relationsEnforced at boot-time
  553. // arming, same reasoning as armRelationsForChild's mirror of
  554. // this construction. `key` is already the qualified child name
  555. // configFor expects.
  556. const bool enforced = config_manager_->configFor(key).relationsEnforced;
  557. byChild[key].push_back(smartbotic::db::storage::RelationRef{
  558. rn.collection, r.childField, pc.collection, r.validateOnWrite, enforced});
  559. keyToProjectCollection[key] = {rc.project, rc.collection};
  560. } catch (const std::exception& e) {
  561. // Advisory per relation: one unparseable declaration must not
  562. // stop the rest from being armed. Loud, because a skipped
  563. // relation silently maintains no reverse index for its child.
  564. spdlog::error("v2.11 relations: could not apply declaration '{}': {}",
  565. r.name, e.what());
  566. }
  567. }
  568. size_t applied = 0;
  569. for (const auto& [key, refs] : byChild) {
  570. const auto& [project, collection] = keyToProjectCollection[key];
  571. try {
  572. // ⚠ v2.11.0 close-out — getOrCreate, not get(), on the POST-MIGRATION
  573. // pass. docStore() is projects_->get(), which does NOT create an env,
  574. // and a `create_relation` migration in a project that has never been
  575. // written to has no env yet: the declaration persisted, this pass
  576. // skipped it, and the relation maintained no reverse index and
  577. // enforced NOTHING until the next restart - the exact
  578. // "declared but not enforcing" state the op was added to avoid, just
  579. // narrowed to new projects. A migration declaring a relation IS a
  580. // statement that the project exists, and the mirror resolver would
  581. // getOrCreate the same env on the first write anyway.
  582. //
  583. // The BOOT pass deliberately keeps get(): creating envs there would
  584. // resurrect the env of a project whose declarations are merely stale.
  585. auto* ds = (afterMigrations && projects_)
  586. ? projects_->getOrCreate(project)
  587. : docStore(project);
  588. auto* lmdb = dynamic_cast<smartbotic::db::storage::LmdbDocumentStore*>(ds);
  589. if (lmdb == nullptr) continue;
  590. lmdb->set_relations(collection, refs);
  591. ++applied;
  592. spdlog::info("v2.11 relations: {} relation(s) active with child '{}:{}'",
  593. refs.size(), project, collection);
  594. // v2.11.0 T7 self-heal — mirrors applyIndexDeclarations()'s
  595. // rebuild of a declared-but-absent index sub-db, for the same
  596. // reason: a relation whose reverse index sub-db is missing is
  597. // exactly the "declared but not enforcing" bug this task exists
  598. // to close. Reachable paths that leave a relation in this state:
  599. // a snapshot/backup restored from before the sub-db existed, or
  600. // an operator manually dropping the `_relidx1_*` sub-db (there
  601. // is no dropRelation-only-the-index tool). CreateRelation's own
  602. // bootstrap scan (see database_grpc_impl.cpp) covers the normal
  603. // declare-time path; this covers everything else.
  604. //
  605. // ⚠ v2.11.0 close-out (review minor 2) — the SAME code path is the
  606. // NORMAL, expected one for a relation a migration just declared:
  607. // migrations run after this pass, so runMigrations() calls it again
  608. // and the brand-new relation legitimately has no sub-db yet. Logging
  609. // "declared but its sub-db was absent (restored snapshot or a manual
  610. // removal)" at WARN there describes a fault that did not happen.
  611. // `afterMigrations` distinguishes the two so the text and the level
  612. // match the event.
  613. for (const auto& ref : refs) {
  614. if (lmdb->relation_index_exists(ref.name)) continue;
  615. const uint64_t rows =
  616. lmdb->build_relation_index(ref.name, collection, ref.childField);
  617. if (afterMigrations) {
  618. spdlog::info("v2.11 relations: built the reverse index for '{}' over {} "
  619. "existing row(s) in '{}:{}' - a migration declared this "
  620. "relation on this boot",
  621. ref.name, rows, project, collection);
  622. } else {
  623. spdlog::warn("v2.11 relations: rebuilt '{}' over {} row(s) in '{}:{}' - the "
  624. "relation was declared but its reverse index sub-db was absent "
  625. "(restored snapshot predating it, or a manual removal)",
  626. ref.name, rows, project, collection);
  627. }
  628. }
  629. } catch (const std::exception& e) {
  630. spdlog::error("v2.11 relations: could not apply declarations for '{}:{}': {}",
  631. project, collection, e.what());
  632. }
  633. }
  634. if (applied > 0) {
  635. spdlog::info("v2.11 relations: applied declarations for {} child collection(s)", applied);
  636. }
  637. }
  638. void DatabaseService::auditSubdbPlacement() {
  639. if (!projects_) return;
  640. for (const auto& name : projects_->listOpen()) {
  641. auto h = projects_->getHandle(name);
  642. if (!h.env) continue;
  643. try {
  644. // audit_env(), NOT audit(). The path-taking overload opens a
  645. // second MDB_env on the same file; closing it would drop the
  646. // POSIX record locks this process already holds for `h.env`
  647. // (POSIX locks are per-process, and closing any fd on a file
  648. // releases all of them), leaving every later mdb_txn_begin
  649. // failing EINVAL. Borrow the open handle instead.
  650. const auto rep = smartbotic::db::storage::audit_env(h.env->raw(), name);
  651. if (rep.misplaced.empty()) {
  652. spdlog::debug("placement audit: project '{}' consistent ({} rows)",
  653. name, rep.rows_scanned);
  654. continue;
  655. }
  656. // Summarise per (physical -> declared home) pair; one line per
  657. // affected document would be unreadable at scale.
  658. std::map<std::string, uint64_t> pairs;
  659. for (const auto& m : rep.misplaced) {
  660. ++pairs[m.physical_subdb + " -> " + m.home_subdb];
  661. }
  662. spdlog::error("placement audit: project '{}' has {} document(s) in the "
  663. "wrong sub-db. Reads served from LMDB will not find them "
  664. "under their own collection. Repair offline with: "
  665. "smartbotic-db-cli reconcile-subdbs --env {} --project {} --apply",
  666. name, rep.misplaced.size(), h.env->path(), name);
  667. for (const auto& [pair, count] : pairs) {
  668. spdlog::error("placement audit: {} document(s) {}", count, pair);
  669. }
  670. } catch (const std::exception& e) {
  671. // Advisory only — a failed audit must not take the service down.
  672. spdlog::warn("placement audit: project '{}' could not be audited: {}",
  673. name, e.what());
  674. }
  675. }
  676. }
  677. bool DatabaseService::backfillIntoDocStore() {
  678. if (!projects_) {
  679. // Registry open failed earlier; nothing to back-fill into.
  680. return true;
  681. }
  682. // v2.3 — MemoryStore stores collection keys in qualified form
  683. // ("project:collection"). For the default project (back-compat
  684. // path) the marker check operates on the default env. Skip when
  685. // already migrated.
  686. auto handle = projects_->getHandle(smartbotic::database::kDefaultProject);
  687. if (handle.env) {
  688. try {
  689. if (smartbotic::db::storage::migration_complete(*handle.env)) {
  690. spdlog::info("v2.3 backfill: default project already migrated, skipping");
  691. return true;
  692. }
  693. } catch (const std::exception& e) {
  694. spdlog::warn("v2.3 backfill: migration_complete probe failed: {} — "
  695. "treating default env as fresh and proceeding",
  696. e.what());
  697. }
  698. }
  699. const auto collections = store_->listCollections();
  700. uint64_t total_docs = 0;
  701. uint64_t total_failures = 0;
  702. auto t0 = std::chrono::steady_clock::now();
  703. for (const auto& qualified : collections) {
  704. if (qualified.empty() || qualified[0] == '_') continue;
  705. // Each qualified key parses to (project, collection); route the
  706. // mirror write to the right project's env.
  707. smartbotic::database::ProjectCollection pc;
  708. try {
  709. pc = smartbotic::database::parseProjectCollection(qualified);
  710. } catch (const std::exception&) {
  711. // Legacy unqualified key — treat as default project.
  712. pc = {std::string(smartbotic::database::kDefaultProject), qualified};
  713. }
  714. auto* ds = projects_->getOrCreate(pc.project);
  715. if (!ds) {
  716. ++total_failures;
  717. continue;
  718. }
  719. const auto docs = store_->getAllDocuments(qualified);
  720. for (const auto& doc : docs) {
  721. std::optional<Document> opt_doc(doc);
  722. const uint64_t drift_before = mirror_drift_count_.load(std::memory_order_relaxed);
  723. // v2.11.0 T13 round 2 — UniqueViolation/MissingParentReference are
  724. // rethrown UNCAUGHT by applyDualWriteMirror (deliberately - see
  725. // both exceptions' header comments), unlike a genuine mirror
  726. // fault which is swallowed and only shows up as a drift bump
  727. // below. Uncaught here would escape this loop, backfillIntoDocStore
  728. // and initialize() itself, refusing startup over one rejected row.
  729. // Practically unreachable today - backfill only runs migrating a
  730. // v1.x dataset, which predates both `_relations` and unique-index
  731. // declarations - but catching removes that reasoning burden for
  732. // the next reader, same as any other row here that fails.
  733. try {
  734. smartbotic::db::storage::applyDualWriteMirror(
  735. ds, mirror_healthy_, mirror_drift_count_,
  736. pc.collection, doc.id, opt_doc, EventType::INSERT);
  737. } catch (const std::exception& e) {
  738. spdlog::error("v2.3 backfill: mirror rejected {}/{}: {} - row stays "
  739. "in MemoryStore only, LMDB does not have it",
  740. qualified, doc.id, e.what());
  741. ++total_failures;
  742. }
  743. if (mirror_drift_count_.load(std::memory_order_relaxed) > drift_before) {
  744. ++total_failures;
  745. }
  746. ++total_docs;
  747. if (total_docs % 10000 == 0) {
  748. sd_notify(0, "EXTEND_TIMEOUT_USEC=600000000");
  749. const std::string status = "STATUS=v2.3 backfill: " +
  750. std::to_string(total_docs) + " docs mirrored";
  751. sd_notify(0, status.c_str());
  752. spdlog::info("v2.3 backfill: {} docs mirrored ({} failures so far)",
  753. total_docs, total_failures);
  754. }
  755. }
  756. }
  757. const auto elapsed = std::chrono::duration_cast<std::chrono::milliseconds>(
  758. std::chrono::steady_clock::now() - t0).count();
  759. if (total_failures == 0) {
  760. // Stamp schema_version=2 on every opened project env. Subsequent
  761. // boots short-circuit at migration_complete().
  762. for (const auto& name : projects_->listOpen()) {
  763. auto h = projects_->getHandle(name);
  764. if (!h.env) continue;
  765. try {
  766. smartbotic::db::storage::mark_migration_complete(*h.env);
  767. } catch (const std::exception& e) {
  768. spdlog::warn("v2.3 backfill: schema_version marker write failed for "
  769. "project '{}': {}", name, e.what());
  770. }
  771. }
  772. spdlog::info("v2.3 backfill: complete — {} docs across {} collections in {} ms",
  773. total_docs, collections.size(), elapsed);
  774. } else {
  775. spdlog::error("v2.3 backfill: completed with {} failures out of {} docs "
  776. "in {} ms — mirror_healthy_=false, reads must NOT flip to LMDB",
  777. total_failures, total_docs, elapsed);
  778. }
  779. return true;
  780. }
  781. void DatabaseService::start() {
  782. if (running_.exchange(true)) {
  783. return;
  784. }
  785. spdlog::info("Starting database service with {} listener(s)", config_.listeners.size());
  786. // Start components
  787. store_->start();
  788. persistence_->start();
  789. events_->start();
  790. files_->start();
  791. if (config_.replicationEnabled) {
  792. replication_->start();
  793. // Set initial local collections for discovery
  794. replication_->setLocalCollections(store_->listCollections());
  795. }
  796. // Start gRPC server
  797. startGrpcServer();
  798. spdlog::info("Database service started");
  799. }
  800. void DatabaseService::stop() {
  801. if (!running_.exchange(false)) {
  802. return;
  803. }
  804. spdlog::info("Stopping database service...");
  805. // Stop gRPC server first
  806. stopGrpcServer();
  807. // Stop components in reverse order
  808. if (replication_) {
  809. replication_->stop();
  810. }
  811. if (files_) {
  812. files_->stop();
  813. }
  814. if (events_) {
  815. events_->stop();
  816. }
  817. if (persistence_ && store_) {
  818. // v1.8.0 — take a final snapshot on graceful shutdown so the next
  819. // boot's recovery is "trivial" (snapshot-loaded) instead of "WAL-only
  820. // replay" which trips auto-readonly mode. This makes ordinary
  821. // restarts pass through cleanly without needing
  822. // `smartbotic-db-cli unlock`.
  823. //
  824. // v1.8.2 — extend systemd's stop watchdog before doing it. On Zoe
  825. // (1.85M docs / ~5 GB tracked) the serialize takes ~70 s and the
  826. // default TimeoutStopSec=30s SIGKILLed the process mid-write, peak
  827. // RSS hit 11 GB (live state + uncompressed serialize buffer), and
  828. // the new snapshot never made it to disk anyway. EXTEND_TIMEOUT_USEC
  829. // tells systemd we're working — same protocol the recovery path
  830. // uses during startup. Paired with TimeoutStopSec=300 in the unit
  831. // file as the static cap.
  832. #ifdef HAVE_SYSTEMD
  833. sd_notify(0, "EXTEND_TIMEOUT_USEC=600000000");
  834. sd_notify(0, "STATUS=Writing final snapshot");
  835. #endif
  836. try {
  837. persistence_->forceSnapshot(*store_);
  838. spdlog::info("Final snapshot taken on shutdown");
  839. } catch (const std::exception& e) {
  840. spdlog::warn("Final snapshot on shutdown failed: {} — recovery on next "
  841. "boot may be WAL-only and trip auto-readonly mode", e.what());
  842. }
  843. }
  844. if (persistence_) {
  845. persistence_->stop();
  846. }
  847. if (store_) {
  848. store_->stop();
  849. }
  850. spdlog::info("Database service stopped");
  851. }
  852. void DatabaseService::wait() {
  853. if (serverThread_.joinable()) {
  854. serverThread_.join();
  855. }
  856. }
  857. void DatabaseService::signalStop() {
  858. stopRequested_ = true;
  859. stop();
  860. }
  861. nlohmann::json DatabaseService::getStats() const {
  862. nlohmann::json stats;
  863. if (store_) {
  864. auto storeStats = store_->getStats();
  865. stats["documents"] = storeStats.totalDocuments;
  866. stats["collections"] = storeStats.totalCollections;
  867. stats["memory_bytes"] = storeStats.estimatedMemoryBytes;
  868. stats["inserts"] = storeStats.insertCount;
  869. stats["updates"] = storeStats.updateCount;
  870. stats["deletes"] = storeStats.deleteCount;
  871. stats["queries"] = storeStats.queryCount;
  872. }
  873. if (persistence_) {
  874. auto persistStats = persistence_->getStats();
  875. stats["wal_sequence"] = persistStats.walSequence;
  876. stats["wal_size_bytes"] = persistStats.walSizeBytes;
  877. stats["snapshot_count"] = persistStats.snapshotCount;
  878. stats["last_snapshot_sequence"] = persistStats.lastSnapshotSequence;
  879. }
  880. if (events_) {
  881. stats["subscriptions"] = events_->subscriptionCount();
  882. }
  883. return stats;
  884. }
  885. DatabaseService::Config DatabaseService::loadConfig(const std::filesystem::path& configPath) {
  886. std::ifstream file(configPath);
  887. if (!file) {
  888. throw std::runtime_error("Failed to open config file: " + configPath.string());
  889. }
  890. nlohmann::json json;
  891. file >> json;
  892. return parseConfig(json);
  893. }
  894. DatabaseService::Config DatabaseService::parseConfig(const nlohmann::json& json) {
  895. Config config;
  896. // Check for both "storage" and "database" keys for backward compatibility
  897. const nlohmann::json* dbConfig = nullptr;
  898. if (json.contains("database")) {
  899. dbConfig = &json["database"];
  900. } else if (json.contains("storage")) {
  901. dbConfig = &json["storage"];
  902. }
  903. if (dbConfig) {
  904. const auto& db = *dbConfig;
  905. config.nodeId = db.value("node_id", config.nodeId);
  906. config.bindAddress = db.value("bind_address", config.bindAddress);
  907. config.rpcPort = db.value("rpc_port", config.rpcPort);
  908. // v2.4 — parse the per-listener fleet. If "listeners" is set,
  909. // it overrides the legacy bind_address+rpc_port. Otherwise we
  910. // synthesise one listener below using those legacy fields.
  911. if (db.contains("listeners") && db["listeners"].is_array()) {
  912. for (const auto& l : db["listeners"]) {
  913. ListenerConfig lc;
  914. lc.bind = l.value("bind", lc.bind);
  915. lc.port = l.value("port", lc.port);
  916. if (l.contains("tls") && l["tls"].is_object()) {
  917. const auto& t = l["tls"];
  918. lc.tls.enabled = t.value("enabled", false);
  919. lc.tls.cert_path = t.value("cert_path", std::string{});
  920. lc.tls.key_path = t.value("key_path", std::string{});
  921. lc.tls.auto_self_signed_if_missing =
  922. t.value("auto_self_signed_if_missing", true);
  923. }
  924. if (l.contains("auth") && l["auth"].is_object()) {
  925. const auto& a = l["auth"];
  926. lc.auth.required = a.value("required", false);
  927. if (a.contains("keys") && a["keys"].is_array()) {
  928. for (const auto& k : a["keys"]) {
  929. // v2.7.0 — a key is either a bare token (v2.4-v2.6
  930. // shape) or {"name","key"}. The name becomes the
  931. // principal that policy is written against; a bare
  932. // token maps to the reserved `unnamed`, so old
  933. // configs keep authenticating but cannot be told
  934. // apart in a policy until they are named.
  935. if (k.is_string()) {
  936. lc.auth.keys.push_back(
  937. {smartbotic::database::auth::kPrincipalUnnamed,
  938. k.get<std::string>()});
  939. } else if (k.is_object()) {
  940. const std::string name = k.value("name", std::string{});
  941. const std::string key = k.value("key", std::string{});
  942. if (key.empty()) {
  943. spdlog::warn("auth: listener {}:{} has a key entry "
  944. "with no `key` value - skipped",
  945. lc.bind, lc.port);
  946. continue;
  947. }
  948. lc.auth.keys.push_back(
  949. {name.empty()
  950. ? std::string(smartbotic::database::auth::kPrincipalUnnamed)
  951. : name,
  952. key});
  953. }
  954. }
  955. }
  956. }
  957. config.listeners.push_back(std::move(lc));
  958. }
  959. }
  960. // Expand environment variables in data directory
  961. std::string dataDir = db.value("data_directory", "");
  962. if (dataDir.find("${HOME}") != std::string::npos) {
  963. const char* home = std::getenv("HOME");
  964. if (home) {
  965. size_t pos = dataDir.find("${HOME}");
  966. dataDir.replace(pos, 7, home);
  967. }
  968. }
  969. config.dataDirectory = dataDir;
  970. // Migrations settings
  971. if (db.contains("migrations")) {
  972. auto& migrations = db["migrations"];
  973. config.migrations.enabled = migrations.value("enabled", config.migrations.enabled);
  974. config.migrations.autoApply = migrations.value("auto_apply", config.migrations.autoApply);
  975. config.migrations.failOnError = migrations.value("fail_on_error", config.migrations.failOnError);
  976. std::string migDir = migrations.value("directory", "");
  977. if (migDir.find("${HOME}") != std::string::npos) {
  978. const char* home = std::getenv("HOME");
  979. if (home) {
  980. size_t pos = migDir.find("${HOME}");
  981. migDir.replace(pos, 7, home);
  982. }
  983. }
  984. config.migrations.directory = migDir;
  985. }
  986. // Memory eviction settings
  987. if (db.contains("memory")) {
  988. auto& memory = db["memory"];
  989. // Deprecation log. These knobs control MemoryStore eviction, which
  990. // is still active for the MemoryStore-side mirror but does NOT bound
  991. // RSS since v2.0 (LMDB mmap is the dominant RSS contributor). The
  992. // substrate-level equivalent is the LMDB env mapsize and the OS page
  993. // cache.
  994. //
  995. // They disappear when the write-handler migration deletes MemoryStore
  996. // (docs/ROADMAP.md, "Pending"). An earlier revision of this comment
  997. // promised a v2.1 rename to `buffer_pool_size_mb` "per the Phase C
  998. // plan"; that plan was superseded by LMDB before it was written, no
  999. // such knob exists, and the rename is not planned.
  1000. for (const char* deprecated : {"max_memory_mb",
  1001. "eviction_threshold_percent",
  1002. "eviction_target_percent",
  1003. "eviction_check_interval_ms",
  1004. "eviction_chunk_size",
  1005. "eviction_chunk_pause_ms",
  1006. "max_eviction_passes_per_trigger",
  1007. "hot_write_floor_ms",
  1008. "memory_priority"}) {
  1009. if (memory.contains(deprecated)) {
  1010. spdlog::warn("v2.0 deprecation: storage.memory.{} is deprecated "
  1011. "and will be removed in v2.1. v2.0 RSS is bounded "
  1012. "by the LMDB env mapsize + OS page cache, not by "
  1013. "MemoryStore eviction. Setting still applied to "
  1014. "the MemoryStore mirror for back-compat.",
  1015. deprecated);
  1016. }
  1017. }
  1018. config.maxMemoryMb = memory.value("max_memory_mb", config.maxMemoryMb);
  1019. config.evictionThresholdPercent = memory.value("eviction_threshold_percent", config.evictionThresholdPercent);
  1020. config.evictionTargetPercent = memory.value("eviction_target_percent", config.evictionTargetPercent);
  1021. config.evictionCheckIntervalMs = memory.value("eviction_check_interval_ms", config.evictionCheckIntervalMs);
  1022. // NEW v1.7.0 eviction tuning (schema-only in T2; runtime use lands in T3/T4)
  1023. config.evictionChunkSize = memory.value("eviction_chunk_size", config.evictionChunkSize);
  1024. config.evictionChunkPauseMs = memory.value("eviction_chunk_pause_ms", config.evictionChunkPauseMs);
  1025. config.maxEvictionPassesPerTrigger = memory.value("max_eviction_passes_per_trigger", config.maxEvictionPassesPerTrigger);
  1026. config.hotWriteFloorMs = memory.value("hot_write_floor_ms", config.hotWriteFloorMs);
  1027. config.memorySoftPercent = memory.value("memory_soft_percent", config.memorySoftPercent);
  1028. config.memoryHardPercent = memory.value("memory_hard_percent", config.memoryHardPercent);
  1029. config.memoryEmergencyPercent = memory.value("memory_emergency_percent", config.memoryEmergencyPercent);
  1030. // v1.7.0 T10 — eviction burst event threshold
  1031. config.evictionBurstThreshold = memory.value("eviction_burst_threshold", config.evictionBurstThreshold);
  1032. // v2.4.3 — eviction drain cap (0 disables)
  1033. config.evictionMaxEpisodePercent = memory.value("eviction_max_episode_percent", config.evictionMaxEpisodePercent);
  1034. }
  1035. // v2.11.0 close-out — TTL sweep budgets. Its own block rather than
  1036. // `memory`, because these bound expiry work per sweep, not memory.
  1037. if (db.contains("ttl")) {
  1038. auto& ttl = db["ttl"];
  1039. config.ttlMaxCandidatesPerSweep =
  1040. ttl.value("max_candidates_per_sweep", config.ttlMaxCandidatesPerSweep);
  1041. config.ttlMaxBlockedRetriesPerSweep =
  1042. ttl.value("max_blocked_retries_per_sweep", config.ttlMaxBlockedRetriesPerSweep);
  1043. config.ttlBlockedRetrySweeps =
  1044. ttl.value("blocked_retry_sweeps", config.ttlBlockedRetrySweeps);
  1045. config.ttlBlockedSummarySweeps =
  1046. ttl.value("blocked_summary_sweeps", config.ttlBlockedSummarySweeps);
  1047. }
  1048. // Persistence settings
  1049. if (db.contains("persistence")) {
  1050. auto& persistence = db["persistence"];
  1051. config.walSyncIntervalMs = persistence.value("wal_sync_interval_ms", config.walSyncIntervalMs);
  1052. config.snapshotIntervalSec = persistence.value("snapshot_interval_sec", config.snapshotIntervalSec);
  1053. config.compressionEnabled = persistence.value("compression", "lz4") != "none";
  1054. // Snapshot durability (NEW in v1.6.1)
  1055. if (persistence.contains("snapshots")) {
  1056. const auto& snap = persistence["snapshots"];
  1057. config.persistenceConfig.validateAfterWrite = snap.value("validate_after_write", config.persistenceConfig.validateAfterWrite);
  1058. config.persistenceConfig.cleanupOnlyIfVerified = snap.value("cleanup_only_if_verified", config.persistenceConfig.cleanupOnlyIfVerified);
  1059. }
  1060. // Recovery (NEW in v1.6.1)
  1061. if (persistence.contains("recovery")) {
  1062. const auto& rec = persistence["recovery"];
  1063. std::string modeStr = rec.value("mode", std::string("normal"));
  1064. try {
  1065. config.persistenceConfig.recoveryMode = recoveryModeFromString(modeStr);
  1066. } catch (const std::exception& e) {
  1067. spdlog::warn("Invalid recovery.mode '{}', defaulting to 'normal': {}",
  1068. modeStr, e.what());
  1069. config.persistenceConfig.recoveryMode = RecoveryMode::Normal;
  1070. }
  1071. config.persistenceConfig.autoEscalate = rec.value("auto_escalate", config.persistenceConfig.autoEscalate);
  1072. config.persistenceConfig.allowEmptyOnFreshInstall = rec.value("allow_empty_on_fresh_install", config.persistenceConfig.allowEmptyOnFreshInstall);
  1073. }
  1074. }
  1075. // File settings
  1076. if (db.contains("files")) {
  1077. auto& files = db["files"];
  1078. config.maxFileSizeMb = files.value("max_file_size_mb", config.maxFileSizeMb);
  1079. // v2.8.0 — this interval now drives the file EXPIRY sweep as well
  1080. // as orphan cleanup, so accept the plainer name too. The original
  1081. // key kept working; nothing read either of them before 2.8.0
  1082. // because there was no sweeper thread at all.
  1083. config.fileCleanupIntervalSec =
  1084. files.value("cleanup_orphans_interval_sec", config.fileCleanupIntervalSec);
  1085. config.fileCleanupIntervalSec =
  1086. files.value("cleanup_interval_sec", config.fileCleanupIntervalSec);
  1087. // v2.8.0 — default retention per file type. The analogue of a
  1088. // collection's defaultTtlSeconds, so retention is an operator
  1089. // setting rather than something every uploader must remember.
  1090. //
  1091. // "default_ttl_seconds": {
  1092. // "acme:generated": 86400, // one project, one type
  1093. // "generated": 604800, // any project
  1094. // "acme:*": 2592000, // every type in one project
  1095. // "*": 0 // everything (0 = keep)
  1096. // }
  1097. if (files.contains("default_ttl_seconds") &&
  1098. files["default_ttl_seconds"].is_object()) {
  1099. for (const auto& [k, v] : files["default_ttl_seconds"].items()) {
  1100. if (v.is_number_unsigned()) {
  1101. config.fileDefaultTtlByType[k] = v.get<uint32_t>();
  1102. }
  1103. }
  1104. }
  1105. if (files.contains("allowed_types")) {
  1106. for (const auto& type : files["allowed_types"]) {
  1107. config.allowedFileTypes.push_back(type.get<std::string>());
  1108. }
  1109. }
  1110. }
  1111. // Encryption settings
  1112. if (db.contains("encryption")) {
  1113. auto& encryption = db["encryption"];
  1114. config.encryptionEnabled = encryption.value("enabled", config.encryptionEnabled);
  1115. config.autoGenerateKey = encryption.value("auto_generate_key", config.autoGenerateKey);
  1116. std::string keyFile = encryption.value("key_file", "");
  1117. if (keyFile.find("${HOME}") != std::string::npos) {
  1118. const char* home = std::getenv("HOME");
  1119. if (home) {
  1120. size_t pos = keyFile.find("${HOME}");
  1121. keyFile.replace(pos, 7, home);
  1122. }
  1123. }
  1124. config.keyFilePath = keyFile;
  1125. }
  1126. // Replication settings
  1127. if (db.contains("replication")) {
  1128. auto& replication = db["replication"];
  1129. config.replicationEnabled = replication.value("enabled", config.replicationEnabled);
  1130. config.conflictResolution = replication.value("conflict_resolution", config.conflictResolution);
  1131. if (replication.contains("peers")) {
  1132. for (const auto& peer : replication["peers"]) {
  1133. config.peerAddresses.push_back(peer.get<std::string>());
  1134. }
  1135. }
  1136. }
  1137. // gRPC settings (v1.6.2 — concurrency cap + configurable message sizes)
  1138. if (db.contains("grpc")) {
  1139. const auto& grpc = db["grpc"];
  1140. config.grpc.maxReceiveMessageSizeMb =
  1141. grpc.value("max_receive_message_size_mb", config.grpc.maxReceiveMessageSizeMb);
  1142. config.grpc.maxSendMessageSizeMb =
  1143. grpc.value("max_send_message_size_mb", config.grpc.maxSendMessageSizeMb);
  1144. config.grpc.resourceQuotaMemoryMb =
  1145. grpc.value("resource_quota_memory_mb", config.grpc.resourceQuotaMemoryMb);
  1146. config.grpc.maxConcurrentSubscribeStreams =
  1147. grpc.value("max_concurrent_subscribe_streams", config.grpc.maxConcurrentSubscribeStreams);
  1148. config.grpc.maxConcurrentFileStreams =
  1149. grpc.value("max_concurrent_file_streams", config.grpc.maxConcurrentFileStreams);
  1150. }
  1151. }
  1152. // v2.4 back-compat shim. If no listeners[] was provided, synthesize
  1153. // one from the legacy bind_address + rpc_port fields. This is what
  1154. // every existing v2.3 deployment will hit on first v2.4 boot —
  1155. // their configs don't know about listeners[] yet, and they get
  1156. // exactly the same listener they had before.
  1157. if (config.listeners.empty()) {
  1158. ListenerConfig legacy;
  1159. legacy.bind = config.bindAddress;
  1160. legacy.port = config.rpcPort;
  1161. legacy.tls.enabled = false;
  1162. legacy.auth.required = false;
  1163. config.listeners.push_back(std::move(legacy));
  1164. }
  1165. return config;
  1166. }
  1167. void DatabaseService::setupComponents() {
  1168. // v2.3 multi-project storage substrate.
  1169. //
  1170. // Each project lives at `<dataDir>/projects/<name>/env/`. The default
  1171. // project always exists and is what v2.0-v2.2 clients (which don't
  1172. // address projects) implicitly use.
  1173. //
  1174. // Pre-2.3 installs have their data at the legacy `<dataDir>/env/`
  1175. // path. migrateLegacyEnvToDefaultProject() detects that layout and
  1176. // atomically moves it under `projects/default/env/` before the
  1177. // registry opens.
  1178. migrateLegacyEnvToDefaultProject();
  1179. try {
  1180. constexpr size_t kMapSize = 4ULL << 30; // 4 GiB per project; LMDB grows lazily.
  1181. constexpr uint32_t kMaxDbs = 1024;
  1182. projects_ = std::make_unique<smartbotic::db::storage::ProjectStoreRegistry>(
  1183. config_.dataDirectory / "projects", kMapSize, kMaxDbs);
  1184. const auto opened = projects_->openExisting();
  1185. spdlog::info("v2.3 storage: opened {} project env(s): [{}]",
  1186. opened.size(),
  1187. [&opened]() {
  1188. std::string s;
  1189. for (size_t i = 0; i < opened.size(); ++i) {
  1190. if (i) s += ", ";
  1191. s += opened[i];
  1192. }
  1193. return s;
  1194. }());
  1195. } catch (const std::exception& e) {
  1196. spdlog::error("v2.3 storage: failed to open project registry at '{}': {} -- "
  1197. "service cannot start without LMDB; bailing",
  1198. (config_.dataDirectory / "projects").string(), e.what());
  1199. projects_.reset();
  1200. throw;
  1201. }
  1202. // Create memory store with eviction config
  1203. MemoryStore::Config storeConfig;
  1204. storeConfig.nodeId = config_.nodeId;
  1205. storeConfig.maxMemoryBytes = config_.maxMemoryMb * 1024ULL * 1024ULL;
  1206. storeConfig.evictionThresholdPercent = config_.evictionThresholdPercent;
  1207. storeConfig.evictionTargetPercent = config_.evictionTargetPercent;
  1208. storeConfig.evictionCheckIntervalMs = config_.evictionCheckIntervalMs;
  1209. // NEW v1.7.0 eviction tuning (wired in T2; consumed in T3/T4).
  1210. storeConfig.evictionChunkSize = config_.evictionChunkSize;
  1211. storeConfig.evictionChunkPauseMs = config_.evictionChunkPauseMs;
  1212. storeConfig.maxEvictionPassesPerTrigger = config_.maxEvictionPassesPerTrigger;
  1213. storeConfig.hotWriteFloorMs = config_.hotWriteFloorMs;
  1214. storeConfig.memorySoftPercent = config_.memorySoftPercent;
  1215. storeConfig.memoryHardPercent = config_.memoryHardPercent;
  1216. storeConfig.memoryEmergencyPercent = config_.memoryEmergencyPercent;
  1217. storeConfig.evictionBurstThreshold = config_.evictionBurstThreshold;
  1218. storeConfig.evictionMaxEpisodePercent = config_.evictionMaxEpisodePercent;
  1219. // v2.11.0 close-out — TTL sweep budgets (storage.ttl).
  1220. storeConfig.ttlMaxCandidatesPerSweep = config_.ttlMaxCandidatesPerSweep;
  1221. storeConfig.ttlMaxBlockedRetriesPerSweep = config_.ttlMaxBlockedRetriesPerSweep;
  1222. storeConfig.ttlBlockedRetrySweeps = config_.ttlBlockedRetrySweeps;
  1223. storeConfig.ttlBlockedSummarySweeps = config_.ttlBlockedSummarySweeps;
  1224. store_ = std::make_unique<MemoryStore>(storeConfig);
  1225. // Create view manager (cache loaded in initialize() after persistence recovery)
  1226. view_manager_ = std::make_unique<ViewManager>(*store_);
  1227. // Create per-collection config manager and attach it to the store so the
  1228. // write paths route document timestamp stamps through it.
  1229. config_manager_ = std::make_unique<CollectionConfigManager>(*store_);
  1230. // v2.8.0 — access policy. Constructed here; its cache is loaded after
  1231. // recovery alongside the view and collection-config caches.
  1232. policy_manager_ = std::make_unique<PolicyManager>(*store_);
  1233. // v2.11.0 T6a — relation declarations. Constructed here; its cache is
  1234. // loaded after recovery alongside the view/config/policy caches.
  1235. relation_manager_ = std::make_unique<RelationManager>(*store_);
  1236. store_->setConfigManager(config_manager_.get());
  1237. // v1.9.0 — disk-resident version history. Replaces the in-heap
  1238. // `CollectionData::versionHistory` deques that were the dominant
  1239. // source of the Zoe untracked-RSS gap. One file per collection at
  1240. // `<dataDir>/history/<collection>.hlog`. Lifecycle: constructed
  1241. // here, wired into the store immediately so snapshot deserializer
  1242. // can migrate pre-v1.9 in-memory history blocks straight to disk.
  1243. HistoryStore::Config histConfig;
  1244. histConfig.dataDir = config_.dataDirectory;
  1245. history_store_ = std::make_unique<HistoryStore>(histConfig);
  1246. store_->setHistoryStore(history_store_.get());
  1247. // v2.3 — wire the LMDB mirror INTO MemoryStore with a project
  1248. // resolver. Every write that flows through store_ commits to the
  1249. // right per-project LMDB env BEFORE the per-collection lock releases.
  1250. if (projects_) {
  1251. auto* registry = projects_.get();
  1252. store_->setDocumentStoreMirror(
  1253. [registry](std::string_view project) -> smartbotic::db::storage::DocumentStore* {
  1254. return registry->getOrCreate(project);
  1255. },
  1256. &mirror_healthy_,
  1257. &mirror_drift_count_);
  1258. }
  1259. // Create persistence manager. Start from any persistenceConfig values the
  1260. // caller (config loader / CLI parser) has already populated — including
  1261. // recoveryMode — and overlay the top-level convenience fields.
  1262. PersistenceManager::Config persistConfig = config_.persistenceConfig;
  1263. persistConfig.dataDir = config_.dataDirectory;
  1264. persistConfig.walSyncIntervalMs = config_.walSyncIntervalMs;
  1265. persistConfig.snapshotIntervalSec = config_.snapshotIntervalSec;
  1266. persistConfig.compressionEnabled = config_.compressionEnabled;
  1267. persistConfig.nodeId = config_.nodeId;
  1268. // Mirror the resolved recoveryMode back into our Config so downstream code
  1269. // (logging, RPC handlers) can inspect config_.persistenceConfig.recoveryMode
  1270. // without having to reach into the PersistenceManager.
  1271. config_.persistenceConfig = persistConfig;
  1272. persistence_ = std::make_unique<PersistenceManager>(persistConfig);
  1273. // Connect store callbacks to persistence - uses setPersistCallback for WAL logging
  1274. store_->setPersistCallback([this](const std::string& collection, const std::string& id,
  1275. const std::optional<Document>& doc, EventType eventType) {
  1276. // Log to WAL based on event type. v1.7.4: capture the assigned WAL
  1277. // sequence and record it on the store so eviction stubs can carry
  1278. // a real seq (instead of doc.version, which was useless as a
  1279. // `fromSequence` hint and forced a full WAL scan on every fault).
  1280. uint64_t walSeq = 0;
  1281. switch (eventType) {
  1282. case EventType::INSERT:
  1283. if (doc) walSeq = persistence_->logInsert(collection, *doc);
  1284. break;
  1285. case EventType::UPDATE:
  1286. if (doc) walSeq = persistence_->logUpdate(collection, *doc);
  1287. break;
  1288. case EventType::DELETE:
  1289. persistence_->logDelete(collection, id);
  1290. break;
  1291. default:
  1292. break;
  1293. }
  1294. if (walSeq > 0 && (eventType == EventType::INSERT || eventType == EventType::UPDATE)) {
  1295. store_->recordWalSequence(collection, id, walSeq);
  1296. }
  1297. // v2.0 dual-write — moved INTO MemoryStore::mirrorWriteToDocStore so
  1298. // that it runs UNDER the per-collection write lock. That closes the
  1299. // race window where readers could see a doc in MemoryStore but not
  1300. // yet in LMDB. The callback now does only WAL + replication + events.
  1301. // v2.11.0 T12 review (C2) — replication + events extracted into
  1302. // notifyReplicationAndEvents() so relations/relation_cascade.cpp can
  1303. // drive them explicitly too. See that method's doc comment.
  1304. notifyReplicationAndEvents(collection, id, doc, eventType);
  1305. });
  1306. // Set up document load callback for LRU eviction recovery.
  1307. //
  1308. // v1.7.4: walSequence is now a real WAL sequence (set by the persist
  1309. // callback below at append time, not the per-doc version counter).
  1310. // Pass `walSequence - 1` as the exclusive `fromSequence` so replay
  1311. // returns entries with seq >= walSequence — the doc's last write is
  1312. // exactly at walSequence, so the very first hit is the one we want.
  1313. // For docs whose seq is unknown (0 — e.g. loaded from a snapshot
  1314. // that predates v1.7.4), fall back to the full-WAL scan from 0.
  1315. store_->setDocumentLoadCallback([this](const std::string& collection,
  1316. const std::string& id,
  1317. uint64_t walSequence) -> std::optional<Document> {
  1318. uint64_t fromSequence = walSequence > 0 ? walSequence - 1 : 0;
  1319. return persistence_->loadDocument(collection, id, fromSequence);
  1320. });
  1321. // Create encryption manager
  1322. EncryptionManager::Config encryptConfig;
  1323. encryptConfig.enabled = config_.encryptionEnabled;
  1324. encryptConfig.keyFilePath = config_.keyFilePath;
  1325. encryptConfig.autoGenerateKey = config_.autoGenerateKey;
  1326. encryption_ = std::make_unique<EncryptionManager>(encryptConfig);
  1327. // Create event manager
  1328. EventManager::Config eventConfig;
  1329. eventConfig.nodeId = config_.nodeId;
  1330. events_ = std::make_unique<EventManager>(eventConfig);
  1331. // Create file manager
  1332. FileManager::Config fileConfig;
  1333. fileConfig.filesDir = config_.dataDirectory / "files";
  1334. fileConfig.maxFileSizeMb = config_.maxFileSizeMb;
  1335. fileConfig.allowedMimeTypes = config_.allowedFileTypes;
  1336. fileConfig.cleanupIntervalSec = config_.fileCleanupIntervalSec;
  1337. fileConfig.defaultTtlSecondsByType = config_.fileDefaultTtlByType;
  1338. files_ = std::make_unique<FileManager>(fileConfig);
  1339. // Create replication manager
  1340. ReplicationManager::Config replConfig;
  1341. replConfig.nodeId = config_.nodeId;
  1342. replConfig.peerAddresses = config_.peerAddresses;
  1343. replConfig.conflictResolution = config_.conflictResolution;
  1344. replication_ = std::make_unique<ReplicationManager>(replConfig);
  1345. // Wire replication entry apply callback
  1346. replication_->setEntryApplyCallback([this](const databasepb::ReplicationEntry& entry) {
  1347. applyReplicatedEntry(entry);
  1348. });
  1349. // Create gRPC implementations
  1350. storageImpl_ = std::make_unique<DatabaseGrpcImpl>(
  1351. *this, *store_, *persistence_, *events_, *files_, *encryption_, *view_manager_, *config_manager_,
  1352. *policy_manager_, *relation_manager_
  1353. );
  1354. // v2.7.0 — hand the impl the union of every listener's keys so it can
  1355. // resolve a principal per call. A token is only accepted if some listener's
  1356. // auth processor accepted it, so the union cannot widen authentication; it
  1357. // only names what was already authenticated.
  1358. {
  1359. std::vector<smartbotic::database::auth::PrincipalResolver::NamedKey> all;
  1360. for (const auto& l : config_.listeners) {
  1361. for (const auto& k : l.auth.keys) all.push_back({k.name, k.key});
  1362. }
  1363. storageImpl_->setPrincipalKeys(std::move(all));
  1364. }
  1365. // v1.6.2 — wire the streaming-RPC concurrency limits from GrpcConfig.
  1366. storageImpl_->setStreamLimits(
  1367. config_.grpc.maxConcurrentSubscribeStreams,
  1368. config_.grpc.maxConcurrentFileStreams);
  1369. replicationImpl_ = std::make_unique<DatabaseReplicationGrpcImpl>(
  1370. *store_, *replication_, *persistence_
  1371. );
  1372. }
  1373. void DatabaseService::notifyReplicationAndEvents(const std::string& collection,
  1374. const std::string& id,
  1375. const std::optional<Document>& doc,
  1376. EventType eventType) {
  1377. // Queue for replication broadcast
  1378. if (replication_ && config_.replicationEnabled) {
  1379. databasepb::ReplicationEntry entry;
  1380. entry.set_collection(collection);
  1381. entry.set_document_id(id);
  1382. // v1.8.0 — stamp our nodeId so peers can attribute the write to
  1383. // its actual origin. Pre-v1.8 broadcasts left this empty, which
  1384. // is now treated by receivers as "legacy / unknown origin".
  1385. entry.set_node_id(config_.nodeId);
  1386. entry.set_global_timestamp(
  1387. std::chrono::duration_cast<std::chrono::milliseconds>(
  1388. std::chrono::system_clock::now().time_since_epoch()
  1389. ).count());
  1390. switch (eventType) {
  1391. case EventType::INSERT:
  1392. entry.set_op(databasepb::OP_INSERT);
  1393. // Send the full Document envelope (id, collection, data,
  1394. // version, timestamps, encryption state). The follower's
  1395. // applyReplicatedEntry uses Document::fromJson() which
  1396. // expects this envelope — sending only doc->data silently
  1397. // dropped every user field on the other side.
  1398. if (doc) entry.set_data(doc->toJson().dump());
  1399. break;
  1400. case EventType::UPDATE:
  1401. entry.set_op(databasepb::OP_UPDATE);
  1402. if (doc) entry.set_data(doc->toJson().dump());
  1403. break;
  1404. case EventType::DELETE:
  1405. entry.set_op(databasepb::OP_DELETE);
  1406. break;
  1407. default:
  1408. break;
  1409. }
  1410. replication_->queueForReplication(entry);
  1411. }
  1412. // Publish event
  1413. if (events_) {
  1414. DatabaseEvent event;
  1415. event.type = eventType;
  1416. event.collection = collection;
  1417. event.documentId = id;
  1418. event.timestamp = std::chrono::duration_cast<std::chrono::milliseconds>(
  1419. std::chrono::system_clock::now().time_since_epoch()
  1420. ).count();
  1421. event.nodeId = config_.nodeId;
  1422. if (doc) {
  1423. event.data = doc->data();
  1424. }
  1425. events_->publish(event);
  1426. }
  1427. }
  1428. void DatabaseService::applyReplicatedEntry(const databasepb::ReplicationEntry& entry) {
  1429. try {
  1430. // v1.8.0 — also append the replicated entry to the local WAL, tagged
  1431. // with the originating node's id. This is what makes follower-side
  1432. // recovery (eviction → WAL fallback, restart → WAL replay) preserve
  1433. // replicated documents. Echo amplification is prevented by:
  1434. // - GetEntriesSince emitting walEntry.nodeId (not local node id)
  1435. // - SyncProtocol skipping entries whose origin is the peer itself
  1436. // both implemented in the same v1.8.0 cycle.
  1437. const std::string origin = entry.node_id().empty()
  1438. ? config_.nodeId // legacy peer (pre-1.8) — best-guess local
  1439. : entry.node_id();
  1440. switch (entry.op()) {
  1441. case databasepb::OP_INSERT:
  1442. case databasepb::OP_UPDATE:
  1443. case databasepb::OP_UPSERT: {
  1444. if (entry.data().empty()) {
  1445. spdlog::warn("Replication entry has no data for op {}", static_cast<int>(entry.op()));
  1446. return;
  1447. }
  1448. auto json = smartbotic::db::parse_to_nlohmann(entry.data());
  1449. Document doc = Document::fromJson(json);
  1450. doc.id = entry.document_id();
  1451. doc.nodeId = entry.node_id();
  1452. // loadDocument bypasses the persist-and-broadcast callback to avoid
  1453. // re-replicating; we explicitly write to the WAL right after with
  1454. // the origin-aware overload so eviction/WAL recovery still works.
  1455. store_->loadDocument(entry.collection(), doc);
  1456. // v2.3.1 — also mirror to LMDB. loadDocument bypasses the
  1457. // MemoryStore persist callback that normally drives the
  1458. // dual-write mirror; in v2.0-v2.2 the mirror was a side
  1459. // channel and reads still hit MemoryStore, so the bypass
  1460. // was harmless. In v2.3 reads are LMDB-first, so a
  1461. // follower that only filled MemoryStore would return
  1462. // empty Find/Get for every replicated doc. Route the
  1463. // entry through the same project-aware mirror the write
  1464. // handlers use.
  1465. if (projects_) {
  1466. auto pc = smartbotic::database::parseProjectCollection(
  1467. entry.collection());
  1468. if (auto* ds = projects_->getOrCreate(pc.project)) {
  1469. std::optional<Document> opt_doc(doc);
  1470. try {
  1471. smartbotic::db::storage::applyDualWriteMirror(
  1472. ds, mirror_healthy_, mirror_drift_count_,
  1473. pc.collection, doc.id, opt_doc, EventType::INSERT);
  1474. } catch (const smartbotic::db::storage::UniqueViolation& e) {
  1475. // v2.11.0 T11 finding 1 — deliberate choice, not an
  1476. // accident of the generic catch below. By this
  1477. // point store_->loadDocument() above has ALREADY
  1478. // put the row into MemoryStore (loadDocument
  1479. // bypasses the persist callback specifically to
  1480. // avoid re-replicating), so unlike every
  1481. // client-facing write path there is no "reject the
  1482. // whole operation" available: the origin node
  1483. // already committed this write and every other
  1484. // follower is expected to converge to it too.
  1485. // Undoing loadDocument here would silently diverge
  1486. // this follower's data from the rest of the
  1487. // cluster over a constraint that may not even
  1488. // exist on the origin (this follower could have
  1489. // declared the unique field locally, after the
  1490. // fact — replication carries no guarantee the
  1491. // origin shares this node's config). So
  1492. // MemoryStore keeps the row — it is the true
  1493. // record of what replication decided — and the
  1494. // mirror is marked unhealthy ON PURPOSE: this is
  1495. // exactly the pre-T11 swallow-and-flip outcome,
  1496. // chosen deliberately for this one caller so an
  1497. // operator sees the drift (and LMDB-first reads
  1498. // fall back to MemoryStore, which has the answer)
  1499. // rather than nothing signalling a permanent gap
  1500. // between the two stores.
  1501. spdlog::error(
  1502. "v2.11 replication: unique constraint violated applying "
  1503. "{}/{}: {} - MemoryStore holds the row, LMDB does not; "
  1504. "marking the mirror unhealthy so reads fall back instead "
  1505. "of silently disagreeing with MemoryStore",
  1506. entry.collection(), doc.id, e.what());
  1507. mirror_healthy_.store(false, std::memory_order_release);
  1508. mirror_drift_count_.fetch_add(1, std::memory_order_relaxed);
  1509. } catch (const smartbotic::db::storage::MissingParentReference& e) {
  1510. // v2.11.0 T13 — same deliberate swallow-and-flip as
  1511. // the UniqueViolation catch just above, for the same
  1512. // reason: loadDocument() has already committed this
  1513. // row into MemoryStore, the origin already accepted
  1514. // the write, and this follower may not even have the
  1515. // same validate_on_write declaration (or the same
  1516. // parent data) as the origin did at the time it
  1517. // wrote. MemoryStore keeps the row; the mirror is
  1518. // marked unhealthy on purpose so an operator sees
  1519. // the drift and reads fall back to MemoryStore
  1520. // rather than the two stores silently disagreeing.
  1521. spdlog::error(
  1522. "v2.11 replication: validate_on_write rejected applying "
  1523. "{}/{}: {} - MemoryStore holds the row, LMDB does not; "
  1524. "marking the mirror unhealthy so reads fall back instead "
  1525. "of silently disagreeing with MemoryStore",
  1526. entry.collection(), doc.id, e.what());
  1527. mirror_healthy_.store(false, std::memory_order_release);
  1528. mirror_drift_count_.fetch_add(1, std::memory_order_relaxed);
  1529. }
  1530. }
  1531. }
  1532. if (persistence_) {
  1533. uint64_t walSeq = (entry.op() == databasepb::OP_INSERT)
  1534. ? persistence_->logInsert(entry.collection(), doc, origin)
  1535. : persistence_->logUpdate(entry.collection(), doc, origin);
  1536. if (walSeq > 0) {
  1537. store_->recordWalSequence(entry.collection(), doc.id, walSeq);
  1538. }
  1539. }
  1540. spdlog::trace("Applied replicated {} to {}/{} from {}",
  1541. entry.op() == databasepb::OP_INSERT ? "insert" :
  1542. entry.op() == databasepb::OP_UPDATE ? "update" : "upsert",
  1543. entry.collection(), entry.document_id(), entry.node_id());
  1544. break;
  1545. }
  1546. case databasepb::OP_DELETE: {
  1547. store_->remove(entry.collection(), entry.document_id());
  1548. if (persistence_) {
  1549. persistence_->logDelete(entry.collection(), entry.document_id(), origin);
  1550. }
  1551. spdlog::trace("Applied replicated delete to {}/{} from {}",
  1552. entry.collection(), entry.document_id(), entry.node_id());
  1553. break;
  1554. }
  1555. case databasepb::OP_CREATE_COLLECTION: {
  1556. CollectionOptions options;
  1557. store_->createCollection(entry.collection(), options);
  1558. if (persistence_) {
  1559. persistence_->logCreateCollection(entry.collection(), options, origin);
  1560. }
  1561. spdlog::trace("Applied replicated create collection {} from {}",
  1562. entry.collection(), entry.node_id());
  1563. break;
  1564. }
  1565. case databasepb::OP_DROP_COLLECTION: {
  1566. store_->dropCollection(entry.collection());
  1567. if (persistence_) {
  1568. persistence_->logDropCollection(entry.collection(), origin);
  1569. }
  1570. spdlog::trace("Applied replicated drop collection {} from {}",
  1571. entry.collection(), entry.node_id());
  1572. break;
  1573. }
  1574. default:
  1575. spdlog::warn("Unknown replication operation type: {}", static_cast<int>(entry.op()));
  1576. break;
  1577. }
  1578. } catch (const std::exception& e) {
  1579. spdlog::error("Failed to apply replicated entry: {}", e.what());
  1580. }
  1581. }
  1582. void DatabaseService::startGrpcServer() {
  1583. if (config_.listeners.empty()) {
  1584. // Should never hit — parseConfig's back-compat shim always
  1585. // synthesises at least one. Defend anyway.
  1586. throw std::runtime_error("startGrpcServer: no listeners configured");
  1587. }
  1588. const int64_t mb = 1024 * 1024;
  1589. spdlog::info("gRPC config: recv={}MB send={}MB quota={}MB subs={} files={}",
  1590. config_.grpc.maxReceiveMessageSizeMb,
  1591. config_.grpc.maxSendMessageSizeMb,
  1592. config_.grpc.resourceQuotaMemoryMb,
  1593. config_.grpc.maxConcurrentSubscribeStreams,
  1594. config_.grpc.maxConcurrentFileStreams);
  1595. grpcServers_.reserve(config_.listeners.size());
  1596. for (const auto& listener : config_.listeners) {
  1597. const std::string addr = listener.bind + ":" + std::to_string(listener.port);
  1598. grpc::ServerBuilder builder;
  1599. // v2.4 — choose credentials per listener. Plaintext for
  1600. // tls.enabled=false; SslServerCredentials with the operator's
  1601. // cert+key (or an auto-generated self-signed pair) for true.
  1602. std::shared_ptr<grpc::ServerCredentials> creds;
  1603. if (listener.tls.enabled) {
  1604. std::filesystem::path cert_path = listener.tls.cert_path;
  1605. std::filesystem::path key_path = listener.tls.key_path;
  1606. // Resolve cert+key. If either is missing/unreadable AND
  1607. // auto_self_signed_if_missing is set, generate a pair.
  1608. const bool cert_ok = !cert_path.empty()
  1609. && std::filesystem::exists(cert_path);
  1610. const bool key_ok = !key_path.empty()
  1611. && std::filesystem::exists(key_path);
  1612. if (!cert_ok || !key_ok) {
  1613. if (!listener.tls.auto_self_signed_if_missing) {
  1614. throw std::runtime_error(
  1615. "TLS listener " + addr +
  1616. ": cert or key missing and auto_self_signed_if_missing=false");
  1617. }
  1618. const auto out_dir = config_.dataDirectory / "tls";
  1619. auto gen = smartbotic::database::tls_util::generateSelfSignedCert(
  1620. out_dir, listener.bind);
  1621. cert_path = gen.cert_path;
  1622. key_path = gen.key_path;
  1623. spdlog::warn(
  1624. "TLS listener {}: no operator cert at '{}' (or key at '{}'); "
  1625. "auto-generated self-signed cert at '{}' (key '{}'). "
  1626. "OK for dev. For production, drop a real cert + key at "
  1627. "the configured cert_path/key_path.",
  1628. addr,
  1629. listener.tls.cert_path.string(),
  1630. listener.tls.key_path.string(),
  1631. cert_path.string(),
  1632. key_path.string());
  1633. }
  1634. auto slurp = [](const std::filesystem::path& p) {
  1635. std::ifstream f(p);
  1636. if (!f) throw std::runtime_error("read " + p.string());
  1637. std::stringstream ss; ss << f.rdbuf();
  1638. return ss.str();
  1639. };
  1640. grpc::SslServerCredentialsOptions sslOpts;
  1641. grpc::SslServerCredentialsOptions::PemKeyCertPair kp;
  1642. kp.private_key = slurp(key_path);
  1643. kp.cert_chain = slurp(cert_path);
  1644. sslOpts.pem_key_cert_pairs.push_back(std::move(kp));
  1645. creds = grpc::SslServerCredentials(sslOpts);
  1646. } else {
  1647. creds = grpc::InsecureServerCredentials();
  1648. }
  1649. // v2.4 Stage D — attach the bearer-token auth processor when
  1650. // the listener requires it. The processor validates the
  1651. // `authorization: Bearer <key>` metadata against the listener's
  1652. // configured keys list. gRPC rejects with UNAUTHENTICATED
  1653. // (no per-handler code needed).
  1654. if (listener.auth.required) {
  1655. if (listener.auth.keys.empty()) {
  1656. throw std::runtime_error(
  1657. "Listener " + addr + ": auth.required=true but auth.keys is empty");
  1658. }
  1659. // v2.7.0 — a key name is a principal, so it must be unambiguous.
  1660. // Refuse at startup rather than resolving a policy against a name
  1661. // that means two different callers.
  1662. {
  1663. std::set<std::string> seen;
  1664. for (const auto& k : listener.auth.keys) {
  1665. if (k.name != smartbotic::database::auth::kPrincipalUnnamed &&
  1666. smartbotic::database::auth::isReservedPrincipal(k.name)) {
  1667. throw std::runtime_error(
  1668. "Listener " + addr + ": auth key name '" + k.name +
  1669. "' is reserved. `anonymous` denotes an unauthenticated "
  1670. "caller and `unnamed` denotes a legacy bare key; a real "
  1671. "key must not be able to impersonate either in a policy.");
  1672. }
  1673. if (!seen.insert(k.name).second &&
  1674. k.name != smartbotic::database::auth::kPrincipalUnnamed) {
  1675. throw std::runtime_error(
  1676. "Listener " + addr + ": duplicate auth key name '" +
  1677. k.name + "'. A principal must identify one caller, "
  1678. "otherwise a policy written against it is ambiguous.");
  1679. }
  1680. }
  1681. }
  1682. if (!listener.tls.enabled) {
  1683. // gRPC ABORTS the process if an auth metadata processor is
  1684. // attached to insecure credentials
  1685. // (insecure_server_credentials.cc: "assertion failed: 0").
  1686. // This used to be a WARN, which meant a plaintext+auth listener
  1687. // looked merely inadvisable in the config and then killed the
  1688. // server on start with an unexplained assert. Refuse it here
  1689. // with a message that says what to do instead.
  1690. throw std::runtime_error(
  1691. "Listener " + addr + ": auth.required=true requires "
  1692. "tls.enabled=true. gRPC does not support an auth metadata "
  1693. "processor on insecure credentials and will abort the "
  1694. "process. Enable TLS for this listener (set "
  1695. "tls.auto_self_signed_if_missing for a local dev cert), or "
  1696. "drop auth and rely on binding to loopback.");
  1697. }
  1698. std::vector<smartbotic::database::auth::BearerAuthProcessor::NamedKey> nk;
  1699. nk.reserve(listener.auth.keys.size());
  1700. for (const auto& k : listener.auth.keys) nk.push_back({k.name, k.key});
  1701. creds->SetAuthMetadataProcessor(
  1702. std::make_shared<smartbotic::database::auth::BearerAuthProcessor>(
  1703. std::move(nk)));
  1704. }
  1705. builder.AddListeningPort(addr, creds);
  1706. builder.RegisterService(storageImpl_.get());
  1707. builder.RegisterService(replicationImpl_.get());
  1708. builder.SetMaxReceiveMessageSize(
  1709. static_cast<int>(config_.grpc.maxReceiveMessageSizeMb * mb));
  1710. builder.SetMaxSendMessageSize(
  1711. static_cast<int>(config_.grpc.maxSendMessageSizeMb * mb));
  1712. grpc::ResourceQuota quota("smartbotic-db-" + addr);
  1713. quota.Resize(config_.grpc.resourceQuotaMemoryMb * mb);
  1714. builder.SetResourceQuota(quota);
  1715. auto server = builder.BuildAndStart();
  1716. if (!server) {
  1717. throw std::runtime_error("Failed to start gRPC server on " + addr);
  1718. }
  1719. // v2.4 Stage B log line — listener identity printed so operators
  1720. // know exactly which interface is exposed with which policy.
  1721. // Stage C/D will extend with TLS + auth status.
  1722. spdlog::info("gRPC listener up: {} (tls={}, auth_required={})",
  1723. addr,
  1724. listener.tls.enabled ? "true" : "false",
  1725. listener.auth.required ? "true" : "false");
  1726. grpcServers_.push_back(std::move(server));
  1727. }
  1728. }
  1729. void DatabaseService::stopGrpcServer() {
  1730. // ⚠ Shutdown() WITHOUT a deadline waits for every in-flight RPC to return,
  1731. // and it does NOT cancel them. `Subscribe` never returns on its own: its
  1732. // handler parks in `while (!context->IsCancelled())`, and that context is
  1733. // cancelled only when the client goes away or a shutdown DEADLINE expires.
  1734. // So a single connected subscriber - the normal state of any consumer using
  1735. // events - hung this call forever.
  1736. //
  1737. // The cost is not just a slow stop. `stop()` calls this FIRST and takes the
  1738. // final snapshot AFTER it, so the hang blocked the snapshot; systemd then
  1739. // SIGKILLed at TimeoutStopSec=300, and the next boot recovered by WAL-only
  1740. // replay, which trips auto read-only mode and needs an operator `unlock`.
  1741. // A graceful restart therefore looked like a crash to the next boot.
  1742. //
  1743. // With a deadline, gRPC cancels whatever is still running once it expires,
  1744. // the Subscribe loops observe IsCancelled() and unwind, and shutdown
  1745. // proceeds. The deadline is ABSOLUTE and shared across listeners, so the
  1746. // total wait is bounded by the grace period rather than by
  1747. // grace × listener count.
  1748. //
  1749. // Five seconds: long enough for an ordinary unary RPC in flight to finish
  1750. // (they are milliseconds), short enough that it is invisible next to the
  1751. // final snapshot, which is the part of shutdown that legitimately takes
  1752. // time (~70 s on the largest known dataset).
  1753. constexpr auto kShutdownGrace = std::chrono::seconds(5);
  1754. const auto deadline = std::chrono::system_clock::now() + kShutdownGrace;
  1755. for (auto& server : grpcServers_) {
  1756. if (server) server->Shutdown(deadline);
  1757. }
  1758. grpcServers_.clear();
  1759. }
  1760. } // namespace smartbotic::database