database_service.hpp 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457
  1. #pragma once
  2. #include "document.hpp"
  3. #include "memory_store.hpp"
  4. #include "database_grpc_impl.hpp"
  5. #include "persistence/persistence_manager.hpp"
  6. #include "persistence/history_store.hpp"
  7. #include "encryption/encryption_manager.hpp"
  8. #include "events/event_manager.hpp"
  9. #include "files/file_manager.hpp"
  10. #include "replication/replication_manager.hpp"
  11. #include "migrations/migration_runner.hpp"
  12. #include "views/view_manager.hpp"
  13. #include "config/collection_config_manager.hpp"
  14. #include "security/policy_manager.hpp"
  15. #include "relations/relation_manager.hpp"
  16. // LMDB storage substrate (v2.0+). DocumentStore + LmdbEnv live alongside the
  17. // v1.x MemoryStore, which is still the write entry point and a bounded read
  18. // cache; writes are mirrored into LMDB under MemoryStore's per-collection lock
  19. // and reads are LMDB-first while the mirror is healthy.
  20. //
  21. // This coexistence is deliberate and still current. Retiring MemoryStore means
  22. // migrating the write handlers onto a WriteCoordinator - see the "Pending"
  23. // section of docs/ROADMAP.md. Older comments call that "Stage 4"; the "Phase C"
  24. // label they sometimes carry refers to a homegrown page engine that was never
  25. // built (LMDB replaced it).
  26. namespace smartbotic::db::storage {
  27. class LmdbEnv;
  28. class DocumentStore;
  29. class ProjectStoreRegistry;
  30. }
  31. #include <nlohmann/json.hpp>
  32. #include <atomic>
  33. #include <filesystem>
  34. #include <memory>
  35. #include <mutex>
  36. #include <string>
  37. #include <string_view>
  38. #include <unordered_map>
  39. #include <thread>
  40. #include <vector>
  41. namespace grpc {
  42. class Server;
  43. }
  44. namespace smartbotic::database {
  45. /**
  46. * Main storage service application.
  47. * Coordinates all components and manages the gRPC server.
  48. */
  49. class DatabaseService {
  50. public:
  51. // v2.4 — per-listener configuration. Each binding can have its own
  52. // TLS material and auth policy. The default 127.0.0.1 listener
  53. // preserves v2.3 behaviour (plaintext, no auth) so existing local
  54. // consumers keep working without code change.
  55. struct ListenerConfig {
  56. std::string bind = "127.0.0.1";
  57. uint16_t port = 9004;
  58. struct TlsConfig {
  59. bool enabled = false;
  60. std::filesystem::path cert_path;
  61. std::filesystem::path key_path;
  62. // When TLS is enabled but cert+key are missing/unreadable, the
  63. // server auto-generates a 10-year self-signed cert to
  64. // <dataDir>/tls/auto_self_signed.{pem,key}. Operators who run
  65. // a real PKI override by dropping their cert at cert_path.
  66. bool auto_self_signed_if_missing = true;
  67. } tls;
  68. struct AuthConfig {
  69. bool required = false;
  70. // v2.7.0 — a key now carries a name, which becomes the
  71. // PRINCIPAL that access policy is written against.
  72. //
  73. // Config accepts either form:
  74. // "keys": [ {"name": "shadowman", "key": "<base64>"} ]
  75. // "keys": [ "<base64>" ] // v2.4-v2.6 shape
  76. //
  77. // A bare string still authenticates and maps to the reserved
  78. // principal `unnamed`, so existing deployments keep working; they
  79. // simply cannot be told apart in a policy until they are named.
  80. struct NamedKey {
  81. std::string name; // principal; "unnamed" when config gave a bare string
  82. std::string key; // the bearer token itself
  83. };
  84. // Bearer tokens accepted by this listener. Constant-time
  85. // compared against the value of the `authorization` gRPC
  86. // metadata header. Rotation = add new key, distribute,
  87. // remove old key. Empty list with required=true is a config
  88. // error caught at startup.
  89. std::vector<NamedKey> keys;
  90. } auth;
  91. };
  92. struct Config {
  93. std::string nodeId = "storage-1";
  94. // v2.4 — listener fleet. Empty means "synthesize a single
  95. // back-compat listener from the legacy bindAddress + rpcPort
  96. // fields below". Once any listener is configured here, the
  97. // legacy fields are ignored.
  98. std::vector<ListenerConfig> listeners;
  99. // v2.3 legacy single-listener fields. Kept so existing
  100. // config.json files continue to work — the config loader
  101. // synthesises a ListenerConfig from these when `listeners` is
  102. // empty. Operators who want TLS or auth must use the
  103. // `listeners[]` schema instead of these.
  104. std::string bindAddress = "0.0.0.0";
  105. uint16_t rpcPort = 9004;
  106. // Bootstrap/service registration
  107. std::string webapiUrl;
  108. std::string serviceKey;
  109. uint32_t heartbeatIntervalSec = 30;
  110. std::filesystem::path dataDirectory;
  111. std::filesystem::path keyFilePath;
  112. // Memory eviction settings
  113. uint64_t maxMemoryMb = 800; // Max memory before eviction
  114. uint32_t evictionThresholdPercent = 80; // Start evicting at this %
  115. uint32_t evictionTargetPercent = 60; // Evict down to this %
  116. uint32_t evictionCheckIntervalMs = 5000; // How often to check
  117. // NEW v1.7.0 eviction tuning — mirrors MemoryStore::Config fields.
  118. // Plumbed here so the config loader can set them and setupComponents()
  119. // can copy them into the MemoryStore. Schema-only in T2; behavior lands
  120. // in T3/T4.
  121. uint32_t evictionChunkSize = 1000;
  122. uint32_t evictionChunkPauseMs = 50;
  123. uint32_t maxEvictionPassesPerTrigger = 20;
  124. uint32_t hotWriteFloorMs = 30000;
  125. uint32_t memorySoftPercent = 70;
  126. uint32_t memoryHardPercent = 85;
  127. uint32_t memoryEmergencyPercent = 95;
  128. // v1.7.0 T10 — eviction burst event threshold (docs per tick).
  129. uint32_t evictionBurstThreshold = 10000;
  130. // v2.4.3 — cap on the share of the resident set one pressure episode
  131. // may evict. 0 disables. See MemoryStore::Config for the rationale.
  132. uint32_t evictionMaxEpisodePercent = 50;
  133. // Persistence settings
  134. uint32_t walSyncIntervalMs = 100;
  135. uint32_t snapshotIntervalSec = 3600;
  136. bool compressionEnabled = true;
  137. // Full persistence manager config (recovery mode, etc.).
  138. // DatabaseService::setupComponents() overlays the simpler fields above
  139. // onto this struct before constructing PersistenceManager.
  140. PersistenceManager::Config persistenceConfig;
  141. // File storage settings
  142. uint64_t maxFileSizeMb = 500;
  143. std::vector<std::string> allowedFileTypes;
  144. uint32_t fileCleanupIntervalSec = 3600;
  145. // v2.8.0 — default file retention per type. See
  146. // FileManager::Config::defaultTtlSecondsByType for the key syntax.
  147. std::unordered_map<std::string, uint32_t> fileDefaultTtlByType;
  148. // Encryption settings
  149. bool encryptionEnabled = true;
  150. bool autoGenerateKey = true;
  151. // Replication settings
  152. bool replicationEnabled = true;
  153. std::vector<std::string> peerAddresses;
  154. std::string conflictResolution = "last_writer_wins";
  155. // Migration settings
  156. struct MigrationsConfig {
  157. bool enabled = true;
  158. std::filesystem::path directory;
  159. bool autoApply = true;
  160. bool failOnError = true;
  161. } migrations;
  162. // gRPC server settings (v1.6.2 — concurrency cap + configurable sizes)
  163. struct GrpcConfig {
  164. uint32_t maxReceiveMessageSizeMb = 100;
  165. uint32_t maxSendMessageSizeMb = 100;
  166. uint32_t resourceQuotaMemoryMb = 256;
  167. uint32_t maxConcurrentSubscribeStreams = 50;
  168. uint32_t maxConcurrentFileStreams = 10;
  169. };
  170. GrpcConfig grpc;
  171. };
  172. explicit DatabaseService(Config config);
  173. ~DatabaseService();
  174. /**
  175. * Initialize the service.
  176. * Loads configuration, initializes components, and prepares for startup.
  177. */
  178. bool initialize();
  179. /**
  180. * Start the service.
  181. * Starts the gRPC server and all background threads.
  182. */
  183. void start();
  184. /**
  185. * Stop the service.
  186. * Gracefully shuts down all components.
  187. */
  188. void stop();
  189. /**
  190. * Wait for the service to stop.
  191. */
  192. void wait();
  193. /**
  194. * Check if the service is running.
  195. */
  196. [[nodiscard]] bool isRunning() const { return running_; }
  197. /**
  198. * Signal the service to stop (for signal handlers).
  199. */
  200. void signalStop();
  201. /**
  202. * Get service statistics.
  203. */
  204. [[nodiscard]] nlohmann::json getStats() const;
  205. /**
  206. * Access the view manager (for RPC handlers, migrations, etc.).
  207. */
  208. ViewManager& viewManager() { return *view_manager_; }
  209. PolicyManager& policyManager() { return *policy_manager_; }
  210. /** Current read-only state (atomic, lock-free read). */
  211. bool isReadOnly() const { return read_only_.load(std::memory_order_acquire); }
  212. /** Human-readable reason the DB is read-only (empty when writable). */
  213. std::string readOnlyReason() const;
  214. /** Toggle read-only state. Called by SetReadOnly RPC and startup logic. */
  215. void setReadOnly(bool value, const std::string& reason);
  216. /** Recovery outcome from the last recover() call. Used by GetReadOnlyStatus RPC. */
  217. const RecoveryOutcome& recoveryOutcome() const { return recovery_outcome_; }
  218. /** Set the force-readwrite flag (from --force-readwrite CLI arg). */
  219. void setForceReadwrite(bool value) { force_readwrite_ = value; }
  220. /**
  221. * Load configuration from a JSON file.
  222. */
  223. static Config loadConfig(const std::filesystem::path& configPath);
  224. /**
  225. * Parse configuration from an already-loaded JSON object. Used by the
  226. * conf.d drop-in loader (see config/config_loader.hpp) after it has
  227. * merged /etc/smartbotic-database/config.json with every *.json in
  228. * /etc/smartbotic-database/conf.d/.
  229. */
  230. static Config parseConfig(const nlohmann::json& json);
  231. // v2.3 — project-aware read accessors. `docStore(project)` returns
  232. // the per-project LmdbDocumentStore or nullptr if the project isn't
  233. // open (rare unless someone dropped it concurrently). The no-arg
  234. // `docStore()` defaults to the "default" project for callers still
  235. // on the v2.0-v2.2 API shape (Stage D will retire it). Health +
  236. // drift atomics gate any LMDB-first read path.
  237. smartbotic::db::storage::DocumentStore* docStore() noexcept;
  238. smartbotic::db::storage::DocumentStore* docStore(std::string_view project) noexcept;
  239. smartbotic::db::storage::ProjectStoreRegistry* projects() noexcept {
  240. return projects_.get();
  241. }
  242. bool mirrorHealthy() const noexcept {
  243. return mirror_healthy_.load(std::memory_order_acquire);
  244. }
  245. uint64_t mirrorDriftCount() const noexcept {
  246. return mirror_drift_count_.load(std::memory_order_relaxed);
  247. }
  248. // v2.11.0 close-out — drift accrued AFTER the boot path finished. What
  249. // destructive relation policies gate on; every READ gate keeps using
  250. // mirrorDriftCount() above. See MemoryStore::markMirrorDriftBaseline().
  251. uint64_t mirrorDriftSinceReady() const noexcept {
  252. return store_ == nullptr ? 0 : store_->mirrorDriftSinceBaseline();
  253. }
  254. uint64_t mirrorDriftAtReady() const noexcept {
  255. return store_ == nullptr ? 0 : store_->mirrorDriftBaseline();
  256. }
  257. // v2.11.0 T12 review (C2) — replication queueing + Subscribe event
  258. // publication, extracted out of the MemoryStore persistCallback_
  259. // lambda (setupComponents()) so a caller that deliberately bypasses
  260. // that callback can still drive both explicitly, in the same shape
  261. // the callback itself uses. This is what a cascade's `executeCascade`
  262. // (relations/relation_cascade.cpp) calls once per child mutation and
  263. // once for the parent delete, AFTER MemoryStore has been updated —
  264. // mirroring how v2.3.1 fixed replicated-entry apply by explicitly
  265. // driving `applyDualWriteMirror` rather than relying on a callback
  266. // that had already been bypassed for the same "don't double up on
  267. // WAL/mirror" reason. `collection` is the QUALIFIED name (what
  268. // `ReplicationEntry`/`DatabaseEvent` both expect); `doc` is null for a
  269. // delete. Does NOT touch WAL or the LMDB mirror — those are the
  270. // caller's job, done separately, and calling this a second time for
  271. // the same mutation would double-queue replication and double-publish
  272. // the event, so callers must call it exactly once per mutation.
  273. void notifyReplicationAndEvents(const std::string& collection,
  274. const std::string& id,
  275. const std::optional<Document>& doc,
  276. EventType eventType);
  277. // v2.3 Stage F — project CRUD entry points (delegate to registry).
  278. // The registry is the source of truth for listing / creating / dropping.
  279. // CreateProject is idempotent. DropProject refuses "default". The
  280. // gRPC handlers in database_grpc_impl.cpp call these.
  281. std::vector<std::string> listProjects() const;
  282. bool createProject(const std::string& name, std::string& error);
  283. bool dropProject(const std::string& name, std::string& error);
  284. private:
  285. void setupComponents();
  286. void startGrpcServer();
  287. void stopGrpcServer();
  288. void applyReplicatedEntry(const databasepb::ReplicationEntry& entry);
  289. bool runMigrations();
  290. // v2.0 Stage 4 — synchronous backfill of MemoryStore into doc_store_.
  291. // Runs once at boot, after recovery + migrations, before gRPC accepts
  292. // traffic. Skips when _meta.schema_version=2 marker exists (the env
  293. // is already populated by migrate_v1_to_v2). Returns false only on
  294. // catastrophic failure; per-doc errors flip mirror_healthy_=false but
  295. // do not abort boot (the read-flip downstream is the gate).
  296. bool backfillIntoDocStore();
  297. // v2.4.4 — audit every open project env for documents whose declared
  298. // `collection` disagrees with the sub-db they occupy. Logs ERROR per
  299. // affected project and leaves a summary; never fails startup.
  300. void auditSubdbPlacement();
  301. // v2.9.0 — arm the write path and query planner with the index declarations
  302. // persisted in _collection_meta. Must run at boot, after config load.
  303. void applyIndexDeclarations();
  304. // v2.11.0 close-out — MemoryStore::TtlExpiryRelationHook. Decides what the
  305. // TTL sweeper must do about an expiring parent's children, and PERFORMS the
  306. // cascade itself when one is required, by calling the ordinary
  307. // relations/relation_cascade.cpp machinery (WAL-before-LMDB) rather than a
  308. // second, sweeper-local implementation - a hand-rolled cascade that wrote
  309. // only LMDB would have its child deletions resurrected by the next boot's
  310. // WAL replay, since MemoryStore is rebuilt from snapshot + WAL and NOT from
  311. // LMDB.
  312. //
  313. // Runs on the cleanup thread with NO MemoryStore lock held (that is
  314. // MemoryStore::expireDocuments()' phase-2 contract, and it is what makes
  315. // calling back into MemoryStore here safe). Installed in setupComponents();
  316. // must never throw (the sweeper catches, but a throw per document would
  317. // stop expiry making progress).
  318. MemoryStore::TtlExpiryAction ttlExpiryRelationDecision(
  319. const std::string& qualifiedCollection, const std::string& id, bool firstAttempt);
  320. // v2.11.0 T6a — arm each project's LmdbDocumentStore with the relation
  321. // declarations loaded by relation_manager_->loadFromStore(), grouped by
  322. // CHILD collection. LmdbDocumentStore has no knowledge of RelationManager
  323. // and maintains the reverse index for nothing until told to via
  324. // set_relations() (see Task 3's report) — this is what tells it, once at
  325. // boot. Task 6b's createRelation/dropRelation RPCs must call
  326. // set_relations() again on live mutation; this only covers what was
  327. // already persisted at startup.
  328. // `afterMigrations` only affects LOGGING (see the self-heal block): a
  329. // relation whose reverse index sub-db is absent is a fault at boot and the
  330. // expected state for a relation a migration just declared, and the same code
  331. // handles both.
  332. void applyRelationDeclarations(bool afterMigrations = false);
  333. // v2.3 Stage C — atomic rename of <dataDir>/env/ into
  334. // <dataDir>/projects/default/env/ when the v2.2 layout is detected
  335. // and the new layout doesn't yet exist. Idempotent. Refuses to start
  336. // if both exist (likely operator hand-mess); the operator must
  337. // resolve manually.
  338. void migrateLegacyEnvToDefaultProject();
  339. Config config_;
  340. std::atomic<bool> running_{false};
  341. std::atomic<bool> stopRequested_{false};
  342. // Components
  343. std::unique_ptr<MemoryStore> store_;
  344. // v2.3 — multi-project storage substrate. Each project is one
  345. // LmdbEnv at `<dataDir>/projects/<name>/env/`. The registry holds
  346. // them in a map keyed by project name, lazy-opens new projects on
  347. // first write, and is the single source of truth for "which projects
  348. // exist." Old v2.0-v2.2 installs auto-migrate the legacy
  349. // `<dataDir>/env/` into `<dataDir>/projects/default/env/` at first
  350. // boot (see migrateLegacyEnvToDefaultProject).
  351. std::unique_ptr<smartbotic::db::storage::ProjectStoreRegistry> projects_;
  352. // v2.0 dual-write health. The persist callback mirrors every user-
  353. // collection write into doc_store_; on failure we log ERROR, bump the
  354. // drift counter, and flip mirror_healthy_ to false. Downstream Stage 4
  355. // tasks that flip reads onto doc_store_ MUST check mirror_healthy_
  356. // before engaging — a degraded mirror means LMDB is missing writes.
  357. std::atomic<bool> mirror_healthy_{true};
  358. std::atomic<uint64_t> mirror_drift_count_{0};
  359. std::unique_ptr<PersistenceManager> persistence_;
  360. std::unique_ptr<HistoryStore> history_store_; // v1.9.0 — disk-resident history
  361. std::unique_ptr<EncryptionManager> encryption_;
  362. std::unique_ptr<EventManager> events_;
  363. std::unique_ptr<FileManager> files_;
  364. std::unique_ptr<ReplicationManager> replication_;
  365. std::unique_ptr<ViewManager> view_manager_;
  366. std::unique_ptr<CollectionConfigManager> config_manager_;
  367. // v2.8.0 — per-project access policy. Owned here so its cache is loaded
  368. // once, after recovery, alongside ViewManager and CollectionConfigManager.
  369. std::unique_ptr<PolicyManager> policy_manager_;
  370. // v2.11.0 T6a — relation declarations. Owned here so its cache is loaded
  371. // once, after recovery, alongside ViewManager/CollectionConfigManager/PolicyManager.
  372. std::unique_ptr<RelationManager> relation_manager_;
  373. // gRPC
  374. std::unique_ptr<DatabaseGrpcImpl> storageImpl_;
  375. std::unique_ptr<DatabaseReplicationGrpcImpl> replicationImpl_;
  376. // v2.4 — one gRPC Server per listener (each one ports.json entry
  377. // builds a separate Server instance with its own TLS / auth
  378. // policy). Pre-2.4 there was a single grpcServer_ wired to
  379. // bindAddress + rpcPort; the back-compat shim in parseConfig
  380. // synthesises exactly one entry for v2.3-shape configs so the
  381. // observable behavior is unchanged on upgrade.
  382. std::vector<std::unique_ptr<grpc::Server>> grpcServers_;
  383. std::thread serverThread_;
  384. // Migration runner
  385. std::unique_ptr<MigrationRunner> migrationRunner_;
  386. // Read-only runtime state (v1.6.1 — auto-readonly on non-trivial recovery)
  387. std::atomic<bool> read_only_{false};
  388. mutable std::mutex reason_mutex_;
  389. std::string read_only_reason_;
  390. RecoveryOutcome recovery_outcome_;
  391. bool force_readwrite_ = false;
  392. };
  393. } // namespace smartbotic::database