database_service.cpp 79 KB

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