database_service.hpp 15 KB

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