database_service.cpp 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696
  1. #include "database_service.hpp"
  2. #include <grpcpp/grpcpp.h>
  3. #include <grpcpp/resource_quota.h>
  4. #include <spdlog/spdlog.h>
  5. #include <fstream>
  6. namespace smartbotic::database {
  7. DatabaseService::DatabaseService(Config config)
  8. : config_(std::move(config))
  9. {
  10. }
  11. DatabaseService::~DatabaseService() {
  12. stop();
  13. }
  14. std::string DatabaseService::readOnlyReason() const {
  15. std::lock_guard<std::mutex> lock(reason_mutex_);
  16. return read_only_reason_;
  17. }
  18. void DatabaseService::setReadOnly(bool value, const std::string& reason) {
  19. {
  20. std::lock_guard<std::mutex> lock(reason_mutex_);
  21. read_only_reason_ = value ? reason : "";
  22. }
  23. bool prev = read_only_.exchange(value, std::memory_order_acq_rel);
  24. if (prev != value) {
  25. if (value) {
  26. spdlog::error("Database entered READ-ONLY mode: {}", reason);
  27. } else {
  28. spdlog::info("Database unlocked -- writes accepted");
  29. }
  30. }
  31. }
  32. bool DatabaseService::initialize() {
  33. spdlog::info("Initializing database service (node: {})", config_.nodeId);
  34. try {
  35. // Create data directory if needed
  36. std::error_code ec;
  37. std::filesystem::create_directories(config_.dataDirectory, ec);
  38. if (ec) {
  39. spdlog::error("Failed to create data directory: {}", ec.message());
  40. return false;
  41. }
  42. setupComponents();
  43. // Initialize encryption
  44. if (config_.encryptionEnabled) {
  45. if (!encryption_->initialize()) {
  46. spdlog::error("Failed to initialize encryption");
  47. return false;
  48. }
  49. }
  50. // Recover from persistence
  51. recovery_outcome_ = persistence_->recover(*store_);
  52. if (recovery_outcome_.isFailure()) {
  53. const std::string modeStr =
  54. recoveryModeToString(config_.persistenceConfig.recoveryMode);
  55. spdlog::error("");
  56. spdlog::error("+------------------------------------------------------------------+");
  57. spdlog::error("| RECOVERY REFUSED |");
  58. spdlog::error("| |");
  59. spdlog::error("| Mode: {}", modeStr);
  60. spdlog::error("| Expected snapshot: {}", recovery_outcome_.expectedSnapshot.string());
  61. spdlog::error("| Reason: {}", recovery_outcome_.failureReason);
  62. spdlog::error("| |");
  63. spdlog::error("| Snapshots available: {}", recovery_outcome_.snapshotsAvailable);
  64. spdlog::error("| Snapshots attempted: {}", recovery_outcome_.snapshotsAttempted);
  65. spdlog::error("| |");
  66. spdlog::error("| To escalate, restart with one of: |");
  67. spdlog::error("| --recovery-mode=snapshot_fallback Try older snapshots |");
  68. spdlog::error("| --recovery-mode=wal_only Replay WAL only (slow) |");
  69. spdlog::error("| --recovery-mode=best_effort Try all of the above |");
  70. spdlog::error("| --recovery-mode=force_empty Start empty (LAST RESORT) |");
  71. spdlog::error("| |");
  72. spdlog::error("| Or set in config.json: \"recovery\": {{ \"mode\": \"<mode>\" }} |");
  73. spdlog::error("| |");
  74. spdlog::error("| PRESERVE /var/lib/smartbotic-database/ BEFORE ESCALATING. |");
  75. spdlog::error("+------------------------------------------------------------------+");
  76. throw std::runtime_error("recovery failed; refusing to start");
  77. }
  78. // Auto-readonly mode on non-trivial recovery
  79. if (recovery_outcome_.isNonTrivial() && !force_readwrite_) {
  80. std::string reason;
  81. switch (recovery_outcome_.kind) {
  82. case RecoveryOutcome::Kind::SnapshotFellBack:
  83. reason = "fell back to snapshot " +
  84. recovery_outcome_.snapshotUsed.filename().string() +
  85. " because " +
  86. recovery_outcome_.expectedSnapshot.filename().string() +
  87. " failed: " + recovery_outcome_.failureReason;
  88. break;
  89. case RecoveryOutcome::Kind::WalOnlyReplay:
  90. reason = "WAL-only replay, no snapshot loaded (" +
  91. std::to_string(recovery_outcome_.walEntriesReplayed) +
  92. " entries)";
  93. break;
  94. case RecoveryOutcome::Kind::ForcedEmpty:
  95. reason = "forced empty by operator (--recovery-mode=force_empty)";
  96. break;
  97. default:
  98. reason = "non-trivial recovery";
  99. break;
  100. }
  101. reason += ". Run `smartbotic-db-cli unlock` to accept this state, "
  102. "or restart with --force-readwrite to bypass this check.";
  103. setReadOnly(true, reason);
  104. spdlog::error("");
  105. spdlog::error("+------------------------------------------------------------------+");
  106. spdlog::error("| [ERROR] Database booted in READ-ONLY mode after non-trivial |");
  107. spdlog::error("| recovery. |");
  108. spdlog::error("| |");
  109. spdlog::error("| Reason: {}", reason);
  110. spdlog::error("| |");
  111. spdlog::error("| Writes will be REJECTED until you acknowledge this state: |");
  112. spdlog::error("| smartbotic-db-cli unlock # live, no restart |");
  113. spdlog::error("| smartbotic-database --force-readwrite # on next restart |");
  114. spdlog::error("+------------------------------------------------------------------+");
  115. }
  116. // Set replication sequence after recovery (WAL sequence is now known)
  117. replication_->setSequence(persistence_->currentWalSequence());
  118. // Load view definitions from the _views system collection (which is now
  119. // populated by the persistence recovery above).
  120. view_manager_->loadFromStore();
  121. // Load per-collection configs from the _collection_meta system collection.
  122. // Must happen AFTER persistence recovery and BEFORE migrations so any
  123. // migration-created documents are stamped with the correct precision.
  124. config_manager_->loadFromStore();
  125. // Run migrations if enabled
  126. if (config_.migrations.enabled && !config_.migrations.directory.empty()) {
  127. if (!runMigrations()) {
  128. if (config_.migrations.failOnError) {
  129. spdlog::error("Failed to run migrations");
  130. return false;
  131. }
  132. spdlog::warn("Some migrations failed, continuing anyway");
  133. }
  134. }
  135. spdlog::info("Database service initialized successfully");
  136. return true;
  137. } catch (const std::exception& e) {
  138. spdlog::error("Failed to initialize database service: {}", e.what());
  139. return false;
  140. }
  141. }
  142. bool DatabaseService::runMigrations() {
  143. if (!config_.migrations.enabled) {
  144. return true;
  145. }
  146. MigrationRunner::Config migrationConfig;
  147. migrationConfig.directory = config_.migrations.directory;
  148. migrationConfig.autoApply = config_.migrations.autoApply;
  149. migrationConfig.failOnError = config_.migrations.failOnError;
  150. migrationRunner_ = std::make_unique<MigrationRunner>(*store_, *view_manager_, migrationConfig);
  151. return migrationRunner_->runMigrations();
  152. }
  153. void DatabaseService::start() {
  154. if (running_.exchange(true)) {
  155. return;
  156. }
  157. spdlog::info("Starting database service on {}:{}", config_.bindAddress, config_.rpcPort);
  158. // Start components
  159. store_->start();
  160. persistence_->start();
  161. events_->start();
  162. files_->start();
  163. if (config_.replicationEnabled) {
  164. replication_->start();
  165. // Set initial local collections for discovery
  166. replication_->setLocalCollections(store_->listCollections());
  167. }
  168. // Start gRPC server
  169. startGrpcServer();
  170. spdlog::info("Database service started");
  171. }
  172. void DatabaseService::stop() {
  173. if (!running_.exchange(false)) {
  174. return;
  175. }
  176. spdlog::info("Stopping database service...");
  177. // Stop gRPC server first
  178. stopGrpcServer();
  179. // Stop components in reverse order
  180. if (replication_) {
  181. replication_->stop();
  182. }
  183. if (files_) {
  184. files_->stop();
  185. }
  186. if (events_) {
  187. events_->stop();
  188. }
  189. if (persistence_) {
  190. persistence_->stop();
  191. }
  192. if (store_) {
  193. store_->stop();
  194. }
  195. spdlog::info("Database service stopped");
  196. }
  197. void DatabaseService::wait() {
  198. if (serverThread_.joinable()) {
  199. serverThread_.join();
  200. }
  201. }
  202. void DatabaseService::signalStop() {
  203. stopRequested_ = true;
  204. stop();
  205. }
  206. nlohmann::json DatabaseService::getStats() const {
  207. nlohmann::json stats;
  208. if (store_) {
  209. auto storeStats = store_->getStats();
  210. stats["documents"] = storeStats.totalDocuments;
  211. stats["collections"] = storeStats.totalCollections;
  212. stats["memory_bytes"] = storeStats.estimatedMemoryBytes;
  213. stats["inserts"] = storeStats.insertCount;
  214. stats["updates"] = storeStats.updateCount;
  215. stats["deletes"] = storeStats.deleteCount;
  216. stats["queries"] = storeStats.queryCount;
  217. }
  218. if (persistence_) {
  219. auto persistStats = persistence_->getStats();
  220. stats["wal_sequence"] = persistStats.walSequence;
  221. stats["wal_size_bytes"] = persistStats.walSizeBytes;
  222. stats["snapshot_count"] = persistStats.snapshotCount;
  223. stats["last_snapshot_sequence"] = persistStats.lastSnapshotSequence;
  224. }
  225. if (events_) {
  226. stats["subscriptions"] = events_->subscriptionCount();
  227. }
  228. return stats;
  229. }
  230. DatabaseService::Config DatabaseService::loadConfig(const std::filesystem::path& configPath) {
  231. std::ifstream file(configPath);
  232. if (!file) {
  233. throw std::runtime_error("Failed to open config file: " + configPath.string());
  234. }
  235. nlohmann::json json;
  236. file >> json;
  237. return parseConfig(json);
  238. }
  239. DatabaseService::Config DatabaseService::parseConfig(const nlohmann::json& json) {
  240. Config config;
  241. // Check for both "storage" and "database" keys for backward compatibility
  242. const nlohmann::json* dbConfig = nullptr;
  243. if (json.contains("database")) {
  244. dbConfig = &json["database"];
  245. } else if (json.contains("storage")) {
  246. dbConfig = &json["storage"];
  247. }
  248. if (dbConfig) {
  249. const auto& db = *dbConfig;
  250. config.nodeId = db.value("node_id", config.nodeId);
  251. config.bindAddress = db.value("bind_address", config.bindAddress);
  252. config.rpcPort = db.value("rpc_port", config.rpcPort);
  253. // Expand environment variables in data directory
  254. std::string dataDir = db.value("data_directory", "");
  255. if (dataDir.find("${HOME}") != std::string::npos) {
  256. const char* home = std::getenv("HOME");
  257. if (home) {
  258. size_t pos = dataDir.find("${HOME}");
  259. dataDir.replace(pos, 7, home);
  260. }
  261. }
  262. config.dataDirectory = dataDir;
  263. // Migrations settings
  264. if (db.contains("migrations")) {
  265. auto& migrations = db["migrations"];
  266. config.migrations.enabled = migrations.value("enabled", config.migrations.enabled);
  267. config.migrations.autoApply = migrations.value("auto_apply", config.migrations.autoApply);
  268. config.migrations.failOnError = migrations.value("fail_on_error", config.migrations.failOnError);
  269. std::string migDir = migrations.value("directory", "");
  270. if (migDir.find("${HOME}") != std::string::npos) {
  271. const char* home = std::getenv("HOME");
  272. if (home) {
  273. size_t pos = migDir.find("${HOME}");
  274. migDir.replace(pos, 7, home);
  275. }
  276. }
  277. config.migrations.directory = migDir;
  278. }
  279. // Memory eviction settings
  280. if (db.contains("memory")) {
  281. auto& memory = db["memory"];
  282. config.maxMemoryMb = memory.value("max_memory_mb", config.maxMemoryMb);
  283. config.evictionThresholdPercent = memory.value("eviction_threshold_percent", config.evictionThresholdPercent);
  284. config.evictionTargetPercent = memory.value("eviction_target_percent", config.evictionTargetPercent);
  285. config.evictionCheckIntervalMs = memory.value("eviction_check_interval_ms", config.evictionCheckIntervalMs);
  286. // NEW v1.7.0 eviction tuning (schema-only in T2; runtime use lands in T3/T4)
  287. config.evictionChunkSize = memory.value("eviction_chunk_size", config.evictionChunkSize);
  288. config.evictionChunkPauseMs = memory.value("eviction_chunk_pause_ms", config.evictionChunkPauseMs);
  289. config.maxEvictionPassesPerTrigger = memory.value("max_eviction_passes_per_trigger", config.maxEvictionPassesPerTrigger);
  290. config.hotWriteFloorMs = memory.value("hot_write_floor_ms", config.hotWriteFloorMs);
  291. config.memorySoftPercent = memory.value("memory_soft_percent", config.memorySoftPercent);
  292. config.memoryHardPercent = memory.value("memory_hard_percent", config.memoryHardPercent);
  293. config.memoryEmergencyPercent = memory.value("memory_emergency_percent", config.memoryEmergencyPercent);
  294. }
  295. // Persistence settings
  296. if (db.contains("persistence")) {
  297. auto& persistence = db["persistence"];
  298. config.walSyncIntervalMs = persistence.value("wal_sync_interval_ms", config.walSyncIntervalMs);
  299. config.snapshotIntervalSec = persistence.value("snapshot_interval_sec", config.snapshotIntervalSec);
  300. config.compressionEnabled = persistence.value("compression", "lz4") != "none";
  301. // Snapshot durability (NEW in v1.6.1)
  302. if (persistence.contains("snapshots")) {
  303. const auto& snap = persistence["snapshots"];
  304. config.persistenceConfig.validateAfterWrite = snap.value("validate_after_write", config.persistenceConfig.validateAfterWrite);
  305. config.persistenceConfig.cleanupOnlyIfVerified = snap.value("cleanup_only_if_verified", config.persistenceConfig.cleanupOnlyIfVerified);
  306. }
  307. // Recovery (NEW in v1.6.1)
  308. if (persistence.contains("recovery")) {
  309. const auto& rec = persistence["recovery"];
  310. std::string modeStr = rec.value("mode", std::string("normal"));
  311. try {
  312. config.persistenceConfig.recoveryMode = recoveryModeFromString(modeStr);
  313. } catch (const std::exception& e) {
  314. spdlog::warn("Invalid recovery.mode '{}', defaulting to 'normal': {}",
  315. modeStr, e.what());
  316. config.persistenceConfig.recoveryMode = RecoveryMode::Normal;
  317. }
  318. config.persistenceConfig.autoEscalate = rec.value("auto_escalate", config.persistenceConfig.autoEscalate);
  319. config.persistenceConfig.allowEmptyOnFreshInstall = rec.value("allow_empty_on_fresh_install", config.persistenceConfig.allowEmptyOnFreshInstall);
  320. }
  321. }
  322. // File settings
  323. if (db.contains("files")) {
  324. auto& files = db["files"];
  325. config.maxFileSizeMb = files.value("max_file_size_mb", config.maxFileSizeMb);
  326. config.fileCleanupIntervalSec = files.value("cleanup_orphans_interval_sec", config.fileCleanupIntervalSec);
  327. if (files.contains("allowed_types")) {
  328. for (const auto& type : files["allowed_types"]) {
  329. config.allowedFileTypes.push_back(type.get<std::string>());
  330. }
  331. }
  332. }
  333. // Encryption settings
  334. if (db.contains("encryption")) {
  335. auto& encryption = db["encryption"];
  336. config.encryptionEnabled = encryption.value("enabled", config.encryptionEnabled);
  337. config.autoGenerateKey = encryption.value("auto_generate_key", config.autoGenerateKey);
  338. std::string keyFile = encryption.value("key_file", "");
  339. if (keyFile.find("${HOME}") != std::string::npos) {
  340. const char* home = std::getenv("HOME");
  341. if (home) {
  342. size_t pos = keyFile.find("${HOME}");
  343. keyFile.replace(pos, 7, home);
  344. }
  345. }
  346. config.keyFilePath = keyFile;
  347. }
  348. // Replication settings
  349. if (db.contains("replication")) {
  350. auto& replication = db["replication"];
  351. config.replicationEnabled = replication.value("enabled", config.replicationEnabled);
  352. config.conflictResolution = replication.value("conflict_resolution", config.conflictResolution);
  353. if (replication.contains("peers")) {
  354. for (const auto& peer : replication["peers"]) {
  355. config.peerAddresses.push_back(peer.get<std::string>());
  356. }
  357. }
  358. }
  359. // gRPC settings (v1.6.2 — concurrency cap + configurable message sizes)
  360. if (db.contains("grpc")) {
  361. const auto& grpc = db["grpc"];
  362. config.grpc.maxReceiveMessageSizeMb =
  363. grpc.value("max_receive_message_size_mb", config.grpc.maxReceiveMessageSizeMb);
  364. config.grpc.maxSendMessageSizeMb =
  365. grpc.value("max_send_message_size_mb", config.grpc.maxSendMessageSizeMb);
  366. config.grpc.resourceQuotaMemoryMb =
  367. grpc.value("resource_quota_memory_mb", config.grpc.resourceQuotaMemoryMb);
  368. config.grpc.maxConcurrentSubscribeStreams =
  369. grpc.value("max_concurrent_subscribe_streams", config.grpc.maxConcurrentSubscribeStreams);
  370. config.grpc.maxConcurrentFileStreams =
  371. grpc.value("max_concurrent_file_streams", config.grpc.maxConcurrentFileStreams);
  372. }
  373. }
  374. return config;
  375. }
  376. void DatabaseService::setupComponents() {
  377. // Create memory store with eviction config
  378. MemoryStore::Config storeConfig;
  379. storeConfig.nodeId = config_.nodeId;
  380. storeConfig.maxMemoryBytes = config_.maxMemoryMb * 1024ULL * 1024ULL;
  381. storeConfig.evictionThresholdPercent = config_.evictionThresholdPercent;
  382. storeConfig.evictionTargetPercent = config_.evictionTargetPercent;
  383. storeConfig.evictionCheckIntervalMs = config_.evictionCheckIntervalMs;
  384. // NEW v1.7.0 eviction tuning (wired in T2; consumed in T3/T4).
  385. storeConfig.evictionChunkSize = config_.evictionChunkSize;
  386. storeConfig.evictionChunkPauseMs = config_.evictionChunkPauseMs;
  387. storeConfig.maxEvictionPassesPerTrigger = config_.maxEvictionPassesPerTrigger;
  388. storeConfig.hotWriteFloorMs = config_.hotWriteFloorMs;
  389. storeConfig.memorySoftPercent = config_.memorySoftPercent;
  390. storeConfig.memoryHardPercent = config_.memoryHardPercent;
  391. storeConfig.memoryEmergencyPercent = config_.memoryEmergencyPercent;
  392. store_ = std::make_unique<MemoryStore>(storeConfig);
  393. // Create view manager (cache loaded in initialize() after persistence recovery)
  394. view_manager_ = std::make_unique<ViewManager>(*store_);
  395. // Create per-collection config manager and attach it to the store so the
  396. // write paths route document timestamp stamps through it.
  397. config_manager_ = std::make_unique<CollectionConfigManager>(*store_);
  398. store_->setConfigManager(config_manager_.get());
  399. // Create persistence manager. Start from any persistenceConfig values the
  400. // caller (config loader / CLI parser) has already populated — including
  401. // recoveryMode — and overlay the top-level convenience fields.
  402. PersistenceManager::Config persistConfig = config_.persistenceConfig;
  403. persistConfig.dataDir = config_.dataDirectory;
  404. persistConfig.walSyncIntervalMs = config_.walSyncIntervalMs;
  405. persistConfig.snapshotIntervalSec = config_.snapshotIntervalSec;
  406. persistConfig.compressionEnabled = config_.compressionEnabled;
  407. // Mirror the resolved recoveryMode back into our Config so downstream code
  408. // (logging, RPC handlers) can inspect config_.persistenceConfig.recoveryMode
  409. // without having to reach into the PersistenceManager.
  410. config_.persistenceConfig = persistConfig;
  411. persistence_ = std::make_unique<PersistenceManager>(persistConfig);
  412. // Connect store callbacks to persistence - uses setPersistCallback for WAL logging
  413. store_->setPersistCallback([this](const std::string& collection, const std::string& id,
  414. const std::optional<Document>& doc, EventType eventType) {
  415. // Log to WAL based on event type
  416. switch (eventType) {
  417. case EventType::INSERT:
  418. if (doc) persistence_->logInsert(collection, *doc);
  419. break;
  420. case EventType::UPDATE:
  421. if (doc) persistence_->logUpdate(collection, *doc);
  422. break;
  423. case EventType::DELETE:
  424. persistence_->logDelete(collection, id);
  425. break;
  426. default:
  427. break;
  428. }
  429. // Queue for replication broadcast
  430. if (replication_ && config_.replicationEnabled) {
  431. databasepb::ReplicationEntry entry;
  432. entry.set_collection(collection);
  433. entry.set_document_id(id);
  434. entry.set_global_timestamp(
  435. std::chrono::duration_cast<std::chrono::milliseconds>(
  436. std::chrono::system_clock::now().time_since_epoch()
  437. ).count());
  438. switch (eventType) {
  439. case EventType::INSERT:
  440. entry.set_op(databasepb::OP_INSERT);
  441. if (doc) entry.set_data(doc->data.dump());
  442. break;
  443. case EventType::UPDATE:
  444. entry.set_op(databasepb::OP_UPDATE);
  445. if (doc) entry.set_data(doc->data.dump());
  446. break;
  447. case EventType::DELETE:
  448. entry.set_op(databasepb::OP_DELETE);
  449. break;
  450. default:
  451. break;
  452. }
  453. replication_->queueForReplication(entry);
  454. }
  455. // Publish event
  456. if (events_) {
  457. DatabaseEvent event;
  458. event.type = eventType;
  459. event.collection = collection;
  460. event.documentId = id;
  461. event.timestamp = std::chrono::duration_cast<std::chrono::milliseconds>(
  462. std::chrono::system_clock::now().time_since_epoch()
  463. ).count();
  464. event.nodeId = config_.nodeId;
  465. if (doc) {
  466. event.data = doc->data;
  467. }
  468. events_->publish(event);
  469. }
  470. });
  471. // Set up document load callback for LRU eviction recovery
  472. store_->setDocumentLoadCallback([this](const std::string& collection,
  473. const std::string& id,
  474. uint64_t walSequence) -> std::optional<Document> {
  475. // Load document from WAL via persistence manager
  476. return persistence_->loadDocument(collection, id, 0); // Search from beginning
  477. });
  478. // Create encryption manager
  479. EncryptionManager::Config encryptConfig;
  480. encryptConfig.enabled = config_.encryptionEnabled;
  481. encryptConfig.keyFilePath = config_.keyFilePath;
  482. encryptConfig.autoGenerateKey = config_.autoGenerateKey;
  483. encryption_ = std::make_unique<EncryptionManager>(encryptConfig);
  484. // Create event manager
  485. EventManager::Config eventConfig;
  486. eventConfig.nodeId = config_.nodeId;
  487. events_ = std::make_unique<EventManager>(eventConfig);
  488. // Create file manager
  489. FileManager::Config fileConfig;
  490. fileConfig.filesDir = config_.dataDirectory / "files";
  491. fileConfig.maxFileSizeMb = config_.maxFileSizeMb;
  492. fileConfig.allowedMimeTypes = config_.allowedFileTypes;
  493. fileConfig.cleanupIntervalSec = config_.fileCleanupIntervalSec;
  494. files_ = std::make_unique<FileManager>(fileConfig);
  495. // Create replication manager
  496. ReplicationManager::Config replConfig;
  497. replConfig.nodeId = config_.nodeId;
  498. replConfig.peerAddresses = config_.peerAddresses;
  499. replConfig.conflictResolution = config_.conflictResolution;
  500. replication_ = std::make_unique<ReplicationManager>(replConfig);
  501. // Wire replication entry apply callback
  502. replication_->setEntryApplyCallback([this](const databasepb::ReplicationEntry& entry) {
  503. applyReplicatedEntry(entry);
  504. });
  505. // Create gRPC implementations
  506. storageImpl_ = std::make_unique<DatabaseGrpcImpl>(
  507. *this, *store_, *persistence_, *events_, *files_, *encryption_, *view_manager_, *config_manager_
  508. );
  509. // v1.6.2 — wire the streaming-RPC concurrency limits from GrpcConfig.
  510. storageImpl_->setStreamLimits(
  511. config_.grpc.maxConcurrentSubscribeStreams,
  512. config_.grpc.maxConcurrentFileStreams);
  513. replicationImpl_ = std::make_unique<DatabaseReplicationGrpcImpl>(
  514. *store_, *replication_, *persistence_
  515. );
  516. }
  517. void DatabaseService::applyReplicatedEntry(const databasepb::ReplicationEntry& entry) {
  518. try {
  519. switch (entry.op()) {
  520. case databasepb::OP_INSERT:
  521. case databasepb::OP_UPDATE:
  522. case databasepb::OP_UPSERT: {
  523. if (entry.data().empty()) {
  524. spdlog::warn("Replication entry has no data for op {}", static_cast<int>(entry.op()));
  525. return;
  526. }
  527. auto json = nlohmann::json::parse(entry.data());
  528. Document doc = Document::fromJson(json);
  529. doc.id = entry.document_id();
  530. doc.nodeId = entry.node_id();
  531. // Use loadDocument to bypass normal callbacks (avoid re-replication)
  532. store_->loadDocument(entry.collection(), doc);
  533. spdlog::trace("Applied replicated {} to {}/{} from {}",
  534. entry.op() == databasepb::OP_INSERT ? "insert" :
  535. entry.op() == databasepb::OP_UPDATE ? "update" : "upsert",
  536. entry.collection(), entry.document_id(), entry.node_id());
  537. break;
  538. }
  539. case databasepb::OP_DELETE: {
  540. store_->remove(entry.collection(), entry.document_id());
  541. spdlog::trace("Applied replicated delete to {}/{} from {}",
  542. entry.collection(), entry.document_id(), entry.node_id());
  543. break;
  544. }
  545. case databasepb::OP_CREATE_COLLECTION: {
  546. CollectionOptions options;
  547. store_->createCollection(entry.collection(), options);
  548. spdlog::trace("Applied replicated create collection {} from {}",
  549. entry.collection(), entry.node_id());
  550. break;
  551. }
  552. case databasepb::OP_DROP_COLLECTION: {
  553. store_->dropCollection(entry.collection());
  554. spdlog::trace("Applied replicated drop collection {} from {}",
  555. entry.collection(), entry.node_id());
  556. break;
  557. }
  558. default:
  559. spdlog::warn("Unknown replication operation type: {}", static_cast<int>(entry.op()));
  560. break;
  561. }
  562. } catch (const std::exception& e) {
  563. spdlog::error("Failed to apply replicated entry: {}", e.what());
  564. }
  565. }
  566. void DatabaseService::startGrpcServer() {
  567. std::string serverAddress = config_.bindAddress + ":" + std::to_string(config_.rpcPort);
  568. grpc::ServerBuilder builder;
  569. builder.AddListeningPort(serverAddress, grpc::InsecureServerCredentials());
  570. builder.RegisterService(storageImpl_.get());
  571. builder.RegisterService(replicationImpl_.get());
  572. // Configurable message sizes (default 100 MB for file uploads)
  573. const int64_t mb = 1024 * 1024;
  574. builder.SetMaxReceiveMessageSize(
  575. static_cast<int>(config_.grpc.maxReceiveMessageSizeMb * mb));
  576. builder.SetMaxSendMessageSize(
  577. static_cast<int>(config_.grpc.maxSendMessageSizeMb * mb));
  578. // ResourceQuota bounds total inbound buffer memory across all RPCs
  579. grpc::ResourceQuota quota("smartbotic-db");
  580. quota.Resize(config_.grpc.resourceQuotaMemoryMb * mb);
  581. builder.SetResourceQuota(quota);
  582. spdlog::info("gRPC config: recv={}MB send={}MB quota={}MB subs={} files={}",
  583. config_.grpc.maxReceiveMessageSizeMb,
  584. config_.grpc.maxSendMessageSizeMb,
  585. config_.grpc.resourceQuotaMemoryMb,
  586. config_.grpc.maxConcurrentSubscribeStreams,
  587. config_.grpc.maxConcurrentFileStreams);
  588. grpcServer_ = builder.BuildAndStart();
  589. if (!grpcServer_) {
  590. throw std::runtime_error("Failed to start gRPC server on " + serverAddress);
  591. }
  592. spdlog::info("gRPC server listening on {}", serverAddress);
  593. }
  594. void DatabaseService::stopGrpcServer() {
  595. if (grpcServer_) {
  596. grpcServer_->Shutdown();
  597. grpcServer_.reset();
  598. }
  599. }
  600. } // namespace smartbotic::database