client.cpp 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882
  1. #include "smartbotic/database/client.hpp"
  2. #include <database.grpc.pb.h>
  3. #include <grpcpp/grpcpp.h>
  4. #include <spdlog/spdlog.h>
  5. #include <atomic>
  6. #include <mutex>
  7. #include <thread>
  8. namespace smartbotic::database {
  9. // ===== PIMPL Implementation =====
  10. class Client::Impl {
  11. public:
  12. explicit Impl(Config config) : config_(std::move(config)) {}
  13. ~Impl() {
  14. disconnect();
  15. }
  16. bool connect() {
  17. try {
  18. auto channelArgs = grpc::ChannelArguments();
  19. // Use longer keepalive intervals to avoid "too_many_pings" errors from server
  20. channelArgs.SetInt(GRPC_ARG_KEEPALIVE_TIME_MS, 60000); // 60 seconds
  21. channelArgs.SetInt(GRPC_ARG_KEEPALIVE_TIMEOUT_MS, 20000); // 20 seconds
  22. channelArgs.SetInt(GRPC_ARG_KEEPALIVE_PERMIT_WITHOUT_CALLS, 0); // Only ping when there are active calls
  23. // Set max message size to match server (100MB for file uploads)
  24. channelArgs.SetMaxReceiveMessageSize(100 * 1024 * 1024);
  25. channelArgs.SetMaxSendMessageSize(100 * 1024 * 1024);
  26. channel_ = grpc::CreateCustomChannel(
  27. config_.address,
  28. grpc::InsecureChannelCredentials(),
  29. channelArgs
  30. );
  31. stub_ = smartbotic::databasepb::DatabaseService::NewStub(channel_);
  32. connected_ = true;
  33. spdlog::info("Database client connected to {}", config_.address);
  34. return true;
  35. } catch (const std::exception& e) {
  36. spdlog::error("Database client connection failed: {}", e.what());
  37. return false;
  38. }
  39. }
  40. void disconnect() {
  41. connected_ = false;
  42. stub_.reset();
  43. channel_.reset();
  44. }
  45. bool isConnected() const {
  46. if (!connected_ || !channel_) {
  47. return false;
  48. }
  49. // Check if channel is in a usable state (not failed or shutdown)
  50. auto state = channel_->GetState(false);
  51. return state == GRPC_CHANNEL_READY ||
  52. state == GRPC_CHANNEL_IDLE ||
  53. state == GRPC_CHANNEL_CONNECTING;
  54. }
  55. // ===== Document Operations =====
  56. std::string insert(const std::string& collection, const nlohmann::json& data,
  57. const std::string& id, uint32_t ttlSeconds,
  58. const std::string& actor) {
  59. smartbotic::databasepb::InsertRequest request;
  60. request.set_collection(collection);
  61. request.set_data(data.dump());
  62. if (!id.empty()) {
  63. request.set_id(id);
  64. }
  65. if (ttlSeconds > 0) {
  66. request.set_ttl_seconds(ttlSeconds);
  67. }
  68. if (!actor.empty()) {
  69. request.set_actor(actor);
  70. }
  71. smartbotic::databasepb::InsertResponse response;
  72. grpc::ClientContext context;
  73. setDeadline(context);
  74. auto status = stub_->Insert(&context, request, &response);
  75. if (!status.ok()) {
  76. spdlog::error("Client::insert failed: {}", status.error_message());
  77. throw std::runtime_error(status.error_message());
  78. }
  79. return response.id();
  80. }
  81. std::optional<nlohmann::json> get(const std::string& collection, const std::string& id) {
  82. smartbotic::databasepb::GetRequest request;
  83. request.set_collection(collection);
  84. request.set_id(id);
  85. smartbotic::databasepb::GetResponse response;
  86. grpc::ClientContext context;
  87. setDeadline(context);
  88. auto status = stub_->Get(&context, request, &response);
  89. if (!status.ok()) {
  90. spdlog::error("Client::get failed: {}", status.error_message());
  91. return std::nullopt;
  92. }
  93. if (!response.found()) {
  94. return std::nullopt;
  95. }
  96. auto json = nlohmann::json::parse(response.document().data());
  97. json["_id"] = response.document().id(); // Include document ID in returned data
  98. return json;
  99. }
  100. bool update(const std::string& collection, const std::string& id, const nlohmann::json& data,
  101. const std::string& actor) {
  102. smartbotic::databasepb::UpdateRequest request;
  103. request.set_collection(collection);
  104. request.set_id(id);
  105. request.set_data(data.dump());
  106. if (!actor.empty()) {
  107. request.set_actor(actor);
  108. }
  109. smartbotic::databasepb::UpdateResponse response;
  110. grpc::ClientContext context;
  111. setDeadline(context);
  112. auto status = stub_->Update(&context, request, &response);
  113. if (!status.ok()) {
  114. spdlog::error("Client::update failed: {}", status.error_message());
  115. return false;
  116. }
  117. return response.success();
  118. }
  119. bool updateIfVersion(const std::string& collection, const std::string& id,
  120. const nlohmann::json& data, uint64_t expectedVersion,
  121. const std::string& actor) {
  122. smartbotic::databasepb::UpdateRequest request;
  123. request.set_collection(collection);
  124. request.set_id(id);
  125. request.set_data(data.dump());
  126. request.set_expected_version(expectedVersion);
  127. if (!actor.empty()) {
  128. request.set_actor(actor);
  129. }
  130. smartbotic::databasepb::UpdateResponse response;
  131. grpc::ClientContext context;
  132. setDeadline(context);
  133. auto status = stub_->Update(&context, request, &response);
  134. if (!status.ok()) {
  135. spdlog::error("Client::updateIfVersion failed: {}", status.error_message());
  136. return false;
  137. }
  138. return response.success();
  139. }
  140. std::pair<std::string, bool> upsert(const std::string& collection, const nlohmann::json& data,
  141. const std::string& id, uint32_t ttlSeconds,
  142. const std::string& actor) {
  143. smartbotic::databasepb::UpsertRequest request;
  144. request.set_collection(collection);
  145. request.set_data(data.dump());
  146. if (!id.empty()) {
  147. request.set_id(id);
  148. }
  149. if (ttlSeconds > 0) {
  150. request.set_ttl_seconds(ttlSeconds);
  151. }
  152. if (!actor.empty()) {
  153. request.set_actor(actor);
  154. }
  155. smartbotic::databasepb::UpsertResponse response;
  156. grpc::ClientContext context;
  157. setDeadline(context);
  158. auto status = stub_->Upsert(&context, request, &response);
  159. if (!status.ok()) {
  160. spdlog::error("Client::upsert failed: {}", status.error_message());
  161. throw std::runtime_error(status.error_message());
  162. }
  163. return {response.id(), response.inserted()};
  164. }
  165. bool remove(const std::string& collection, const std::string& id) {
  166. smartbotic::databasepb::DeleteRequest request;
  167. request.set_collection(collection);
  168. request.set_id(id);
  169. smartbotic::databasepb::DeleteResponse response;
  170. grpc::ClientContext context;
  171. setDeadline(context);
  172. auto status = stub_->Delete(&context, request, &response);
  173. if (!status.ok()) {
  174. spdlog::error("Client::remove failed: {}", status.error_message());
  175. return false;
  176. }
  177. return response.deleted();
  178. }
  179. bool exists(const std::string& collection, const std::string& id) {
  180. smartbotic::databasepb::ExistsRequest request;
  181. request.set_collection(collection);
  182. request.set_id(id);
  183. smartbotic::databasepb::ExistsResponse response;
  184. grpc::ClientContext context;
  185. setDeadline(context);
  186. auto status = stub_->Exists(&context, request, &response);
  187. if (!status.ok()) {
  188. spdlog::error("Client::exists failed: {}", status.error_message());
  189. return false;
  190. }
  191. return response.exists();
  192. }
  193. // ===== Query Operations =====
  194. std::vector<nlohmann::json> find(const std::string& collection,
  195. const Client::QueryOptions& options) {
  196. smartbotic::databasepb::FindRequest request;
  197. request.set_collection(collection);
  198. // Set filters
  199. for (const auto& [field, value] : options.filters) {
  200. auto* filter = request.add_filters();
  201. filter->set_field(field);
  202. filter->set_value(value.dump());
  203. // Use SEARCH op for _search field, EQ for others
  204. if (field == "_search") {
  205. filter->set_op(smartbotic::databasepb::FILTER_OP_SEARCH);
  206. } else {
  207. filter->set_op(smartbotic::databasepb::FILTER_OP_EQ);
  208. }
  209. }
  210. // Set sorting
  211. if (!options.sortField.empty()) {
  212. auto* sort = request.mutable_sort();
  213. sort->set_field(options.sortField);
  214. sort->set_descending(options.sortDescending);
  215. }
  216. // Set pagination
  217. request.set_limit(options.limit);
  218. request.set_offset(options.offset);
  219. smartbotic::databasepb::FindResponse response;
  220. grpc::ClientContext context;
  221. setDeadline(context);
  222. auto status = stub_->Find(&context, request, &response);
  223. if (!status.ok()) {
  224. spdlog::error("Client::find failed: {}", status.error_message());
  225. return {};
  226. }
  227. std::vector<nlohmann::json> results;
  228. results.reserve(response.documents_size());
  229. for (const auto& doc : response.documents()) {
  230. auto json = nlohmann::json::parse(doc.data());
  231. json["_id"] = doc.id(); // Include document ID in returned data
  232. json["_created_at"] = doc.created_at(); // Creation timestamp (ms)
  233. json["_updated_at"] = doc.updated_at(); // Last update timestamp (ms)
  234. json["_created_by"] = doc.created_by(); // User ID who created
  235. json["_updated_by"] = doc.updated_by(); // User ID who last updated
  236. results.push_back(json);
  237. }
  238. return results;
  239. }
  240. uint64_t count(const std::string& collection,
  241. const std::vector<std::pair<std::string, nlohmann::json>>& filters) {
  242. smartbotic::databasepb::CountRequest request;
  243. request.set_collection(collection);
  244. for (const auto& [field, value] : filters) {
  245. auto* filter = request.add_filters();
  246. filter->set_field(field);
  247. filter->set_value(value.dump());
  248. // Use SEARCH op for _search field, EQ for others
  249. if (field == "_search") {
  250. filter->set_op(smartbotic::databasepb::FILTER_OP_SEARCH);
  251. } else {
  252. filter->set_op(smartbotic::databasepb::FILTER_OP_EQ);
  253. }
  254. }
  255. smartbotic::databasepb::CountResponse response;
  256. grpc::ClientContext context;
  257. setDeadline(context);
  258. auto status = stub_->Count(&context, request, &response);
  259. if (!status.ok()) {
  260. spdlog::error("Client::count failed: {}", status.error_message());
  261. return 0;
  262. }
  263. return response.count();
  264. }
  265. // ===== Set Operations =====
  266. bool setAdd(const std::string& collection, const std::string& setId, const std::string& member) {
  267. smartbotic::databasepb::SetAddRequest request;
  268. request.set_collection(collection);
  269. request.set_set_id(setId);
  270. request.set_member(member);
  271. smartbotic::databasepb::SetAddResponse response;
  272. grpc::ClientContext context;
  273. setDeadline(context);
  274. auto status = stub_->SetAdd(&context, request, &response);
  275. if (!status.ok()) {
  276. spdlog::error("Client::setAdd failed: {}", status.error_message());
  277. return false;
  278. }
  279. return response.added();
  280. }
  281. bool setRemove(const std::string& collection, const std::string& setId, const std::string& member) {
  282. smartbotic::databasepb::SetRemoveRequest request;
  283. request.set_collection(collection);
  284. request.set_set_id(setId);
  285. request.set_member(member);
  286. smartbotic::databasepb::SetRemoveResponse response;
  287. grpc::ClientContext context;
  288. setDeadline(context);
  289. auto status = stub_->SetRemove(&context, request, &response);
  290. if (!status.ok()) {
  291. spdlog::error("Client::setRemove failed: {}", status.error_message());
  292. return false;
  293. }
  294. return response.removed();
  295. }
  296. std::vector<std::string> setMembers(const std::string& collection, const std::string& setId) {
  297. smartbotic::databasepb::SetMembersRequest request;
  298. request.set_collection(collection);
  299. request.set_set_id(setId);
  300. smartbotic::databasepb::SetMembersResponse response;
  301. grpc::ClientContext context;
  302. setDeadline(context);
  303. auto status = stub_->SetMembers(&context, request, &response);
  304. if (!status.ok()) {
  305. spdlog::error("Client::setMembers failed: {}", status.error_message());
  306. return {};
  307. }
  308. return {response.members().begin(), response.members().end()};
  309. }
  310. bool setIsMember(const std::string& collection, const std::string& setId, const std::string& member) {
  311. smartbotic::databasepb::SetIsMemberRequest request;
  312. request.set_collection(collection);
  313. request.set_set_id(setId);
  314. request.set_member(member);
  315. smartbotic::databasepb::SetIsMemberResponse response;
  316. grpc::ClientContext context;
  317. setDeadline(context);
  318. auto status = stub_->SetIsMember(&context, request, &response);
  319. if (!status.ok()) {
  320. spdlog::error("Client::setIsMember failed: {}", status.error_message());
  321. return false;
  322. }
  323. return response.is_member();
  324. }
  325. // ===== Collection Management =====
  326. bool createCollection(const std::string& name, uint32_t defaultTtlSeconds, bool encrypted,
  327. uint32_t maxVersions) {
  328. smartbotic::databasepb::CreateCollectionRequest request;
  329. request.set_name(name);
  330. auto* options = request.mutable_options();
  331. if (defaultTtlSeconds > 0) {
  332. options->set_default_ttl_seconds(defaultTtlSeconds);
  333. }
  334. options->set_encrypted(encrypted);
  335. if (maxVersions > 0) {
  336. options->set_max_versions(maxVersions);
  337. }
  338. smartbotic::databasepb::CreateCollectionResponse response;
  339. grpc::ClientContext context;
  340. setDeadline(context);
  341. auto status = stub_->CreateCollection(&context, request, &response);
  342. if (!status.ok()) {
  343. spdlog::error("Client::createCollection failed: {}", status.error_message());
  344. return false;
  345. }
  346. return response.created();
  347. }
  348. bool dropCollection(const std::string& name) {
  349. smartbotic::databasepb::DropCollectionRequest request;
  350. request.set_name(name);
  351. smartbotic::databasepb::DropCollectionResponse response;
  352. grpc::ClientContext context;
  353. setDeadline(context);
  354. auto status = stub_->DropCollection(&context, request, &response);
  355. if (!status.ok()) {
  356. spdlog::error("Client::dropCollection failed: {}", status.error_message());
  357. return false;
  358. }
  359. return response.dropped();
  360. }
  361. std::vector<std::string> listCollections() {
  362. smartbotic::databasepb::ListCollectionsRequest request;
  363. smartbotic::databasepb::ListCollectionsResponse response;
  364. grpc::ClientContext context;
  365. setDeadline(context);
  366. auto status = stub_->ListCollections(&context, request, &response);
  367. if (!status.ok()) {
  368. spdlog::error("Client::listCollections failed: {}", status.error_message());
  369. return {};
  370. }
  371. return {response.names().begin(), response.names().end()};
  372. }
  373. std::optional<Client::CollectionInfo> getCollectionInfo(const std::string& name) {
  374. smartbotic::databasepb::GetCollectionInfoRequest request;
  375. request.set_name(name);
  376. smartbotic::databasepb::GetCollectionInfoResponse response;
  377. grpc::ClientContext context;
  378. setDeadline(context);
  379. auto status = stub_->GetCollectionInfo(&context, request, &response);
  380. if (!status.ok()) {
  381. spdlog::error("Client::getCollectionInfo failed: {}", status.error_message());
  382. return std::nullopt;
  383. }
  384. if (!response.found()) {
  385. return std::nullopt;
  386. }
  387. const auto& info = response.info();
  388. Client::CollectionInfo result;
  389. result.name = info.name();
  390. result.documentCount = info.document_count();
  391. result.sizeBytes = info.size_bytes();
  392. result.defaultTtlSeconds = info.options().default_ttl_seconds();
  393. result.encrypted = info.options().encrypted();
  394. result.maxVersions = info.options().max_versions();
  395. result.createdAt = info.created_at();
  396. result.updatedAt = info.updated_at();
  397. return result;
  398. }
  399. // ===== Event Subscription =====
  400. class SubscriptionHandle {
  401. public:
  402. SubscriptionHandle(std::shared_ptr<grpc::ClientContext> ctx,
  403. std::unique_ptr<grpc::ClientReader<smartbotic::databasepb::DatabaseEvent>> reader,
  404. std::thread readerThread)
  405. : context_(std::move(ctx))
  406. , reader_(std::move(reader))
  407. , readerThread_(std::move(readerThread))
  408. , active_(true) {}
  409. ~SubscriptionHandle() {
  410. cancel();
  411. }
  412. void cancel() {
  413. if (active_.exchange(false)) {
  414. context_->TryCancel();
  415. if (readerThread_.joinable()) {
  416. readerThread_.join();
  417. }
  418. }
  419. }
  420. private:
  421. std::shared_ptr<grpc::ClientContext> context_;
  422. std::unique_ptr<grpc::ClientReader<smartbotic::databasepb::DatabaseEvent>> reader_;
  423. std::thread readerThread_;
  424. std::atomic<bool> active_;
  425. };
  426. std::shared_ptr<void> subscribe(const std::vector<std::string>& collections,
  427. Client::EventCallback callback) {
  428. auto context = std::make_shared<grpc::ClientContext>();
  429. smartbotic::databasepb::SubscribeRequest request;
  430. for (const auto& coll : collections) {
  431. request.add_collections(coll);
  432. }
  433. request.set_include_data(true);
  434. auto reader = stub_->Subscribe(context.get(), request);
  435. // Create reader thread
  436. auto readerThread = std::thread([reader = reader.get(), callback = std::move(callback)]() {
  437. smartbotic::databasepb::DatabaseEvent event;
  438. while (reader->Read(&event)) {
  439. std::optional<nlohmann::json> data;
  440. if (!event.data().empty()) {
  441. data = nlohmann::json::parse(event.data());
  442. }
  443. // Convert proto event type to string
  444. std::string eventType;
  445. switch (event.type()) {
  446. case smartbotic::databasepb::EVENT_INSERT: eventType = "insert"; break;
  447. case smartbotic::databasepb::EVENT_UPDATE: eventType = "update"; break;
  448. case smartbotic::databasepb::EVENT_DELETE: eventType = "delete"; break;
  449. case smartbotic::databasepb::EVENT_EXPIRE: eventType = "expire"; break;
  450. case smartbotic::databasepb::EVENT_INVALIDATE: eventType = "invalidate"; break;
  451. default: eventType = "unknown"; break;
  452. }
  453. callback(
  454. event.collection(),
  455. event.document_id(),
  456. eventType,
  457. data
  458. );
  459. }
  460. });
  461. return std::make_shared<SubscriptionHandle>(
  462. std::move(context),
  463. std::move(reader),
  464. std::move(readerThread)
  465. );
  466. }
  467. // ===== Version History =====
  468. Client::VersionHistoryResult getVersionHistory(const std::string& collection,
  469. const std::string& id, uint32_t limit, uint32_t offset) {
  470. smartbotic::databasepb::GetVersionHistoryRequest request;
  471. request.set_collection(collection);
  472. request.set_id(id);
  473. if (limit > 0) request.set_limit(limit);
  474. if (offset > 0) request.set_offset(offset);
  475. smartbotic::databasepb::GetVersionHistoryResponse response;
  476. grpc::ClientContext context;
  477. setDeadline(context);
  478. auto status = stub_->GetVersionHistory(&context, request, &response);
  479. if (!status.ok()) {
  480. spdlog::error("Client::getVersionHistory failed: {}", status.error_message());
  481. return {};
  482. }
  483. Client::VersionHistoryResult result;
  484. result.currentVersion = response.current_version();
  485. result.totalCount = response.total_count();
  486. result.documentDeleted = response.document_deleted();
  487. result.versions.reserve(response.versions_size());
  488. for (const auto& ver : response.versions()) {
  489. Client::VersionEntry entry;
  490. entry.version = ver.version();
  491. entry.timestamp = ver.timestamp();
  492. entry.updatedBy = ver.updated_by();
  493. if (!ver.data().empty()) {
  494. entry.data = nlohmann::json::parse(ver.data());
  495. }
  496. result.versions.push_back(std::move(entry));
  497. }
  498. return result;
  499. }
  500. std::optional<Client::VersionEntry> getDocumentVersion(const std::string& collection,
  501. const std::string& id, uint64_t version) {
  502. smartbotic::databasepb::GetDocumentVersionRequest request;
  503. request.set_collection(collection);
  504. request.set_id(id);
  505. request.set_version(version);
  506. smartbotic::databasepb::GetDocumentVersionResponse response;
  507. grpc::ClientContext context;
  508. setDeadline(context);
  509. auto status = stub_->GetDocumentVersion(&context, request, &response);
  510. if (!status.ok()) {
  511. spdlog::error("Client::getDocumentVersion failed: {}", status.error_message());
  512. return std::nullopt;
  513. }
  514. if (!response.found()) {
  515. return std::nullopt;
  516. }
  517. Client::VersionEntry entry;
  518. entry.version = response.version_entry().version();
  519. entry.timestamp = response.version_entry().timestamp();
  520. entry.updatedBy = response.version_entry().updated_by();
  521. if (!response.version_entry().data().empty()) {
  522. entry.data = nlohmann::json::parse(response.version_entry().data());
  523. }
  524. return entry;
  525. }
  526. uint64_t restoreVersion(const std::string& collection, const std::string& id,
  527. uint64_t version, const std::string& actor) {
  528. smartbotic::databasepb::RestoreVersionRequest request;
  529. request.set_collection(collection);
  530. request.set_id(id);
  531. request.set_version(version);
  532. if (!actor.empty()) request.set_actor(actor);
  533. smartbotic::databasepb::RestoreVersionResponse response;
  534. grpc::ClientContext context;
  535. setDeadline(context);
  536. auto status = stub_->RestoreVersion(&context, request, &response);
  537. if (!status.ok()) {
  538. spdlog::error("Client::restoreVersion failed: {}", status.error_message());
  539. return 0;
  540. }
  541. if (!response.success()) {
  542. spdlog::error("Client::restoreVersion: {}", response.error());
  543. return 0;
  544. }
  545. return response.new_version();
  546. }
  547. std::pair<uint64_t, uint64_t> restoreToDate(const std::string& collection,
  548. const std::string& id, uint64_t timestamp, const std::string& actor) {
  549. smartbotic::databasepb::RestoreToDateRequest request;
  550. request.set_collection(collection);
  551. request.set_id(id);
  552. request.set_timestamp(timestamp);
  553. if (!actor.empty()) request.set_actor(actor);
  554. smartbotic::databasepb::RestoreToDateResponse response;
  555. grpc::ClientContext context;
  556. setDeadline(context);
  557. auto status = stub_->RestoreToDate(&context, request, &response);
  558. if (!status.ok()) {
  559. spdlog::error("Client::restoreToDate failed: {}", status.error_message());
  560. return {0, 0};
  561. }
  562. if (!response.success()) {
  563. spdlog::error("Client::restoreToDate: {}", response.error());
  564. return {0, 0};
  565. }
  566. return {response.restored_version(), response.new_version()};
  567. }
  568. // ===== Health =====
  569. bool healthCheck() {
  570. smartbotic::databasepb::HealthCheckRequest request;
  571. smartbotic::databasepb::HealthCheckResponse response;
  572. grpc::ClientContext context;
  573. setDeadline(context);
  574. auto status = stub_->HealthCheck(&context, request, &response);
  575. if (!status.ok()) {
  576. return false;
  577. }
  578. return response.healthy();
  579. }
  580. std::optional<Client::HealthInfo> getHealthInfo() {
  581. smartbotic::databasepb::HealthCheckRequest request;
  582. smartbotic::databasepb::HealthCheckResponse response;
  583. grpc::ClientContext context;
  584. setDeadline(context);
  585. auto status = stub_->HealthCheck(&context, request, &response);
  586. if (!status.ok()) {
  587. return std::nullopt;
  588. }
  589. Client::HealthInfo info;
  590. info.healthy = response.healthy();
  591. info.uptimeMs = response.uptime_seconds() * 1000;
  592. info.documentCount = response.document_count();
  593. info.memoryUsedBytes = response.memory_used_bytes();
  594. info.walSizeBytes = response.wal_size_bytes();
  595. return info;
  596. }
  597. private:
  598. void setDeadline(grpc::ClientContext& context) {
  599. context.set_deadline(
  600. std::chrono::system_clock::now() + std::chrono::milliseconds(config_.timeoutMs)
  601. );
  602. }
  603. Config config_;
  604. std::shared_ptr<grpc::Channel> channel_;
  605. std::unique_ptr<smartbotic::databasepb::DatabaseService::Stub> stub_;
  606. std::atomic<bool> connected_{false};
  607. };
  608. // ===== Client Public Interface Implementation =====
  609. Client::Client(Config config)
  610. : impl_(std::make_unique<Impl>(std::move(config))) {}
  611. Client::~Client() = default;
  612. Client::Client(Client&&) noexcept = default;
  613. Client& Client::operator=(Client&&) noexcept = default;
  614. bool Client::connect() {
  615. return impl_->connect();
  616. }
  617. bool Client::isConnected() const {
  618. return impl_->isConnected();
  619. }
  620. std::string Client::insert(const std::string& collection, const nlohmann::json& data,
  621. const std::string& id, uint32_t ttlSeconds,
  622. const std::string& actor) {
  623. return impl_->insert(collection, data, id, ttlSeconds, actor);
  624. }
  625. std::optional<nlohmann::json> Client::get(const std::string& collection, const std::string& id) {
  626. return impl_->get(collection, id);
  627. }
  628. bool Client::update(const std::string& collection, const std::string& id, const nlohmann::json& data,
  629. const std::string& actor) {
  630. return impl_->update(collection, id, data, actor);
  631. }
  632. bool Client::updateIfVersion(const std::string& collection, const std::string& id,
  633. const nlohmann::json& data, uint64_t expectedVersion,
  634. const std::string& actor) {
  635. return impl_->updateIfVersion(collection, id, data, expectedVersion, actor);
  636. }
  637. std::pair<std::string, bool> Client::upsert(const std::string& collection, const nlohmann::json& data,
  638. const std::string& id, uint32_t ttlSeconds,
  639. const std::string& actor) {
  640. return impl_->upsert(collection, data, id, ttlSeconds, actor);
  641. }
  642. bool Client::remove(const std::string& collection, const std::string& id) {
  643. return impl_->remove(collection, id);
  644. }
  645. bool Client::exists(const std::string& collection, const std::string& id) {
  646. return impl_->exists(collection, id);
  647. }
  648. Client::VersionHistoryResult Client::getVersionHistory(const std::string& collection,
  649. const std::string& id, uint32_t limit, uint32_t offset) {
  650. return impl_->getVersionHistory(collection, id, limit, offset);
  651. }
  652. std::optional<Client::VersionEntry> Client::getDocumentVersion(const std::string& collection,
  653. const std::string& id, uint64_t version) {
  654. return impl_->getDocumentVersion(collection, id, version);
  655. }
  656. uint64_t Client::restoreVersion(const std::string& collection, const std::string& id,
  657. uint64_t version, const std::string& actor) {
  658. return impl_->restoreVersion(collection, id, version, actor);
  659. }
  660. std::pair<uint64_t, uint64_t> Client::restoreToDate(const std::string& collection,
  661. const std::string& id, uint64_t timestamp, const std::string& actor) {
  662. return impl_->restoreToDate(collection, id, timestamp, actor);
  663. }
  664. std::vector<nlohmann::json> Client::find(const std::string& collection,
  665. const QueryOptions& options) {
  666. return impl_->find(collection, options);
  667. }
  668. std::vector<nlohmann::json> Client::find(const std::string& collection) {
  669. return impl_->find(collection, QueryOptions{});
  670. }
  671. uint64_t Client::count(const std::string& collection,
  672. const std::vector<std::pair<std::string, nlohmann::json>>& filters) {
  673. return impl_->count(collection, filters);
  674. }
  675. uint64_t Client::count(const std::string& collection) {
  676. return impl_->count(collection, {});
  677. }
  678. bool Client::setAdd(const std::string& collection, const std::string& setId, const std::string& member) {
  679. return impl_->setAdd(collection, setId, member);
  680. }
  681. bool Client::setRemove(const std::string& collection, const std::string& setId, const std::string& member) {
  682. return impl_->setRemove(collection, setId, member);
  683. }
  684. std::vector<std::string> Client::setMembers(const std::string& collection, const std::string& setId) {
  685. return impl_->setMembers(collection, setId);
  686. }
  687. bool Client::setIsMember(const std::string& collection, const std::string& setId, const std::string& member) {
  688. return impl_->setIsMember(collection, setId, member);
  689. }
  690. bool Client::createCollection(const std::string& name, uint32_t defaultTtlSeconds,
  691. bool encrypted, uint32_t maxVersions) {
  692. return impl_->createCollection(name, defaultTtlSeconds, encrypted, maxVersions);
  693. }
  694. bool Client::dropCollection(const std::string& name) {
  695. return impl_->dropCollection(name);
  696. }
  697. std::vector<std::string> Client::listCollections() {
  698. return impl_->listCollections();
  699. }
  700. std::optional<Client::CollectionInfo> Client::getCollectionInfo(const std::string& name) {
  701. return impl_->getCollectionInfo(name);
  702. }
  703. std::shared_ptr<void> Client::subscribe(const std::vector<std::string>& collections, EventCallback callback) {
  704. return impl_->subscribe(collections, std::move(callback));
  705. }
  706. bool Client::healthCheck() {
  707. return impl_->healthCheck();
  708. }
  709. std::optional<Client::HealthInfo> Client::getHealthInfo() {
  710. return impl_->getHealthInfo();
  711. }
  712. } // namespace smartbotic::database