database_service.cpp 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609
  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.max_message_size_mb = cfg.getOr<int>("max_message_size_mb", 64);
  308. config.wal.sync_interval_ms = cfg.getOr<int>("persistence.wal_sync_interval_ms", 100);
  309. config.snapshot.interval_sec = cfg.getOr<int>("persistence.snapshot_interval_sec", 3600);
  310. }
  311. return config;
  312. }
  313. void DatabaseService::start() {
  314. if (running_) {
  315. return;
  316. }
  317. LOG_INFO("Starting database service...");
  318. // Create data directory
  319. std::filesystem::create_directories(config_.data_directory);
  320. // Recovery from persistence
  321. recoveryFromPersistence();
  322. // Create predefined collections
  323. createPredefinedCollections();
  324. // Start WAL
  325. wal_->start();
  326. // Start gRPC server
  327. service_impl_ = std::make_unique<DatabaseServiceImpl>(store_);
  328. // Disable health check to avoid async completion queue issues
  329. // grpc::EnableDefaultHealthCheckService(true);
  330. grpc::ServerBuilder builder;
  331. builder.AddListeningPort("0.0.0.0:" + std::to_string(config_.grpc_port),
  332. grpc::InsecureServerCredentials());
  333. builder.RegisterService(service_impl_.get());
  334. // Set max message size to handle large execution results
  335. int max_msg_size = config_.max_message_size_mb * 1024 * 1024;
  336. builder.SetMaxReceiveMessageSize(max_msg_size);
  337. builder.SetMaxSendMessageSize(max_msg_size);
  338. LOG_INFO("Max gRPC message size: {} MB", config_.max_message_size_mb);
  339. server_ = builder.BuildAndStart();
  340. LOG_INFO("Database gRPC server listening on port {}", config_.grpc_port);
  341. running_ = true;
  342. // Start background threads
  343. snapshot_thread_ = std::thread(&DatabaseService::snapshotLoop, this);
  344. ttl_thread_ = std::thread(&DatabaseService::ttlLoop, this);
  345. }
  346. void DatabaseService::stop() {
  347. if (!running_) {
  348. return;
  349. }
  350. LOG_INFO("Stopping database service...");
  351. // Signal shutdown to background threads
  352. {
  353. std::lock_guard<std::mutex> lock(shutdown_mutex_);
  354. running_ = false;
  355. }
  356. shutdown_cv_.notify_all();
  357. // Stop background threads (they will wake up immediately now)
  358. if (snapshot_thread_.joinable()) {
  359. snapshot_thread_.join();
  360. }
  361. if (ttl_thread_.joinable()) {
  362. ttl_thread_.join();
  363. }
  364. // Final snapshot before shutdown
  365. LOG_INFO("Creating final snapshot before shutdown...");
  366. createSnapshot();
  367. // Stop WAL
  368. wal_->stop();
  369. // Stop gRPC server
  370. if (server_) {
  371. server_->Shutdown();
  372. }
  373. LOG_INFO("Database service stopped");
  374. }
  375. void DatabaseService::createPredefinedCollections() {
  376. // Users collection
  377. store_.createCollection(CollectionConfig("users"));
  378. // Workflows collection
  379. store_.createCollection(CollectionConfig("workflows"));
  380. // Workflow groups collection
  381. store_.createCollection(CollectionConfig("workflow_groups"));
  382. // Executions collection with 7-day TTL
  383. store_.createCollection(
  384. CollectionConfig("executions").withTTL(7 * 24 * 60 * 60 * 1000));
  385. // Sessions collection with 1-day TTL
  386. store_.createCollection(
  387. CollectionConfig("sessions").withTTL(24 * 60 * 60 * 1000));
  388. // API keys collection
  389. store_.createCollection(CollectionConfig("api_keys"));
  390. // Credentials collection
  391. store_.createCollection(CollectionConfig("credentials"));
  392. // Runners collection
  393. store_.createCollection(CollectionConfig("runners"));
  394. // Nodes collection - stores node definitions
  395. store_.createCollection(CollectionConfig("nodes"));
  396. LOG_INFO("Created predefined collections");
  397. }
  398. void DatabaseService::createSnapshot() {
  399. auto snapshot = store_.createSnapshot();
  400. auto result = snapshot_manager_->save(snapshot);
  401. if (result.ok()) {
  402. // Truncate WAL after successful snapshot
  403. wal_->truncate(snapshot.value("sequence", 0));
  404. }
  405. }
  406. void DatabaseService::recoveryFromPersistence() {
  407. LOG_INFO("Starting recovery from persistence...");
  408. // Try to load latest snapshot
  409. auto snapshot_result = snapshot_manager_->loadLatest();
  410. if (snapshot_result.ok()) {
  411. store_.loadSnapshot(snapshot_result.value());
  412. LOG_INFO("Loaded snapshot");
  413. }
  414. // Replay WAL entries after snapshot
  415. int64_t last_sequence = snapshot_manager_->getLatestSequence();
  416. wal_->replay([this, last_sequence](const WalEntry& entry) {
  417. if (entry.sequence <= last_sequence) {
  418. return; // Skip entries already in snapshot
  419. }
  420. // Apply WAL entry
  421. switch (entry.type) {
  422. case WalEntryType::Insert: {
  423. auto doc = Document::fromJson(entry.data);
  424. store_.getCollection(entry.collection)->loadFromSnapshot({doc});
  425. break;
  426. }
  427. case WalEntryType::Update: {
  428. auto doc = Document::fromJson(entry.data);
  429. store_.getCollection(entry.collection)->loadFromSnapshot({doc});
  430. break;
  431. }
  432. case WalEntryType::Delete: {
  433. auto* col = store_.getCollection(entry.collection);
  434. if (col) {
  435. col->remove(entry.document_id, 0);
  436. }
  437. break;
  438. }
  439. case WalEntryType::CreateCollection: {
  440. CollectionConfig config(entry.collection);
  441. if (entry.data.contains("default_ttl_ms")) {
  442. config.default_ttl_ms = entry.data["default_ttl_ms"];
  443. }
  444. if (entry.data.contains("versioning")) {
  445. config.versioning = VersioningConfig::fromJson(entry.data["versioning"]);
  446. }
  447. store_.createCollection(config);
  448. break;
  449. }
  450. case WalEntryType::DropCollection: {
  451. store_.dropCollection(entry.collection);
  452. break;
  453. }
  454. case WalEntryType::StoreVersion: {
  455. // Version entries are stored in-memory, no special recovery needed
  456. // They are reconstructed from document updates during normal operation
  457. break;
  458. }
  459. case WalEntryType::DeleteVersion: {
  460. // Handled similarly to StoreVersion
  461. break;
  462. }
  463. }
  464. });
  465. LOG_INFO("Recovery complete");
  466. }
  467. void DatabaseService::snapshotLoop() {
  468. while (running_) {
  469. // Wait for shutdown signal or timeout
  470. {
  471. std::unique_lock<std::mutex> lock(shutdown_mutex_);
  472. if (shutdown_cv_.wait_for(lock,
  473. std::chrono::seconds(config_.snapshot.interval_sec),
  474. [this] { return !running_.load(); })) {
  475. // Shutdown signaled, exit loop
  476. break;
  477. }
  478. }
  479. if (!running_) break;
  480. // Check if WAL is large enough to trigger snapshot
  481. if (wal_->getSize() >= config_.snapshot.wal_size_trigger) {
  482. createSnapshot();
  483. }
  484. }
  485. }
  486. void DatabaseService::ttlLoop() {
  487. constexpr int ttl_check_interval_sec = 60; // Check every minute
  488. while (running_) {
  489. // Wait for shutdown signal or timeout
  490. {
  491. std::unique_lock<std::mutex> lock(shutdown_mutex_);
  492. if (shutdown_cv_.wait_for(lock,
  493. std::chrono::seconds(ttl_check_interval_sec),
  494. [this] { return !running_.load(); })) {
  495. // Shutdown signaled, exit loop
  496. break;
  497. }
  498. }
  499. if (!running_) break;
  500. store_.expireAllDocuments();
  501. store_.expireAllVersions();
  502. }
  503. }
  504. } // namespace smartbotic::database