database_service.cpp 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602
  1. #include "database_service.hpp"
  2. #include "logging/logger.hpp"
  3. #include "common/time_utils.hpp"
  4. #include <grpcpp/health_check_service_interface.h>
  5. namespace smartbotic::database {
  6. using namespace common;
  7. // DatabaseServiceImpl implementation
  8. DatabaseServiceImpl::DatabaseServiceImpl(MemoryStore& store)
  9. : store_(store) {}
  10. grpc::Status DatabaseServiceImpl::Get(grpc::ServerContext* context,
  11. const proto::GetRequest* request,
  12. proto::Document* response) {
  13. auto result = store_.get(request->collection(), request->id());
  14. if (result.failed()) {
  15. return grpc::Status(grpc::StatusCode::NOT_FOUND, result.error().message());
  16. }
  17. documentToProto(result.value(), response);
  18. return grpc::Status::OK;
  19. }
  20. grpc::Status DatabaseServiceImpl::Query(grpc::ServerContext* context,
  21. const proto::QueryRequest* request,
  22. proto::QueryResponse* response) {
  23. struct Query q(request->collection());
  24. // Build filters
  25. for (const auto& f : request->filters()) {
  26. Filter filter;
  27. filter.field = f.field();
  28. filter.op = protoToFilterOp(f.op());
  29. filter.value = nlohmann::json::parse(f.value());
  30. q.filters.push_back(std::move(filter));
  31. }
  32. // Build sorts
  33. for (const auto& s : request->sorts()) {
  34. Sort sort;
  35. sort.field = s.field();
  36. sort.direction = s.direction() == proto::SORT_DIRECTION_DESC ?
  37. SortDirection::Desc : SortDirection::Asc;
  38. q.sorts.push_back(std::move(sort));
  39. }
  40. // Pagination
  41. if (request->has_pagination()) {
  42. q.offset = (request->pagination().page() - 1) * request->pagination().page_size();
  43. q.limit = request->pagination().page_size();
  44. }
  45. // Projection
  46. for (const auto& field : request->fields()) {
  47. q.fields.push_back(field);
  48. }
  49. auto result = store_.query(q);
  50. for (const auto& doc : result.documents) {
  51. documentToProto(doc, response->add_documents());
  52. }
  53. auto* pagination = response->mutable_pagination();
  54. pagination->set_total_count(result.total_count);
  55. pagination->set_has_more(result.has_more);
  56. return grpc::Status::OK;
  57. }
  58. grpc::Status DatabaseServiceImpl::Insert(grpc::ServerContext* context,
  59. const proto::InsertRequest* request,
  60. proto::MutationResponse* response) {
  61. Document doc;
  62. doc.id = request->id();
  63. doc.collection = request->collection();
  64. doc.data = nlohmann::json::parse(request->data());
  65. if (request->ttl_ms() > 0) {
  66. doc.expires_at = TimeUtils::nowMs() + request->ttl_ms();
  67. }
  68. auto result = store_.insert(request->collection(), std::move(doc));
  69. if (result.failed()) {
  70. response->set_success(false);
  71. auto* error = response->mutable_error();
  72. error->set_code(static_cast<int32_t>(result.error().code()));
  73. error->set_message(result.error().message());
  74. return grpc::Status::OK;
  75. }
  76. response->set_id(result.value().id);
  77. response->set_version(result.value().version);
  78. response->set_success(true);
  79. return grpc::Status::OK;
  80. }
  81. grpc::Status DatabaseServiceImpl::Update(grpc::ServerContext* context,
  82. const proto::UpdateRequest* request,
  83. proto::MutationResponse* response) {
  84. auto data = nlohmann::json::parse(request->data());
  85. auto result = store_.update(request->collection(), request->id(), data,
  86. request->expected_version(), request->partial());
  87. if (result.failed()) {
  88. response->set_success(false);
  89. auto* error = response->mutable_error();
  90. error->set_code(static_cast<int32_t>(result.error().code()));
  91. error->set_message(result.error().message());
  92. return grpc::Status::OK;
  93. }
  94. response->set_id(result.value().id);
  95. response->set_version(result.value().version);
  96. response->set_success(true);
  97. return grpc::Status::OK;
  98. }
  99. grpc::Status DatabaseServiceImpl::Delete(grpc::ServerContext* context,
  100. const proto::DeleteRequest* request,
  101. proto::MutationResponse* response) {
  102. auto result = store_.remove(request->collection(), request->id(),
  103. request->expected_version());
  104. if (result.failed()) {
  105. response->set_success(false);
  106. auto* error = response->mutable_error();
  107. error->set_code(static_cast<int32_t>(result.error().code()));
  108. error->set_message(result.error().message());
  109. return grpc::Status::OK;
  110. }
  111. response->set_id(request->id());
  112. response->set_success(true);
  113. return grpc::Status::OK;
  114. }
  115. grpc::Status DatabaseServiceImpl::BatchInsert(grpc::ServerContext* context,
  116. const proto::BatchInsertRequest* request,
  117. proto::BatchMutationResponse* response) {
  118. int32_t success_count = 0;
  119. int32_t failure_count = 0;
  120. for (const auto& req : request->documents()) {
  121. Document doc;
  122. doc.id = req.id();
  123. doc.collection = request->collection();
  124. doc.data = nlohmann::json::parse(req.data());
  125. if (req.ttl_ms() > 0) {
  126. doc.expires_at = TimeUtils::nowMs() + req.ttl_ms();
  127. }
  128. auto result = store_.insert(request->collection(), std::move(doc));
  129. auto* mutation_result = response->add_results();
  130. if (result.ok()) {
  131. mutation_result->set_id(result.value().id);
  132. mutation_result->set_version(result.value().version);
  133. mutation_result->set_success(true);
  134. ++success_count;
  135. } else {
  136. mutation_result->set_success(false);
  137. auto* error = mutation_result->mutable_error();
  138. error->set_code(static_cast<int32_t>(result.error().code()));
  139. error->set_message(result.error().message());
  140. ++failure_count;
  141. }
  142. }
  143. response->set_success_count(success_count);
  144. response->set_failure_count(failure_count);
  145. return grpc::Status::OK;
  146. }
  147. grpc::Status DatabaseServiceImpl::BatchDelete(grpc::ServerContext* context,
  148. const proto::BatchDeleteRequest* request,
  149. proto::BatchMutationResponse* response) {
  150. int32_t success_count = 0;
  151. int32_t failure_count = 0;
  152. for (const auto& id : request->ids()) {
  153. auto result = store_.remove(request->collection(), id, 0);
  154. auto* mutation_result = response->add_results();
  155. mutation_result->set_id(id);
  156. if (result.ok()) {
  157. mutation_result->set_success(true);
  158. ++success_count;
  159. } else {
  160. mutation_result->set_success(false);
  161. auto* error = mutation_result->mutable_error();
  162. error->set_code(static_cast<int32_t>(result.error().code()));
  163. error->set_message(result.error().message());
  164. ++failure_count;
  165. }
  166. }
  167. response->set_success_count(success_count);
  168. response->set_failure_count(failure_count);
  169. return grpc::Status::OK;
  170. }
  171. grpc::Status DatabaseServiceImpl::CreateCollection(grpc::ServerContext* context,
  172. const proto::CreateCollectionRequest* request,
  173. proto::Empty* response) {
  174. CollectionConfig config(request->name());
  175. config.default_ttl_ms = request->default_ttl_ms();
  176. if (!request->schema().empty()) {
  177. config.schema = nlohmann::json::parse(request->schema());
  178. }
  179. // Handle versioning configuration
  180. if (request->has_versioning()) {
  181. VersioningConfig versioning;
  182. versioning.enabled = request->versioning().enabled();
  183. versioning.max_versions = request->versioning().max_versions();
  184. versioning.version_ttl_ms = request->versioning().version_ttl_ms();
  185. versioning.keep_on_delete = request->versioning().keep_on_delete();
  186. config.versioning = versioning;
  187. }
  188. auto result = store_.createCollection(config);
  189. if (result.failed()) {
  190. return grpc::Status(grpc::StatusCode::ALREADY_EXISTS, result.error().message());
  191. }
  192. return grpc::Status::OK;
  193. }
  194. grpc::Status DatabaseServiceImpl::DropCollection(grpc::ServerContext* context,
  195. const proto::DropCollectionRequest* request,
  196. proto::Empty* response) {
  197. auto result = store_.dropCollection(request->name());
  198. if (result.failed()) {
  199. return grpc::Status(grpc::StatusCode::NOT_FOUND, result.error().message());
  200. }
  201. return grpc::Status::OK;
  202. }
  203. grpc::Status DatabaseServiceImpl::ListCollections(grpc::ServerContext* context,
  204. const proto::ListCollectionsRequest* request,
  205. proto::ListCollectionsResponse* response) {
  206. auto collections = store_.listCollections();
  207. for (const auto& name : collections) {
  208. response->add_collections(name);
  209. }
  210. return grpc::Status::OK;
  211. }
  212. grpc::Status DatabaseServiceImpl::GetVersion(grpc::ServerContext* context,
  213. const proto::GetVersionRequest* request,
  214. proto::Document* response) {
  215. auto result = store_.getVersion(request->collection(), request->id(), request->version());
  216. if (result.failed()) {
  217. return grpc::Status(grpc::StatusCode::NOT_FOUND, result.error().message());
  218. }
  219. const auto& ver = result.value();
  220. response->set_id(ver.document_id);
  221. response->set_collection(request->collection());
  222. response->set_data(ver.data.dump());
  223. response->set_version(ver.version);
  224. response->set_created_at(ver.created_at);
  225. response->set_updated_at(ver.created_at);
  226. response->set_expires_at(ver.expires_at);
  227. return grpc::Status::OK;
  228. }
  229. grpc::Status DatabaseServiceImpl::ListVersions(grpc::ServerContext* context,
  230. const proto::ListVersionsRequest* request,
  231. proto::ListVersionsResponse* response) {
  232. int32_t limit = request->limit() > 0 ? request->limit() : 100;
  233. int32_t offset = request->offset();
  234. auto result = store_.listVersions(request->collection(), request->id(), limit, offset);
  235. for (const auto& doc : result.documents) {
  236. auto ver = DocumentVersion::fromJson(doc.data);
  237. auto* proto_ver = response->add_versions();
  238. documentVersionToProto(ver, proto_ver);
  239. }
  240. response->set_total_count(result.total_count);
  241. response->set_has_more(result.has_more);
  242. return grpc::Status::OK;
  243. }
  244. void DatabaseServiceImpl::documentToProto(const Document& doc, proto::Document* proto) {
  245. proto->set_id(doc.id);
  246. proto->set_collection(doc.collection);
  247. proto->set_data(doc.data.dump());
  248. proto->set_version(doc.version);
  249. proto->set_created_at(doc.created_at);
  250. proto->set_updated_at(doc.updated_at);
  251. proto->set_expires_at(doc.expires_at);
  252. }
  253. void DatabaseServiceImpl::documentVersionToProto(const DocumentVersion& ver, proto::DocumentVersion* proto) {
  254. proto->set_id(ver.id);
  255. proto->set_document_id(ver.document_id);
  256. proto->set_version(ver.version);
  257. proto->set_data(ver.data.dump());
  258. proto->set_created_at(ver.created_at);
  259. proto->set_expires_at(ver.expires_at);
  260. }
  261. FilterOp DatabaseServiceImpl::protoToFilterOp(proto::FilterOp op) {
  262. switch (op) {
  263. case proto::FILTER_OP_EQ: return FilterOp::Eq;
  264. case proto::FILTER_OP_NE: return FilterOp::Ne;
  265. case proto::FILTER_OP_GT: return FilterOp::Gt;
  266. case proto::FILTER_OP_GTE: return FilterOp::Gte;
  267. case proto::FILTER_OP_LT: return FilterOp::Lt;
  268. case proto::FILTER_OP_LTE: return FilterOp::Lte;
  269. case proto::FILTER_OP_IN: return FilterOp::In;
  270. case proto::FILTER_OP_NIN: return FilterOp::Nin;
  271. case proto::FILTER_OP_CONTAINS: return FilterOp::Contains;
  272. case proto::FILTER_OP_REGEX: return FilterOp::Regex;
  273. case proto::FILTER_OP_EXISTS: return FilterOp::Exists;
  274. default: return FilterOp::Eq;
  275. }
  276. }
  277. // DatabaseService implementation
  278. DatabaseService::DatabaseService(const Config& config)
  279. : config_(config) {
  280. // Setup WAL
  281. WAL::Config wal_config;
  282. wal_config.directory = std::filesystem::path(config_.data_directory) / "wal";
  283. wal_config.sync_interval_ms = config_.wal.sync_interval_ms;
  284. wal_config.enabled = true;
  285. wal_ = std::make_unique<WAL>(wal_config);
  286. // Setup snapshot manager
  287. SnapshotManager::Config snapshot_config;
  288. snapshot_config.directory = std::filesystem::path(config_.data_directory) / "snapshots";
  289. snapshot_config.interval_sec = config_.snapshot.interval_sec;
  290. snapshot_config.enabled = true;
  291. snapshot_manager_ = std::make_unique<SnapshotManager>(snapshot_config);
  292. // Setup mutation callback for WAL
  293. store_.setMutationCallback([this](const WalEntry& entry) {
  294. wal_->append(entry);
  295. });
  296. }
  297. DatabaseService::~DatabaseService() {
  298. stop();
  299. }
  300. DatabaseService::Config DatabaseService::loadConfig(const std::filesystem::path& path) {
  301. Config config;
  302. auto result = config::Config::fromFile(path);
  303. if (result.ok()) {
  304. auto& cfg = result.value();
  305. config.grpc_port = cfg.getOr<int>("grpc_port", 9001);
  306. config.data_directory = cfg.getOr<std::string>("data_directory", "./data/database");
  307. config.wal.sync_interval_ms = cfg.getOr<int>("persistence.wal_sync_interval_ms", 100);
  308. config.snapshot.interval_sec = cfg.getOr<int>("persistence.snapshot_interval_sec", 3600);
  309. }
  310. return config;
  311. }
  312. void DatabaseService::start() {
  313. if (running_) {
  314. return;
  315. }
  316. LOG_INFO("Starting database service...");
  317. // Create data directory
  318. std::filesystem::create_directories(config_.data_directory);
  319. // Recovery from persistence
  320. recoveryFromPersistence();
  321. // Create predefined collections
  322. createPredefinedCollections();
  323. // Start WAL
  324. wal_->start();
  325. // Start gRPC server
  326. service_impl_ = std::make_unique<DatabaseServiceImpl>(store_);
  327. // Disable health check to avoid async completion queue issues
  328. // grpc::EnableDefaultHealthCheckService(true);
  329. grpc::ServerBuilder builder;
  330. builder.AddListeningPort("0.0.0.0:" + std::to_string(config_.grpc_port),
  331. grpc::InsecureServerCredentials());
  332. builder.RegisterService(service_impl_.get());
  333. server_ = builder.BuildAndStart();
  334. LOG_INFO("Database gRPC server listening on port {}", config_.grpc_port);
  335. running_ = true;
  336. // Start background threads
  337. snapshot_thread_ = std::thread(&DatabaseService::snapshotLoop, this);
  338. ttl_thread_ = std::thread(&DatabaseService::ttlLoop, this);
  339. }
  340. void DatabaseService::stop() {
  341. if (!running_) {
  342. return;
  343. }
  344. LOG_INFO("Stopping database service...");
  345. // Signal shutdown to background threads
  346. {
  347. std::lock_guard<std::mutex> lock(shutdown_mutex_);
  348. running_ = false;
  349. }
  350. shutdown_cv_.notify_all();
  351. // Stop background threads (they will wake up immediately now)
  352. if (snapshot_thread_.joinable()) {
  353. snapshot_thread_.join();
  354. }
  355. if (ttl_thread_.joinable()) {
  356. ttl_thread_.join();
  357. }
  358. // Final snapshot before shutdown
  359. LOG_INFO("Creating final snapshot before shutdown...");
  360. createSnapshot();
  361. // Stop WAL
  362. wal_->stop();
  363. // Stop gRPC server
  364. if (server_) {
  365. server_->Shutdown();
  366. }
  367. LOG_INFO("Database service stopped");
  368. }
  369. void DatabaseService::createPredefinedCollections() {
  370. // Users collection
  371. store_.createCollection(CollectionConfig("users"));
  372. // Workflows collection
  373. store_.createCollection(CollectionConfig("workflows"));
  374. // Workflow groups collection
  375. store_.createCollection(CollectionConfig("workflow_groups"));
  376. // Executions collection with 7-day TTL
  377. store_.createCollection(
  378. CollectionConfig("executions").withTTL(7 * 24 * 60 * 60 * 1000));
  379. // Sessions collection with 1-day TTL
  380. store_.createCollection(
  381. CollectionConfig("sessions").withTTL(24 * 60 * 60 * 1000));
  382. // API keys collection
  383. store_.createCollection(CollectionConfig("api_keys"));
  384. // Credentials collection
  385. store_.createCollection(CollectionConfig("credentials"));
  386. // Runners collection
  387. store_.createCollection(CollectionConfig("runners"));
  388. // Nodes collection - stores node definitions
  389. store_.createCollection(CollectionConfig("nodes"));
  390. LOG_INFO("Created predefined collections");
  391. }
  392. void DatabaseService::createSnapshot() {
  393. auto snapshot = store_.createSnapshot();
  394. auto result = snapshot_manager_->save(snapshot);
  395. if (result.ok()) {
  396. // Truncate WAL after successful snapshot
  397. wal_->truncate(snapshot.value("sequence", 0));
  398. }
  399. }
  400. void DatabaseService::recoveryFromPersistence() {
  401. LOG_INFO("Starting recovery from persistence...");
  402. // Try to load latest snapshot
  403. auto snapshot_result = snapshot_manager_->loadLatest();
  404. if (snapshot_result.ok()) {
  405. store_.loadSnapshot(snapshot_result.value());
  406. LOG_INFO("Loaded snapshot");
  407. }
  408. // Replay WAL entries after snapshot
  409. int64_t last_sequence = snapshot_manager_->getLatestSequence();
  410. wal_->replay([this, last_sequence](const WalEntry& entry) {
  411. if (entry.sequence <= last_sequence) {
  412. return; // Skip entries already in snapshot
  413. }
  414. // Apply WAL entry
  415. switch (entry.type) {
  416. case WalEntryType::Insert: {
  417. auto doc = Document::fromJson(entry.data);
  418. store_.getCollection(entry.collection)->loadFromSnapshot({doc});
  419. break;
  420. }
  421. case WalEntryType::Update: {
  422. auto doc = Document::fromJson(entry.data);
  423. store_.getCollection(entry.collection)->loadFromSnapshot({doc});
  424. break;
  425. }
  426. case WalEntryType::Delete: {
  427. auto* col = store_.getCollection(entry.collection);
  428. if (col) {
  429. col->remove(entry.document_id, 0);
  430. }
  431. break;
  432. }
  433. case WalEntryType::CreateCollection: {
  434. CollectionConfig config(entry.collection);
  435. if (entry.data.contains("default_ttl_ms")) {
  436. config.default_ttl_ms = entry.data["default_ttl_ms"];
  437. }
  438. if (entry.data.contains("versioning")) {
  439. config.versioning = VersioningConfig::fromJson(entry.data["versioning"]);
  440. }
  441. store_.createCollection(config);
  442. break;
  443. }
  444. case WalEntryType::DropCollection: {
  445. store_.dropCollection(entry.collection);
  446. break;
  447. }
  448. case WalEntryType::StoreVersion: {
  449. // Version entries are stored in-memory, no special recovery needed
  450. // They are reconstructed from document updates during normal operation
  451. break;
  452. }
  453. case WalEntryType::DeleteVersion: {
  454. // Handled similarly to StoreVersion
  455. break;
  456. }
  457. }
  458. });
  459. LOG_INFO("Recovery complete");
  460. }
  461. void DatabaseService::snapshotLoop() {
  462. while (running_) {
  463. // Wait for shutdown signal or timeout
  464. {
  465. std::unique_lock<std::mutex> lock(shutdown_mutex_);
  466. if (shutdown_cv_.wait_for(lock,
  467. std::chrono::seconds(config_.snapshot.interval_sec),
  468. [this] { return !running_.load(); })) {
  469. // Shutdown signaled, exit loop
  470. break;
  471. }
  472. }
  473. if (!running_) break;
  474. // Check if WAL is large enough to trigger snapshot
  475. if (wal_->getSize() >= config_.snapshot.wal_size_trigger) {
  476. createSnapshot();
  477. }
  478. }
  479. }
  480. void DatabaseService::ttlLoop() {
  481. constexpr int ttl_check_interval_sec = 60; // Check every minute
  482. while (running_) {
  483. // Wait for shutdown signal or timeout
  484. {
  485. std::unique_lock<std::mutex> lock(shutdown_mutex_);
  486. if (shutdown_cv_.wait_for(lock,
  487. std::chrono::seconds(ttl_check_interval_sec),
  488. [this] { return !running_.load(); })) {
  489. // Shutdown signaled, exit loop
  490. break;
  491. }
  492. }
  493. if (!running_) break;
  494. store_.expireAllDocuments();
  495. store_.expireAllVersions();
  496. }
  497. }
  498. } // namespace smartbotic::database