client.cpp 84 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117
  1. #include "smartbotic/database/client.hpp"
  2. #include <database.grpc.pb.h>
  3. #include <grpcpp/grpcpp.h>
  4. #include <grpcpp/security/credentials.h>
  5. #include <grpcpp/security/tls_certificate_verifier.h>
  6. #include <grpcpp/security/tls_credentials_options.h>
  7. #include <spdlog/spdlog.h>
  8. #include <atomic>
  9. #include <chrono>
  10. #include <cstdlib>
  11. #include <fstream>
  12. #include <mutex>
  13. #include <sstream>
  14. #include <thread>
  15. namespace smartbotic::database {
  16. namespace {
  17. // v1.7.0 T11 — transient transport errors that are safe to retry for idempotent
  18. // writes (insert with explicit ID, updateIfVersion, upsert, remove, patch).
  19. // All other codes (INVALID_ARGUMENT, FAILED_PRECONDITION, NOT_FOUND, ...) are
  20. // final — no retry.
  21. bool retryableStatus(const grpc::Status& s) {
  22. switch (s.error_code()) {
  23. case grpc::StatusCode::DEADLINE_EXCEEDED: // server queued behind eviction/WAL
  24. case grpc::StatusCode::RESOURCE_EXHAUSTED: // admission control or gRPC concurrency cap
  25. case grpc::StatusCode::UNAVAILABLE: // transient connection issue
  26. return true;
  27. default:
  28. return false;
  29. }
  30. }
  31. // Exponential backoff with jitter. attempt=0 => baseMs; attempt=1 => baseMs*2; ...
  32. uint32_t computeBackoffMs(uint32_t baseMs, uint32_t attempt, uint32_t capMs, double jitterPct) {
  33. uint64_t exp = static_cast<uint64_t>(baseMs) << attempt;
  34. if (exp > capMs) exp = capMs;
  35. // ±jitterPct jitter; std::rand() is fine here — not a security sensitive RNG.
  36. double jitter = 1.0 + ((double(std::rand()) / RAND_MAX) * 2.0 - 1.0) * jitterPct;
  37. if (jitter < 0.1) jitter = 0.1;
  38. return static_cast<uint32_t>(exp * jitter);
  39. }
  40. // v2.4 — slurp a PEM-encoded CA cert file into a string. Returns empty
  41. // string on read failure; caller logs the surrounding context.
  42. std::string readCaCertFile(const std::string& path) {
  43. std::ifstream in(path, std::ios::binary);
  44. if (!in.is_open()) return {};
  45. std::ostringstream buf;
  46. buf << in.rdbuf();
  47. return buf.str();
  48. }
  49. } // anonymous namespace
  50. // ===== PIMPL Implementation =====
  51. class Client::Impl {
  52. public:
  53. explicit Impl(Config config) : config_(std::move(config)) {
  54. if (config_.project.empty()) config_.project = "default";
  55. }
  56. // v2.3 — qualify a collection name with Config::project unless the
  57. // caller already supplied a "<project>:<collection>" qualified form.
  58. // The qualified form always wins. Server-side parseProjectCollection
  59. // is strict on edge cases (leading colon, double colon, etc.) so
  60. // we just delegate the validation to the server.
  61. std::string qualify(std::string_view collection) const {
  62. if (collection.find(':') != std::string_view::npos) {
  63. return std::string(collection);
  64. }
  65. // v2.7.0 — system collections (`_views`, `_collection_meta`,
  66. // `_policies`, `_files`) are GLOBAL, not per-project: their rows are
  67. // already keyed by a project-qualified id. Qualifying the collection
  68. // name too produced `default:_policies`, which does not exist, so
  69. // reads silently found nothing while writes succeeded (the server
  70. // tolerates both spellings on the policy write path). That asymmetry
  71. // is exactly how a management command can look like it worked and
  72. // then show stale state.
  73. if (!collection.empty() && collection.front() == '_') {
  74. return std::string(collection);
  75. }
  76. std::string out;
  77. out.reserve(config_.project.size() + 1 + collection.size());
  78. out.append(config_.project);
  79. out.push_back(':');
  80. out.append(collection.data(), collection.size());
  81. return out;
  82. }
  83. // Inverse of qualify() for values coming back off the wire. Strips only
  84. // OUR project's prefix, so a name from another namespace stays visibly
  85. // qualified rather than being silently flattened into ours.
  86. std::string unqualify(const std::string& qualified) const {
  87. const std::string prefix = config_.project + ":";
  88. if (qualified.rfind(prefix, 0) == 0) {
  89. return qualified.substr(prefix.size());
  90. }
  91. return qualified;
  92. }
  93. ~Impl() {
  94. disconnect();
  95. }
  96. bool connect() {
  97. try {
  98. auto channelArgs = grpc::ChannelArguments();
  99. // Use longer keepalive intervals to avoid "too_many_pings" errors from server
  100. channelArgs.SetInt(GRPC_ARG_KEEPALIVE_TIME_MS, 60000); // 60 seconds
  101. channelArgs.SetInt(GRPC_ARG_KEEPALIVE_TIMEOUT_MS, 20000); // 20 seconds
  102. channelArgs.SetInt(GRPC_ARG_KEEPALIVE_PERMIT_WITHOUT_CALLS, 0); // Only ping when there are active calls
  103. // Set max message size to match server (100MB for file uploads)
  104. channelArgs.SetMaxReceiveMessageSize(100 * 1024 * 1024);
  105. channelArgs.SetMaxSendMessageSize(100 * 1024 * 1024);
  106. std::shared_ptr<grpc::ChannelCredentials> creds = buildChannelCredentials();
  107. if (!creds) {
  108. spdlog::error("Database client connection failed: could not build channel credentials");
  109. return false;
  110. }
  111. channel_ = grpc::CreateCustomChannel(
  112. config_.address,
  113. creds,
  114. channelArgs
  115. );
  116. stub_ = smartbotic::databasepb::DatabaseService::NewStub(channel_);
  117. connected_ = true;
  118. const char* scheme = config_.tls_enabled ? "tls" : "plaintext";
  119. const char* authed = config_.auth_token.empty() ? "no-auth" : "bearer-auth";
  120. spdlog::info("Database client connected to {} ({}, {})",
  121. config_.address, scheme, authed);
  122. return true;
  123. } catch (const std::exception& e) {
  124. spdlog::error("Database client connection failed: {}", e.what());
  125. return false;
  126. }
  127. }
  128. // v2.4 — build the ChannelCredentials matching Config's TLS settings.
  129. //
  130. // Three cases:
  131. // 1. !tls_enabled → InsecureChannelCredentials
  132. // (v2.3 back-compat default).
  133. // 2. tls_enabled && !insecure_skip_verify → SslCredentials with
  134. // pem_root_certs read from tls_ca_cert_path, or empty (= system
  135. // trust roots) when the path is unset.
  136. // 3. tls_enabled && insecure_skip_verify → experimental
  137. // TlsChannelCredentialsOptions with NoOpCertificateVerifier and
  138. // set_verify_server_certs(false). Loud WARN on every connect.
  139. std::shared_ptr<grpc::ChannelCredentials> buildChannelCredentials() {
  140. if (!config_.tls_enabled) {
  141. return grpc::InsecureChannelCredentials();
  142. }
  143. if (config_.tls_insecure_skip_verify) {
  144. spdlog::warn(
  145. "Database client: connecting with tls_insecure_skip_verify — "
  146. "DEV ONLY. Server certificate will NOT be validated against any "
  147. "trust root. Drop the server's self-signed cert at "
  148. "tls_ca_cert_path for any production-shaped use.");
  149. // Route through gRPC's experimental TLS API. This is the path
  150. // gRPC documents for "accept any cert" — combine
  151. // NoOpCertificateVerifier with set_verify_server_certs(false).
  152. grpc::experimental::TlsChannelCredentialsOptions tlsOpts;
  153. tlsOpts.set_certificate_verifier(
  154. std::make_shared<grpc::experimental::NoOpCertificateVerifier>());
  155. tlsOpts.set_verify_server_certs(false);
  156. tlsOpts.set_check_call_host(false);
  157. auto creds = grpc::experimental::TlsCredentials(tlsOpts);
  158. if (!creds) {
  159. // No fallback that "trusts everything" without
  160. // experimental TLS — if the build/runtime can't produce
  161. // experimental TlsCredentials, refuse to silently degrade
  162. // to system-trust mode (which would reject the
  163. // self-signed cert anyway). The operator should bake the
  164. // server's self-signed cert into tls_ca_cert_path.
  165. spdlog::error(
  166. "Database client: tls_insecure_skip_verify requested but "
  167. "experimental TlsCredentials returned null. Set "
  168. "tls_ca_cert_path to the server's self-signed PEM "
  169. "instead.");
  170. return nullptr;
  171. }
  172. return creds;
  173. }
  174. grpc::SslCredentialsOptions sslOpts;
  175. if (!config_.tls_ca_cert_path.empty()) {
  176. sslOpts.pem_root_certs = readCaCertFile(config_.tls_ca_cert_path);
  177. if (sslOpts.pem_root_certs.empty()) {
  178. spdlog::error(
  179. "Database client: could not read tls_ca_cert_path '{}'",
  180. config_.tls_ca_cert_path);
  181. return nullptr;
  182. }
  183. }
  184. // Else leave pem_root_certs empty → gRPC uses system trust roots
  185. // (or GRPC_DEFAULT_SSL_ROOTS_FILE_PATH env var, per gRPC docs).
  186. return grpc::SslCredentials(sslOpts);
  187. }
  188. void disconnect() {
  189. connected_ = false;
  190. stub_.reset();
  191. channel_.reset();
  192. }
  193. bool isConnected() const {
  194. if (!connected_ || !channel_) {
  195. return false;
  196. }
  197. // Check if channel is in a usable state (not failed or shutdown)
  198. auto state = channel_->GetState(false);
  199. return state == GRPC_CHANNEL_READY ||
  200. state == GRPC_CHANNEL_IDLE ||
  201. state == GRPC_CHANNEL_CONNECTING;
  202. }
  203. // ===== Document Operations =====
  204. std::string insert(const std::string& collection, const nlohmann::json& data,
  205. const std::string& id, uint32_t ttlSeconds,
  206. const std::string& actor) {
  207. smartbotic::databasepb::InsertRequest request;
  208. request.set_collection(qualify(collection));
  209. request.set_data(data.dump());
  210. if (!id.empty()) {
  211. request.set_id(id);
  212. }
  213. if (ttlSeconds > 0) {
  214. request.set_ttl_seconds(ttlSeconds);
  215. }
  216. if (!actor.empty()) {
  217. request.set_actor(actor);
  218. }
  219. smartbotic::databasepb::InsertResponse response;
  220. grpc::Status status;
  221. for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
  222. grpc::ClientContext context;
  223. setDeadline(context);
  224. response.Clear();
  225. status = stub_->Insert(&context, request, &response);
  226. if (status.ok() || !retryableStatus(status)) {
  227. break;
  228. }
  229. if (attempt < config_.writeRetries) {
  230. uint32_t backoffMs = computeBackoffMs(
  231. config_.writeRetryBackoffMs, attempt,
  232. config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
  233. spdlog::warn("Client::insert {}; retrying in {}ms (attempt {}/{})",
  234. status.error_message(), backoffMs,
  235. attempt + 1, config_.writeRetries);
  236. std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
  237. }
  238. }
  239. if (!status.ok()) {
  240. spdlog::error("Client::insert failed after retries: {}", status.error_message());
  241. throw std::runtime_error(status.error_message());
  242. }
  243. return response.id();
  244. }
  245. std::optional<nlohmann::json> get(const std::string& collection, const std::string& id) {
  246. smartbotic::databasepb::GetRequest request;
  247. request.set_collection(qualify(collection));
  248. request.set_id(id);
  249. smartbotic::databasepb::GetResponse response;
  250. grpc::ClientContext context;
  251. setDeadline(context);
  252. auto status = stub_->Get(&context, request, &response);
  253. if (!status.ok()) {
  254. // v2.8.0 — a permission denial must NOT look like "no data".
  255. // Returning empty here would leave an application unable to tell
  256. // "you may not see this" from "there is nothing to see", so it
  257. // would silently take the wrong branch. Same reasoning as the
  258. // v2.4.2 change that stopped Find returning empty for a missing
  259. // collection.
  260. if (status.error_code() == grpc::StatusCode::PERMISSION_DENIED) {
  261. throw std::runtime_error("access denied: get");
  262. }
  263. spdlog::error("Client::get failed: {}", status.error_message());
  264. return std::nullopt;
  265. }
  266. if (!response.found()) {
  267. return std::nullopt;
  268. }
  269. auto json = nlohmann::json::parse(response.document().data());
  270. json["_id"] = response.document().id();
  271. json["_version"] = response.document().version();
  272. json["_created_at"] = response.document().created_at();
  273. json["_updated_at"] = response.document().updated_at();
  274. json["_created_by"] = response.document().created_by();
  275. json["_updated_by"] = response.document().updated_by();
  276. return json;
  277. }
  278. bool update(const std::string& collection, const std::string& id, const nlohmann::json& data,
  279. const std::string& actor) {
  280. // Strip metadata fields from data — callers pass the JSON from get() which
  281. // includes _version, _id, etc. We use _version for optimistic locking.
  282. auto cleanData = data;
  283. cleanData.erase("_id");
  284. cleanData.erase("_version");
  285. cleanData.erase("_created_at");
  286. cleanData.erase("_updated_at");
  287. cleanData.erase("_created_by");
  288. cleanData.erase("_updated_by");
  289. // Optimistic locking with automatic retry:
  290. // 1. Read current version
  291. // 2. Call updateIfVersion
  292. // 3. On version conflict (another writer), re-read and retry
  293. for (uint32_t attempt = 0; attempt <= config_.maxRetries; ++attempt) {
  294. // Get current version
  295. auto current = get(collection, id);
  296. if (!current) {
  297. return false; // document doesn't exist
  298. }
  299. uint64_t version = (*current)["_version"].get<uint64_t>();
  300. // Attempt version-checked update
  301. if (updateIfVersion(collection, id, cleanData, version, actor)) {
  302. return true;
  303. }
  304. // Version conflict — another client wrote between our get() and update
  305. if (attempt < config_.maxRetries) {
  306. spdlog::debug("Client::update version conflict on {}/{}, retry {}/{}",
  307. collection, id, attempt + 1, config_.maxRetries);
  308. }
  309. }
  310. spdlog::warn("Client::update failed after {} retries due to version conflicts on {}/{}",
  311. config_.maxRetries, collection, id);
  312. return false;
  313. }
  314. bool updateIfVersion(const std::string& collection, const std::string& id,
  315. const nlohmann::json& data, uint64_t expectedVersion,
  316. const std::string& actor) {
  317. smartbotic::databasepb::UpdateRequest request;
  318. request.set_collection(qualify(collection));
  319. request.set_id(id);
  320. request.set_data(data.dump());
  321. request.set_expected_version(expectedVersion);
  322. if (!actor.empty()) {
  323. request.set_actor(actor);
  324. }
  325. smartbotic::databasepb::UpdateResponse response;
  326. grpc::Status status;
  327. for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
  328. grpc::ClientContext context;
  329. setDeadline(context);
  330. response.Clear();
  331. status = stub_->Update(&context, request, &response);
  332. if (status.ok() || !retryableStatus(status)) {
  333. break;
  334. }
  335. if (attempt < config_.writeRetries) {
  336. uint32_t backoffMs = computeBackoffMs(
  337. config_.writeRetryBackoffMs, attempt,
  338. config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
  339. spdlog::warn("Client::updateIfVersion {}; retrying in {}ms (attempt {}/{})",
  340. status.error_message(), backoffMs,
  341. attempt + 1, config_.writeRetries);
  342. std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
  343. }
  344. }
  345. if (!status.ok()) {
  346. spdlog::error("Client::updateIfVersion failed after retries: {}", status.error_message());
  347. return false;
  348. }
  349. return response.success();
  350. }
  351. uint64_t patch(const std::string& collection, const std::string& id,
  352. const nlohmann::json& fields, const std::string& actor) {
  353. smartbotic::databasepb::PatchDocumentRequest request;
  354. request.set_collection(qualify(collection));
  355. request.set_id(id);
  356. request.set_patch_json(fields.dump());
  357. if (!actor.empty()) {
  358. request.set_actor(actor);
  359. }
  360. smartbotic::databasepb::PatchDocumentResponse response;
  361. grpc::Status status;
  362. for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
  363. grpc::ClientContext context;
  364. setDeadline(context);
  365. response.Clear();
  366. status = stub_->PatchDocument(&context, request, &response);
  367. if (status.ok() || !retryableStatus(status)) {
  368. break;
  369. }
  370. if (attempt < config_.writeRetries) {
  371. uint32_t backoffMs = computeBackoffMs(
  372. config_.writeRetryBackoffMs, attempt,
  373. config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
  374. spdlog::warn("Client::patch {}; retrying in {}ms (attempt {}/{})",
  375. status.error_message(), backoffMs,
  376. attempt + 1, config_.writeRetries);
  377. std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
  378. }
  379. }
  380. if (!status.ok()) {
  381. spdlog::error("Client::patch failed after retries: {}", status.error_message());
  382. return 0;
  383. }
  384. if (!response.success()) {
  385. spdlog::error("Client::patch failed: {}", response.error());
  386. return 0;
  387. }
  388. return response.new_version();
  389. }
  390. std::pair<std::string, bool> upsert(const std::string& collection, const nlohmann::json& data,
  391. const std::string& id, uint32_t ttlSeconds,
  392. const std::string& actor) {
  393. smartbotic::databasepb::UpsertRequest request;
  394. request.set_collection(qualify(collection));
  395. request.set_data(data.dump());
  396. if (!id.empty()) {
  397. request.set_id(id);
  398. }
  399. if (ttlSeconds > 0) {
  400. request.set_ttl_seconds(ttlSeconds);
  401. }
  402. if (!actor.empty()) {
  403. request.set_actor(actor);
  404. }
  405. smartbotic::databasepb::UpsertResponse response;
  406. grpc::Status status;
  407. for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
  408. grpc::ClientContext context;
  409. setDeadline(context);
  410. response.Clear();
  411. status = stub_->Upsert(&context, request, &response);
  412. if (status.ok() || !retryableStatus(status)) {
  413. break;
  414. }
  415. if (attempt < config_.writeRetries) {
  416. uint32_t backoffMs = computeBackoffMs(
  417. config_.writeRetryBackoffMs, attempt,
  418. config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
  419. spdlog::warn("Client::upsert {}; retrying in {}ms (attempt {}/{})",
  420. status.error_message(), backoffMs,
  421. attempt + 1, config_.writeRetries);
  422. std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
  423. }
  424. }
  425. if (!status.ok()) {
  426. spdlog::error("Client::upsert failed after retries: {}", status.error_message());
  427. throw std::runtime_error(status.error_message());
  428. }
  429. return {response.id(), response.inserted()};
  430. }
  431. bool remove(const std::string& collection, const std::string& id) {
  432. smartbotic::databasepb::DeleteRequest request;
  433. request.set_collection(qualify(collection));
  434. request.set_id(id);
  435. smartbotic::databasepb::DeleteResponse response;
  436. grpc::Status status;
  437. for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
  438. grpc::ClientContext context;
  439. setDeadline(context);
  440. response.Clear();
  441. status = stub_->Delete(&context, request, &response);
  442. if (status.ok() || !retryableStatus(status)) {
  443. break;
  444. }
  445. if (attempt < config_.writeRetries) {
  446. uint32_t backoffMs = computeBackoffMs(
  447. config_.writeRetryBackoffMs, attempt,
  448. config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
  449. spdlog::warn("Client::remove {}; retrying in {}ms (attempt {}/{})",
  450. status.error_message(), backoffMs,
  451. attempt + 1, config_.writeRetries);
  452. std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
  453. }
  454. }
  455. if (!status.ok()) {
  456. spdlog::error("Client::remove failed after retries: {}", status.error_message());
  457. return false;
  458. }
  459. return response.deleted();
  460. }
  461. bool exists(const std::string& collection, const std::string& id) {
  462. smartbotic::databasepb::ExistsRequest request;
  463. request.set_collection(qualify(collection));
  464. request.set_id(id);
  465. smartbotic::databasepb::ExistsResponse response;
  466. grpc::ClientContext context;
  467. setDeadline(context);
  468. auto status = stub_->Exists(&context, request, &response);
  469. if (!status.ok()) {
  470. // v2.8.0 — a permission denial must NOT look like "no data".
  471. // Returning empty here would leave an application unable to tell
  472. // "you may not see this" from "there is nothing to see", so it
  473. // would silently take the wrong branch. Same reasoning as the
  474. // v2.4.2 change that stopped Find returning empty for a missing
  475. // collection.
  476. if (status.error_code() == grpc::StatusCode::PERMISSION_DENIED) {
  477. throw std::runtime_error("access denied: exists");
  478. }
  479. spdlog::error("Client::exists failed: {}", status.error_message());
  480. return false;
  481. }
  482. return response.exists();
  483. }
  484. // ===== Query Operations =====
  485. std::vector<nlohmann::json> find(const std::string& collection,
  486. const Client::QueryOptions& options,
  487. const std::vector<std::string>& projection = {}) {
  488. smartbotic::databasepb::FindRequest request;
  489. request.set_collection(qualify(collection));
  490. // Set filters
  491. for (const auto& filter : options.filters) {
  492. auto* pb = request.add_filters();
  493. pb->set_field(filter.field);
  494. pb->set_value(filter.value.dump());
  495. pb->set_op(static_cast<smartbotic::databasepb::FilterOp>(filter.op));
  496. }
  497. // Set sorting
  498. if (!options.sortField.empty()) {
  499. auto* sort = request.mutable_sort();
  500. sort->set_field(options.sortField);
  501. sort->set_descending(options.sortDescending);
  502. }
  503. // Set pagination
  504. request.set_limit(options.limit);
  505. request.set_offset(options.offset);
  506. // v2.7.1 — field projection, passed alongside QueryOptions rather than
  507. // inside it: adding a member would change the struct size, and with an
  508. // unchanged soname an installed consumer crashes. See client.hpp.
  509. for (const auto& f : projection) {
  510. request.add_projection(f);
  511. }
  512. smartbotic::databasepb::FindResponse response;
  513. grpc::ClientContext context;
  514. setDeadline(context);
  515. auto status = stub_->Find(&context, request, &response);
  516. if (status.error_code() == grpc::StatusCode::NOT_FOUND) {
  517. // Surfaced rather than logged: an unknown collection or view is a
  518. // caller bug, and returning {} here is indistinguishable from a
  519. // legitimately empty result. Other failures keep the old
  520. // log-and-return-empty behaviour.
  521. throw std::runtime_error(status.error_message());
  522. }
  523. if (!status.ok()) {
  524. // v2.8.0 — a permission denial must NOT look like "no data".
  525. // Returning empty here would leave an application unable to tell
  526. // "you may not see this" from "there is nothing to see", so it
  527. // would silently take the wrong branch. Same reasoning as the
  528. // v2.4.2 change that stopped Find returning empty for a missing
  529. // collection.
  530. if (status.error_code() == grpc::StatusCode::PERMISSION_DENIED) {
  531. throw std::runtime_error("access denied: find");
  532. }
  533. spdlog::error("Client::find failed: {}", status.error_message());
  534. return {};
  535. }
  536. std::vector<nlohmann::json> results;
  537. results.reserve(response.documents_size());
  538. for (const auto& doc : response.documents()) {
  539. auto json = nlohmann::json::parse(doc.data());
  540. json["_id"] = doc.id();
  541. json["_version"] = doc.version();
  542. json["_created_at"] = doc.created_at();
  543. json["_updated_at"] = doc.updated_at();
  544. json["_created_by"] = doc.created_by();
  545. json["_updated_by"] = doc.updated_by();
  546. results.push_back(json);
  547. }
  548. return results;
  549. }
  550. Client::FindResult findWithMetrics(const std::string& collection,
  551. const Client::QueryOptions& options,
  552. const std::vector<std::string>& projection = {}) {
  553. smartbotic::databasepb::FindRequest request;
  554. request.set_collection(qualify(collection));
  555. // Set filters
  556. for (const auto& filter : options.filters) {
  557. auto* pb = request.add_filters();
  558. pb->set_field(filter.field);
  559. pb->set_value(filter.value.dump());
  560. pb->set_op(static_cast<smartbotic::databasepb::FilterOp>(filter.op));
  561. }
  562. // Set sorting
  563. if (!options.sortField.empty()) {
  564. auto* sort = request.mutable_sort();
  565. sort->set_field(options.sortField);
  566. sort->set_descending(options.sortDescending);
  567. }
  568. // Set pagination
  569. request.set_limit(options.limit);
  570. request.set_offset(options.offset);
  571. // v2.7.1 — field projection, passed alongside QueryOptions rather than
  572. // inside it: adding a member would change the struct size, and with an
  573. // unchanged soname an installed consumer crashes. See client.hpp.
  574. for (const auto& f : projection) {
  575. request.add_projection(f);
  576. }
  577. smartbotic::databasepb::FindResponse response;
  578. grpc::ClientContext context;
  579. setDeadline(context);
  580. auto status = stub_->Find(&context, request, &response);
  581. if (status.error_code() == grpc::StatusCode::NOT_FOUND) {
  582. throw std::runtime_error(status.error_message()); // see find()
  583. }
  584. if (!status.ok()) {
  585. // v2.8.0 — a permission denial must NOT look like "no data".
  586. // Returning empty here would leave an application unable to tell
  587. // "you may not see this" from "there is nothing to see", so it
  588. // would silently take the wrong branch. Same reasoning as the
  589. // v2.4.2 change that stopped Find returning empty for a missing
  590. // collection.
  591. if (status.error_code() == grpc::StatusCode::PERMISSION_DENIED) {
  592. throw std::runtime_error("access denied: findWithMetrics");
  593. }
  594. spdlog::error("Client::findWithMetrics failed: {}", status.error_message());
  595. return {};
  596. }
  597. Client::FindResult result;
  598. result.totalCount = response.total_count();
  599. result.hasMore = response.has_more();
  600. result.usedWalFallback = response.used_wal_fallback();
  601. result.memoryMatchCount = response.memory_match_count();
  602. result.walMatchCount = response.wal_match_count();
  603. result.memorySearchMicros = response.memory_search_micros();
  604. result.walSearchMicros = response.wal_search_micros();
  605. result.documents.reserve(response.documents_size());
  606. for (const auto& doc : response.documents()) {
  607. auto json = nlohmann::json::parse(doc.data());
  608. json["_id"] = doc.id();
  609. json["_version"] = doc.version();
  610. json["_created_at"] = doc.created_at();
  611. json["_updated_at"] = doc.updated_at();
  612. json["_created_by"] = doc.created_by();
  613. json["_updated_by"] = doc.updated_by();
  614. result.documents.push_back(json);
  615. }
  616. return result;
  617. }
  618. uint64_t count(const std::string& collection,
  619. const std::vector<Client::Filter>& filters) {
  620. smartbotic::databasepb::CountRequest request;
  621. request.set_collection(qualify(collection));
  622. for (const auto& filter : filters) {
  623. auto* pb = request.add_filters();
  624. pb->set_field(filter.field);
  625. pb->set_value(filter.value.dump());
  626. pb->set_op(static_cast<smartbotic::databasepb::FilterOp>(filter.op));
  627. }
  628. smartbotic::databasepb::CountResponse response;
  629. grpc::ClientContext context;
  630. setDeadline(context);
  631. auto status = stub_->Count(&context, request, &response);
  632. if (!status.ok()) {
  633. // v2.8.0 — a permission denial must NOT look like "no data".
  634. // Returning empty here would leave an application unable to tell
  635. // "you may not see this" from "there is nothing to see", so it
  636. // would silently take the wrong branch. Same reasoning as the
  637. // v2.4.2 change that stopped Find returning empty for a missing
  638. // collection.
  639. if (status.error_code() == grpc::StatusCode::PERMISSION_DENIED) {
  640. throw std::runtime_error("access denied: count");
  641. }
  642. spdlog::error("Client::count failed: {}", status.error_message());
  643. return 0;
  644. }
  645. return response.count();
  646. }
  647. // ===== Set Operations =====
  648. bool setAdd(const std::string& collection, const std::string& setId, const std::string& member) {
  649. smartbotic::databasepb::SetAddRequest request;
  650. request.set_collection(qualify(collection));
  651. request.set_set_id(setId);
  652. request.set_member(member);
  653. smartbotic::databasepb::SetAddResponse response;
  654. grpc::ClientContext context;
  655. setDeadline(context);
  656. auto status = stub_->SetAdd(&context, request, &response);
  657. if (!status.ok()) {
  658. spdlog::error("Client::setAdd failed: {}", status.error_message());
  659. return false;
  660. }
  661. return response.added();
  662. }
  663. bool setRemove(const std::string& collection, const std::string& setId, const std::string& member) {
  664. smartbotic::databasepb::SetRemoveRequest request;
  665. request.set_collection(qualify(collection));
  666. request.set_set_id(setId);
  667. request.set_member(member);
  668. smartbotic::databasepb::SetRemoveResponse response;
  669. grpc::ClientContext context;
  670. setDeadline(context);
  671. auto status = stub_->SetRemove(&context, request, &response);
  672. if (!status.ok()) {
  673. spdlog::error("Client::setRemove failed: {}", status.error_message());
  674. return false;
  675. }
  676. return response.removed();
  677. }
  678. std::vector<std::string> setMembers(const std::string& collection, const std::string& setId) {
  679. smartbotic::databasepb::SetMembersRequest request;
  680. request.set_collection(qualify(collection));
  681. request.set_set_id(setId);
  682. smartbotic::databasepb::SetMembersResponse response;
  683. grpc::ClientContext context;
  684. setDeadline(context);
  685. auto status = stub_->SetMembers(&context, request, &response);
  686. if (!status.ok()) {
  687. spdlog::error("Client::setMembers failed: {}", status.error_message());
  688. return {};
  689. }
  690. return {response.members().begin(), response.members().end()};
  691. }
  692. bool setIsMember(const std::string& collection, const std::string& setId, const std::string& member) {
  693. smartbotic::databasepb::SetIsMemberRequest request;
  694. request.set_collection(qualify(collection));
  695. request.set_set_id(setId);
  696. request.set_member(member);
  697. smartbotic::databasepb::SetIsMemberResponse response;
  698. grpc::ClientContext context;
  699. setDeadline(context);
  700. auto status = stub_->SetIsMember(&context, request, &response);
  701. if (!status.ok()) {
  702. spdlog::error("Client::setIsMember failed: {}", status.error_message());
  703. return false;
  704. }
  705. return response.is_member();
  706. }
  707. // ===== Collection Management =====
  708. std::vector<Client::SimilarityResult> similaritySearch(
  709. const std::string& collection, const std::vector<float>& queryVector,
  710. uint32_t topK, float minScore) {
  711. smartbotic::databasepb::SimilaritySearchRequest request;
  712. request.set_collection(qualify(collection));
  713. for (float v : queryVector) {
  714. request.add_query_vector(v);
  715. }
  716. request.set_top_k(topK);
  717. request.set_min_score(minScore);
  718. smartbotic::databasepb::SimilaritySearchResponse response;
  719. grpc::ClientContext context;
  720. setDeadline(context);
  721. auto status = stub_->SimilaritySearch(&context, request, &response);
  722. if (!status.ok()) {
  723. spdlog::error("Client::similaritySearch failed: {}", status.error_message());
  724. throw std::runtime_error(status.error_message());
  725. }
  726. std::vector<Client::SimilarityResult> results;
  727. results.reserve(response.results_size());
  728. for (const auto& r : response.results()) {
  729. Client::SimilarityResult entry;
  730. entry.id = r.id();
  731. entry.score = r.score();
  732. if (!r.data().empty()) {
  733. entry.data = nlohmann::json::parse(r.data());
  734. }
  735. results.push_back(std::move(entry));
  736. }
  737. return results;
  738. }
  739. bool createCollection(const std::string& name, uint32_t defaultTtlSeconds, bool encrypted,
  740. uint32_t maxVersions, uint32_t vectorDimension) {
  741. smartbotic::databasepb::CreateCollectionRequest request;
  742. // v2.4.5 — MUST qualify. Sending a bare name here while insert/get/
  743. // find all qualify created TWO collections from one call site: the
  744. // bare one received the options (encrypted, maxVersions,
  745. // vectorDimension, TTL) and stayed empty, while the qualified one
  746. // was created implicitly by the first insert with DEFAULTS. Silent,
  747. // and unrecoverable for vectorDimension, which is immutable after
  748. // creation. Same class of bug as the v2.4.2 createView break.
  749. request.set_name(qualify(name));
  750. auto* options = request.mutable_options();
  751. if (defaultTtlSeconds > 0) {
  752. options->set_default_ttl_seconds(defaultTtlSeconds);
  753. }
  754. options->set_encrypted(encrypted);
  755. if (maxVersions > 0) {
  756. options->set_max_versions(maxVersions);
  757. }
  758. if (vectorDimension > 0) {
  759. options->set_vector_dimension(vectorDimension);
  760. }
  761. smartbotic::databasepb::CreateCollectionResponse response;
  762. grpc::ClientContext context;
  763. setDeadline(context);
  764. auto status = stub_->CreateCollection(&context, request, &response);
  765. if (!status.ok()) {
  766. spdlog::error("Client::createCollection failed: {}", status.error_message());
  767. return false;
  768. }
  769. return response.created();
  770. }
  771. bool dropCollection(const std::string& name) {
  772. smartbotic::databasepb::DropCollectionRequest request;
  773. // v2.4.5 — MUST qualify. Unqualified, this dropped the empty phantom
  774. // created by createCollection and returned TRUE while every document
  775. // in <project>:<name> remained readable — a deletion request that
  776. // reports success and deletes nothing.
  777. request.set_name(qualify(name));
  778. smartbotic::databasepb::DropCollectionResponse response;
  779. grpc::ClientContext context;
  780. setDeadline(context);
  781. auto status = stub_->DropCollection(&context, request, &response);
  782. if (!status.ok()) {
  783. spdlog::error("Client::dropCollection failed: {}", status.error_message());
  784. return false;
  785. }
  786. return response.dropped();
  787. }
  788. std::vector<std::string> listCollections() {
  789. smartbotic::databasepb::ListCollectionsRequest request;
  790. smartbotic::databasepb::ListCollectionsResponse response;
  791. grpc::ClientContext context;
  792. setDeadline(context);
  793. auto status = stub_->ListCollections(&context, request, &response);
  794. if (!status.ok()) {
  795. spdlog::error("Client::listCollections failed: {}", status.error_message());
  796. return {};
  797. }
  798. // v2.4.5 — round-trip the naming: callers pass bare names in, so they
  799. // get bare names back for their OWN project. Names belonging to other
  800. // projects stay visibly qualified rather than being flattened into
  801. // ours (same rule as listViews). Note the server still returns every
  802. // project's collections; ListCollectionsRequest has no project filter
  803. // yet, unlike ListViewsRequest.
  804. std::vector<std::string> out;
  805. out.reserve(static_cast<size_t>(response.names_size()));
  806. for (const auto& n : response.names()) {
  807. out.push_back(unqualify(n));
  808. }
  809. return out;
  810. }
  811. std::optional<Client::CollectionInfo> getCollectionInfo(const std::string& name) {
  812. smartbotic::databasepb::GetCollectionInfoRequest request;
  813. // v2.4.5 — MUST qualify, else this reports on the phantom (0 docs)
  814. // rather than the collection the caller reads and writes.
  815. request.set_name(qualify(name));
  816. smartbotic::databasepb::GetCollectionInfoResponse response;
  817. grpc::ClientContext context;
  818. setDeadline(context);
  819. auto status = stub_->GetCollectionInfo(&context, request, &response);
  820. if (!status.ok()) {
  821. spdlog::error("Client::getCollectionInfo failed: {}", status.error_message());
  822. return std::nullopt;
  823. }
  824. if (!response.found()) {
  825. return std::nullopt;
  826. }
  827. const auto& info = response.info();
  828. Client::CollectionInfo result;
  829. result.name = info.name();
  830. result.documentCount = info.document_count();
  831. result.sizeBytes = info.size_bytes();
  832. result.defaultTtlSeconds = info.options().default_ttl_seconds();
  833. result.encrypted = info.options().encrypted();
  834. result.maxVersions = info.options().max_versions();
  835. result.createdAt = info.created_at();
  836. result.updatedAt = info.updated_at();
  837. return result;
  838. }
  839. // ===== Project Management (v2.3 Stage F) =====
  840. std::vector<std::string> listProjects() {
  841. smartbotic::databasepb::ListProjectsRequest request;
  842. smartbotic::databasepb::ListProjectsResponse response;
  843. grpc::ClientContext context;
  844. setDeadline(context);
  845. auto status = stub_->ListProjects(&context, request, &response);
  846. if (!status.ok()) {
  847. spdlog::error("Client::listProjects failed: {}", status.error_message());
  848. return {};
  849. }
  850. return {response.projects().begin(), response.projects().end()};
  851. }
  852. bool createProject(const std::string& name) {
  853. smartbotic::databasepb::CreateProjectRequest request;
  854. request.set_name(name);
  855. smartbotic::databasepb::CreateProjectResponse response;
  856. grpc::ClientContext context;
  857. setDeadline(context);
  858. auto status = stub_->CreateProject(&context, request, &response);
  859. if (!status.ok()) {
  860. spdlog::error("Client::createProject failed: {}", status.error_message());
  861. throw std::runtime_error("createProject transport failure: " +
  862. status.error_message());
  863. }
  864. if (!response.error().empty()) {
  865. // Validation error (invalid name etc.) — surfaced via the
  866. // response field, not the gRPC status.
  867. throw std::runtime_error("createProject('" + name + "') rejected: " +
  868. response.error());
  869. }
  870. return response.created();
  871. }
  872. void dropProject(const std::string& name) {
  873. smartbotic::databasepb::DropProjectRequest request;
  874. request.set_name(name);
  875. smartbotic::databasepb::DropProjectResponse response;
  876. grpc::ClientContext context;
  877. setDeadline(context);
  878. auto status = stub_->DropProject(&context, request, &response);
  879. if (!status.ok()) {
  880. spdlog::error("Client::dropProject failed: {}", status.error_message());
  881. throw std::runtime_error("dropProject transport failure: " +
  882. status.error_message());
  883. }
  884. if (!response.error().empty()) {
  885. throw std::runtime_error("dropProject('" + name + "') refused: " +
  886. response.error());
  887. }
  888. // dropped() may still be false here defensively — but the
  889. // server-side contract is dropped==true iff error is empty.
  890. }
  891. // ===== Collection Configuration =====
  892. bool configureCollection(const std::string& collection, const Client::CollectionConfig& cfg) {
  893. smartbotic::databasepb::ConfigureCollectionRequest request;
  894. request.set_collection(qualify(collection));
  895. request.mutable_config()->set_timestamp_precision(cfg.timestampPrecision);
  896. // v2.4.5 — only send versioning when the caller actually set it, so
  897. // the server's partial update leaves it alone otherwise.
  898. if (cfg.versioningEnabled.has_value()) {
  899. request.mutable_config()->set_versioning_enabled(*cfg.versioningEnabled);
  900. }
  901. smartbotic::databasepb::ConfigureCollectionResponse response;
  902. grpc::ClientContext context;
  903. setDeadline(context);
  904. auto status = stub_->ConfigureCollection(&context, request, &response);
  905. if (!status.ok()) {
  906. spdlog::error("Client::configureCollection failed: {}", status.error_message());
  907. return false;
  908. }
  909. if (!response.success()) {
  910. spdlog::error("Client::configureCollection rejected: {}", response.error());
  911. return false;
  912. }
  913. return true;
  914. }
  915. Client::CollectionConfig getCollectionConfig(const std::string& collection) {
  916. smartbotic::databasepb::GetCollectionConfigRequest request;
  917. request.set_collection(qualify(collection));
  918. smartbotic::databasepb::GetCollectionConfigResponse response;
  919. grpc::ClientContext context;
  920. setDeadline(context);
  921. Client::CollectionConfig out;
  922. auto status = stub_->GetCollectionConfig(&context, request, &response);
  923. if (!status.ok()) {
  924. spdlog::error("Client::getCollectionConfig failed: {}", status.error_message());
  925. return out;
  926. }
  927. out.timestampPrecision = response.config().timestamp_precision();
  928. if (out.timestampPrecision.empty()) out.timestampPrecision = "ms";
  929. // v2.4.5 — always populated on read, so callers get the effective
  930. // value rather than the "leave unchanged" sentinel. A server older
  931. // than 2.4.5 omits the field; treat that as versioning on, which is
  932. // what those builds always did.
  933. out.versioningEnabled = response.config().has_versioning_enabled()
  934. ? response.config().versioning_enabled()
  935. : true;
  936. return out;
  937. }
  938. bool hasCollectionConfig(const std::string& collection) {
  939. smartbotic::databasepb::GetCollectionConfigRequest request;
  940. request.set_collection(qualify(collection));
  941. smartbotic::databasepb::GetCollectionConfigResponse response;
  942. grpc::ClientContext context;
  943. setDeadline(context);
  944. auto status = stub_->GetCollectionConfig(&context, request, &response);
  945. if (!status.ok()) return false;
  946. return response.found();
  947. }
  948. Client::TimestampMigrationResult migrateCollectionTimestamps(
  949. const std::string& collection,
  950. const std::string& fromPrecision,
  951. const std::string& toPrecision
  952. ) {
  953. smartbotic::databasepb::MigrateCollectionTimestampsRequest request;
  954. request.set_collection(qualify(collection));
  955. request.set_from_precision(fromPrecision);
  956. request.set_to_precision(toPrecision);
  957. smartbotic::databasepb::MigrateCollectionTimestampsResponse response;
  958. grpc::ClientContext context;
  959. // Longer deadline — can iterate many rows
  960. auto deadline = std::chrono::system_clock::now() + std::chrono::minutes(10);
  961. context.set_deadline(deadline);
  962. attachAuth(context);
  963. Client::TimestampMigrationResult out;
  964. auto status = stub_->MigrateCollectionTimestamps(&context, request, &response);
  965. if (!status.ok()) {
  966. out.success = false;
  967. out.error = status.error_message();
  968. return out;
  969. }
  970. out.success = response.success();
  971. out.error = response.error();
  972. out.rowsMigrated = response.rows_migrated();
  973. out.rowsSkipped = response.rows_skipped();
  974. return out;
  975. }
  976. // ===== View Management =====
  977. bool createView(const std::string& name,
  978. const std::string& collection,
  979. const std::vector<std::string>& include,
  980. const std::vector<std::string>& exclude,
  981. const std::vector<Client::Filter>& where,
  982. const std::optional<Client::Sort>& defaultSort) {
  983. smartbotic::databasepb::CreateViewRequest request;
  984. // Both fields are namespaced. Qualifying `collection` but not `name`
  985. // is what broke view lookups in v2.3: the registry keyed on the bare
  986. // name while every read path sent the qualified one, so the keys
  987. // could never meet.
  988. request.set_name(qualify(name));
  989. request.set_collection(qualify(collection));
  990. for (const auto& p : include) request.add_include(p);
  991. for (const auto& p : exclude) request.add_exclude(p);
  992. for (const auto& f : where) {
  993. auto* pb = request.add_where();
  994. pb->set_field(f.field);
  995. pb->set_op(static_cast<smartbotic::databasepb::FilterOp>(f.op));
  996. pb->set_value(f.value.dump());
  997. }
  998. if (defaultSort) {
  999. auto* s = request.mutable_default_sort();
  1000. s->set_field(defaultSort->field);
  1001. s->set_descending(defaultSort->descending);
  1002. }
  1003. smartbotic::databasepb::CreateViewResponse response;
  1004. grpc::ClientContext context;
  1005. setDeadline(context);
  1006. auto status = stub_->CreateView(&context, request, &response);
  1007. if (!status.ok()) {
  1008. spdlog::error("Client::createView failed: {}", status.error_message());
  1009. return false;
  1010. }
  1011. if (!response.success()) {
  1012. spdlog::error("Client::createView rejected: {}", response.error());
  1013. return false;
  1014. }
  1015. return true;
  1016. }
  1017. bool dropView(const std::string& name) {
  1018. smartbotic::databasepb::DropViewRequest request;
  1019. request.set_name(qualify(name));
  1020. smartbotic::databasepb::DropViewResponse response;
  1021. grpc::ClientContext context;
  1022. setDeadline(context);
  1023. auto status = stub_->DropView(&context, request, &response);
  1024. if (!status.ok()) {
  1025. spdlog::error("Client::dropView failed: {}", status.error_message());
  1026. return false;
  1027. }
  1028. return response.success();
  1029. }
  1030. std::vector<Client::ViewDefinition> listViews() {
  1031. smartbotic::databasepb::ListViewsRequest request;
  1032. // Scope the listing to this client's workspace.
  1033. request.set_project(config_.project);
  1034. smartbotic::databasepb::ListViewsResponse response;
  1035. grpc::ClientContext context;
  1036. setDeadline(context);
  1037. std::vector<Client::ViewDefinition> out;
  1038. auto status = stub_->ListViews(&context, request, &response);
  1039. if (!status.ok()) {
  1040. spdlog::error("Client::listViews failed: {}", status.error_message());
  1041. return out;
  1042. }
  1043. for (const auto& pb : response.views()) {
  1044. Client::ViewDefinition v;
  1045. // Round-trip: the caller created "adults", so list it as "adults".
  1046. v.name = unqualify(pb.name());
  1047. v.collection = unqualify(pb.collection());
  1048. for (const auto& p : pb.include()) v.include.push_back(p);
  1049. for (const auto& p : pb.exclude()) v.exclude.push_back(p);
  1050. for (const auto& pbf : pb.where()) {
  1051. Client::Filter f;
  1052. f.field = pbf.field();
  1053. f.op = static_cast<Client::FilterOp>(pbf.op());
  1054. try { f.value = nlohmann::json::parse(pbf.value()); }
  1055. catch (...) { f.value = pbf.value(); }
  1056. v.where.push_back(f);
  1057. }
  1058. if (pb.has_default_sort()) {
  1059. Client::Sort s;
  1060. s.field = pb.default_sort().field();
  1061. s.descending = pb.default_sort().descending();
  1062. v.defaultSort = s;
  1063. }
  1064. v.createdAt = pb.created_at();
  1065. v.updatedAt = pb.updated_at();
  1066. out.push_back(std::move(v));
  1067. }
  1068. return out;
  1069. }
  1070. std::optional<Client::ViewDefinition> getViewInfo(const std::string& name) {
  1071. smartbotic::databasepb::GetViewInfoRequest request;
  1072. request.set_name(qualify(name));
  1073. smartbotic::databasepb::GetViewInfoResponse response;
  1074. grpc::ClientContext context;
  1075. setDeadline(context);
  1076. auto status = stub_->GetViewInfo(&context, request, &response);
  1077. if (!status.ok() || !response.found()) {
  1078. return std::nullopt;
  1079. }
  1080. Client::ViewDefinition v;
  1081. v.name = unqualify(response.view().name());
  1082. v.collection = unqualify(response.view().collection());
  1083. for (const auto& p : response.view().include()) v.include.push_back(p);
  1084. for (const auto& p : response.view().exclude()) v.exclude.push_back(p);
  1085. for (const auto& pbf : response.view().where()) {
  1086. Client::Filter f;
  1087. f.field = pbf.field();
  1088. f.op = static_cast<Client::FilterOp>(pbf.op());
  1089. try { f.value = nlohmann::json::parse(pbf.value()); }
  1090. catch (...) { f.value = pbf.value(); }
  1091. v.where.push_back(f);
  1092. }
  1093. if (response.view().has_default_sort()) {
  1094. Client::Sort s;
  1095. s.field = response.view().default_sort().field();
  1096. s.descending = response.view().default_sort().descending();
  1097. v.defaultSort = s;
  1098. }
  1099. v.createdAt = response.view().created_at();
  1100. v.updatedAt = response.view().updated_at();
  1101. return v;
  1102. }
  1103. // ===== Event Subscription =====
  1104. class SubscriptionHandle {
  1105. public:
  1106. SubscriptionHandle(std::shared_ptr<grpc::ClientContext> ctx,
  1107. std::unique_ptr<grpc::ClientReader<smartbotic::databasepb::DatabaseEvent>> reader,
  1108. std::thread readerThread)
  1109. : context_(std::move(ctx))
  1110. , reader_(std::move(reader))
  1111. , readerThread_(std::move(readerThread))
  1112. , active_(true) {}
  1113. ~SubscriptionHandle() {
  1114. cancel();
  1115. }
  1116. void cancel() {
  1117. if (active_.exchange(false)) {
  1118. context_->TryCancel();
  1119. if (readerThread_.joinable()) {
  1120. readerThread_.join();
  1121. }
  1122. }
  1123. }
  1124. private:
  1125. std::shared_ptr<grpc::ClientContext> context_;
  1126. std::unique_ptr<grpc::ClientReader<smartbotic::databasepb::DatabaseEvent>> reader_;
  1127. std::thread readerThread_;
  1128. std::atomic<bool> active_;
  1129. };
  1130. std::shared_ptr<void> subscribe(const std::vector<std::string>& collections,
  1131. Client::EventCallback callback) {
  1132. auto context = std::make_shared<grpc::ClientContext>();
  1133. // Subscribe is a long-running stream — no deadline — but the auth
  1134. // metadata still has to land on the initial request headers.
  1135. attachAuth(*context);
  1136. smartbotic::databasepb::SubscribeRequest request;
  1137. // v2.8.0 — subscribe was the one call left unqualified while every
  1138. // document call qualified. Two consequences, both cross-project leaks:
  1139. //
  1140. // * a bare name matched nothing, because events carry the qualified
  1141. // collection - so subscriptions silently never fired;
  1142. // * an EMPTY list means "every collection" server-side, which meant
  1143. // every collection in every PROJECT.
  1144. //
  1145. // Named collections are qualified like any other. For the "everything"
  1146. // case we send a `<project>:*` pattern instead of an empty request, so
  1147. // "all" means all of MY project. That needs no proto change - the
  1148. // pattern field already existed.
  1149. for (const auto& coll : collections) {
  1150. request.add_collections(qualify(coll));
  1151. }
  1152. if (collections.empty()) {
  1153. request.add_patterns(config_.project + ":*");
  1154. }
  1155. request.set_include_data(true);
  1156. auto reader = stub_->Subscribe(context.get(), request);
  1157. // Create reader thread
  1158. auto readerThread = std::thread([reader = reader.get(), callback = std::move(callback)]() {
  1159. smartbotic::databasepb::DatabaseEvent event;
  1160. while (reader->Read(&event)) {
  1161. std::optional<nlohmann::json> data;
  1162. if (!event.data().empty()) {
  1163. data = nlohmann::json::parse(event.data());
  1164. }
  1165. // Convert proto event type to string
  1166. std::string eventType;
  1167. switch (event.type()) {
  1168. case smartbotic::databasepb::EVENT_INSERT: eventType = "insert"; break;
  1169. case smartbotic::databasepb::EVENT_UPDATE: eventType = "update"; break;
  1170. case smartbotic::databasepb::EVENT_DELETE: eventType = "delete"; break;
  1171. case smartbotic::databasepb::EVENT_EXPIRE: eventType = "expire"; break;
  1172. case smartbotic::databasepb::EVENT_INVALIDATE: eventType = "invalidate"; break;
  1173. default: eventType = "unknown"; break;
  1174. }
  1175. callback(
  1176. event.collection(),
  1177. event.document_id(),
  1178. eventType,
  1179. data
  1180. );
  1181. }
  1182. });
  1183. return std::make_shared<SubscriptionHandle>(
  1184. std::move(context),
  1185. std::move(reader),
  1186. std::move(readerThread)
  1187. );
  1188. }
  1189. // ===== Version History =====
  1190. Client::VersionHistoryResult getVersionHistory(const std::string& collection,
  1191. const std::string& id, uint32_t limit, uint32_t offset) {
  1192. smartbotic::databasepb::GetVersionHistoryRequest request;
  1193. request.set_collection(qualify(collection));
  1194. request.set_id(id);
  1195. if (limit > 0) request.set_limit(limit);
  1196. if (offset > 0) request.set_offset(offset);
  1197. smartbotic::databasepb::GetVersionHistoryResponse response;
  1198. grpc::ClientContext context;
  1199. setDeadline(context);
  1200. auto status = stub_->GetVersionHistory(&context, request, &response);
  1201. if (!status.ok()) {
  1202. spdlog::error("Client::getVersionHistory failed: {}", status.error_message());
  1203. return {};
  1204. }
  1205. Client::VersionHistoryResult result;
  1206. result.currentVersion = response.current_version();
  1207. result.totalCount = response.total_count();
  1208. result.documentDeleted = response.document_deleted();
  1209. result.versions.reserve(response.versions_size());
  1210. for (const auto& ver : response.versions()) {
  1211. Client::VersionEntry entry;
  1212. entry.version = ver.version();
  1213. entry.timestamp = ver.timestamp();
  1214. entry.updatedBy = ver.updated_by();
  1215. if (!ver.data().empty()) {
  1216. entry.data = nlohmann::json::parse(ver.data());
  1217. }
  1218. result.versions.push_back(std::move(entry));
  1219. }
  1220. return result;
  1221. }
  1222. std::optional<Client::VersionEntry> getDocumentVersion(const std::string& collection,
  1223. const std::string& id, uint64_t version) {
  1224. smartbotic::databasepb::GetDocumentVersionRequest request;
  1225. request.set_collection(qualify(collection));
  1226. request.set_id(id);
  1227. request.set_version(version);
  1228. smartbotic::databasepb::GetDocumentVersionResponse response;
  1229. grpc::ClientContext context;
  1230. setDeadline(context);
  1231. auto status = stub_->GetDocumentVersion(&context, request, &response);
  1232. if (!status.ok()) {
  1233. spdlog::error("Client::getDocumentVersion failed: {}", status.error_message());
  1234. return std::nullopt;
  1235. }
  1236. if (!response.found()) {
  1237. return std::nullopt;
  1238. }
  1239. Client::VersionEntry entry;
  1240. entry.version = response.version_entry().version();
  1241. entry.timestamp = response.version_entry().timestamp();
  1242. entry.updatedBy = response.version_entry().updated_by();
  1243. if (!response.version_entry().data().empty()) {
  1244. entry.data = nlohmann::json::parse(response.version_entry().data());
  1245. }
  1246. return entry;
  1247. }
  1248. uint64_t restoreVersion(const std::string& collection, const std::string& id,
  1249. uint64_t version, const std::string& actor) {
  1250. smartbotic::databasepb::RestoreVersionRequest request;
  1251. request.set_collection(qualify(collection));
  1252. request.set_id(id);
  1253. request.set_version(version);
  1254. if (!actor.empty()) request.set_actor(actor);
  1255. smartbotic::databasepb::RestoreVersionResponse response;
  1256. grpc::ClientContext context;
  1257. setDeadline(context);
  1258. auto status = stub_->RestoreVersion(&context, request, &response);
  1259. if (!status.ok()) {
  1260. spdlog::error("Client::restoreVersion failed: {}", status.error_message());
  1261. return 0;
  1262. }
  1263. if (!response.success()) {
  1264. spdlog::error("Client::restoreVersion: {}", response.error());
  1265. return 0;
  1266. }
  1267. return response.new_version();
  1268. }
  1269. std::pair<uint64_t, uint64_t> restoreToDate(const std::string& collection,
  1270. const std::string& id, uint64_t timestamp, const std::string& actor) {
  1271. smartbotic::databasepb::RestoreToDateRequest request;
  1272. request.set_collection(qualify(collection));
  1273. request.set_id(id);
  1274. request.set_timestamp(timestamp);
  1275. if (!actor.empty()) request.set_actor(actor);
  1276. smartbotic::databasepb::RestoreToDateResponse response;
  1277. grpc::ClientContext context;
  1278. setDeadline(context);
  1279. auto status = stub_->RestoreToDate(&context, request, &response);
  1280. if (!status.ok()) {
  1281. spdlog::error("Client::restoreToDate failed: {}", status.error_message());
  1282. return {0, 0};
  1283. }
  1284. if (!response.success()) {
  1285. spdlog::error("Client::restoreToDate: {}", response.error());
  1286. return {0, 0};
  1287. }
  1288. return {response.restored_version(), response.new_version()};
  1289. }
  1290. // ===== Health =====
  1291. bool healthCheck() {
  1292. smartbotic::databasepb::HealthCheckRequest request;
  1293. smartbotic::databasepb::HealthCheckResponse response;
  1294. grpc::ClientContext context;
  1295. setDeadline(context);
  1296. auto status = stub_->HealthCheck(&context, request, &response);
  1297. if (!status.ok()) {
  1298. return false;
  1299. }
  1300. return response.healthy();
  1301. }
  1302. std::optional<Client::HealthInfo> getHealthInfo() {
  1303. smartbotic::databasepb::HealthCheckRequest request;
  1304. smartbotic::databasepb::HealthCheckResponse response;
  1305. grpc::ClientContext context;
  1306. setDeadline(context);
  1307. auto status = stub_->HealthCheck(&context, request, &response);
  1308. if (!status.ok()) {
  1309. return std::nullopt;
  1310. }
  1311. Client::HealthInfo info;
  1312. info.healthy = response.healthy();
  1313. info.uptimeMs = response.uptime_seconds() * 1000;
  1314. info.documentCount = response.document_count();
  1315. info.memoryUsedBytes = response.memory_used_bytes();
  1316. info.walSizeBytes = response.wal_size_bytes();
  1317. return info;
  1318. }
  1319. std::optional<Client::StatsInfo> getStats() {
  1320. smartbotic::databasepb::GetStatsRequest request;
  1321. smartbotic::databasepb::GetStatsResponse response;
  1322. grpc::ClientContext context;
  1323. setDeadline(context);
  1324. auto status = stub_->GetStats(&context, request, &response);
  1325. if (!status.ok()) {
  1326. return std::nullopt;
  1327. }
  1328. Client::StatsInfo info;
  1329. info.totalDocuments = response.total_documents();
  1330. info.totalCollections = response.total_collections();
  1331. info.memoryUsedBytes = response.memory_used_bytes();
  1332. info.walSequence = response.wal_sequence();
  1333. info.walSizeBytes = response.wal_size_bytes();
  1334. info.snapshotCount = response.snapshot_count();
  1335. info.lastSnapshotSequence = response.last_snapshot_sequence();
  1336. info.insertCount = response.insert_count();
  1337. info.updateCount = response.update_count();
  1338. info.deleteCount = response.delete_count();
  1339. info.queryCount = response.query_count();
  1340. // Memory eviction stats
  1341. info.evictedDocuments = response.evicted_documents();
  1342. info.totalEvictions = response.total_evictions();
  1343. info.recoveryCount = response.recovery_count();
  1344. // Memory configuration
  1345. info.maxMemoryBytes = response.max_memory_bytes();
  1346. info.evictionThresholdPercent = response.eviction_threshold_percent();
  1347. info.evictionTargetPercent = response.eviction_target_percent();
  1348. // Operation timing (microseconds) - for performance monitoring
  1349. info.getCount = response.get_count();
  1350. info.getTotalMicros = response.get_total_micros();
  1351. info.getMaxMicros = response.get_max_micros();
  1352. info.insertTotalMicros = response.insert_total_micros();
  1353. info.insertMaxMicros = response.insert_max_micros();
  1354. info.updateTotalMicros = response.update_total_micros();
  1355. info.updateMaxMicros = response.update_max_micros();
  1356. info.queryTotalMicros = response.query_total_micros();
  1357. info.queryMaxMicros = response.query_max_micros();
  1358. return info;
  1359. }
  1360. Client::MemoryStats getMemoryStats() {
  1361. smartbotic::databasepb::GetMemoryStatsRequest request;
  1362. smartbotic::databasepb::GetMemoryStatsResponse response;
  1363. grpc::ClientContext context;
  1364. setDeadline(context);
  1365. Client::MemoryStats out;
  1366. auto status = stub_->GetMemoryStats(&context, request, &response);
  1367. if (!status.ok()) {
  1368. spdlog::error("Client::getMemoryStats failed: {}", status.error_message());
  1369. return out;
  1370. }
  1371. out.totalMemoryBytes = response.total_memory_bytes();
  1372. out.maxMemoryBytes = response.max_memory_bytes();
  1373. out.pressurePercent = response.pressure_percent();
  1374. out.pressureLevel = response.pressure_level();
  1375. out.lastEvictionTimestamp = response.last_eviction_timestamp();
  1376. out.lastEvictionDocs = response.last_eviction_docs();
  1377. out.lastEvictionBytesFreed = response.last_eviction_bytes_freed();
  1378. out.collections.reserve(response.collections_size());
  1379. for (const auto& pbc : response.collections()) {
  1380. Client::MemoryCollectionStats c;
  1381. c.collection = pbc.collection();
  1382. c.documentCount = pbc.document_count();
  1383. c.estimatedBytes = pbc.estimated_bytes();
  1384. c.evictedStubCount = pbc.evicted_stub_count();
  1385. switch (pbc.priority()) {
  1386. case smartbotic::databasepb::MEMORY_PRIORITY_LOW: c.priority = "low"; break;
  1387. case smartbotic::databasepb::MEMORY_PRIORITY_NORMAL: c.priority = "normal"; break;
  1388. case smartbotic::databasepb::MEMORY_PRIORITY_HIGH: c.priority = "high"; break;
  1389. default: c.priority = "normal"; break;
  1390. }
  1391. out.collections.push_back(std::move(c));
  1392. }
  1393. return out;
  1394. }
  1395. // ===== Read-Only Control =====
  1396. bool setReadOnly(bool readOnly) {
  1397. smartbotic::databasepb::SetReadOnlyRequest request;
  1398. request.set_read_only(readOnly);
  1399. smartbotic::databasepb::SetReadOnlyResponse response;
  1400. grpc::ClientContext context;
  1401. setDeadline(context);
  1402. auto status = stub_->SetReadOnly(&context, request, &response);
  1403. if (!status.ok()) {
  1404. spdlog::error("Client::setReadOnly failed: {}", status.error_message());
  1405. return false;
  1406. }
  1407. return response.success();
  1408. }
  1409. Client::ReadOnlyStatus getReadOnlyStatus() {
  1410. smartbotic::databasepb::GetReadOnlyStatusRequest request;
  1411. smartbotic::databasepb::GetReadOnlyStatusResponse response;
  1412. grpc::ClientContext context;
  1413. setDeadline(context);
  1414. Client::ReadOnlyStatus out;
  1415. auto status = stub_->GetReadOnlyStatus(&context, request, &response);
  1416. if (!status.ok()) {
  1417. spdlog::error("Client::getReadOnlyStatus failed: {}", status.error_message());
  1418. return out;
  1419. }
  1420. out.readOnly = response.read_only();
  1421. out.reason = response.reason();
  1422. out.recoveryOutcome = response.recovery_outcome();
  1423. out.expectedSnapshot = response.expected_snapshot();
  1424. out.snapshotUsed = response.snapshot_used();
  1425. out.failureReason = response.failure_reason();
  1426. out.walEntriesReplayed = response.wal_entries_replayed();
  1427. out.snapshotsAttempted = response.snapshots_attempted();
  1428. return out;
  1429. }
  1430. // ===== File Operations =====
  1431. Client::FileUploadResult uploadFile(const std::vector<uint8_t>& data,
  1432. const Client::FileUploadMeta& meta,
  1433. uint32_t ttlSeconds = 0,
  1434. bool explicitTtl = false) {
  1435. grpc::ClientContext context;
  1436. auto timeout_ms = config_.timeoutMs + (data.size() / (1024 * 1024)) * 1000;
  1437. context.set_deadline(
  1438. std::chrono::system_clock::now() + std::chrono::milliseconds(timeout_ms));
  1439. attachAuth(context);
  1440. smartbotic::databasepb::UploadFileResponse response;
  1441. auto writer = stub_->UploadFile(&context, &response);
  1442. // First chunk: metadata
  1443. smartbotic::databasepb::FileChunk metaChunk;
  1444. auto* m = metaChunk.mutable_metadata();
  1445. m->set_name(meta.name);
  1446. m->set_mime_type(meta.mime_type);
  1447. m->set_file_type(meta.file_type);
  1448. m->set_related_id(meta.related_id);
  1449. m->set_is_public(meta.is_public);
  1450. // v2.6.0 — files are namespaced like collections. Filled from
  1451. // Config::project so callers keep passing bare ids and names.
  1452. m->set_project(config_.project);
  1453. // v2.8.0 — relative TTL; the server converts it to an absolute expiry.
  1454. // Set unconditionally on this overload: calling it IS the explicit
  1455. // choice, and 0 therefore means "never expire, ignore any configured
  1456. // default for this file type". The two-argument uploadFile() leaves the
  1457. // field absent so the default applies.
  1458. if (explicitTtl) m->set_ttl_seconds(ttlSeconds);
  1459. for (const auto& [key, value] : meta.metadata) {
  1460. (*m->mutable_metadata())[key] = value;
  1461. }
  1462. writer->Write(metaChunk);
  1463. // Data chunks (64KB each)
  1464. constexpr size_t CHUNK_SIZE = 64 * 1024;
  1465. for (size_t offset = 0; offset < data.size(); offset += CHUNK_SIZE) {
  1466. smartbotic::databasepb::FileChunk dataChunk;
  1467. size_t chunkSize = std::min(CHUNK_SIZE, data.size() - offset);
  1468. dataChunk.set_data(data.data() + offset, chunkSize);
  1469. if (!writer->Write(dataChunk)) break;
  1470. }
  1471. writer->WritesDone();
  1472. auto status = writer->Finish();
  1473. if (!status.ok()) {
  1474. spdlog::error("Client::uploadFile failed: {}", status.error_message());
  1475. throw std::runtime_error(status.error_message());
  1476. }
  1477. return {response.id(), response.size(), response.checksum(), response.deduplicated()};
  1478. }
  1479. std::vector<uint8_t> downloadFile(const std::string& id) {
  1480. // v2.6.0 — see note in uploadFile; project scopes every file call.
  1481. smartbotic::databasepb::DownloadFileRequest request;
  1482. request.set_id(id);
  1483. request.set_project(config_.project);
  1484. grpc::ClientContext context;
  1485. context.set_deadline(
  1486. std::chrono::system_clock::now() + std::chrono::milliseconds(config_.timeoutMs * 10));
  1487. attachAuth(context);
  1488. auto reader = stub_->DownloadFile(&context, request);
  1489. std::vector<uint8_t> fileData;
  1490. smartbotic::databasepb::FileChunk chunk;
  1491. while (reader->Read(&chunk)) {
  1492. if (chunk.has_data()) {
  1493. const auto& d = chunk.data();
  1494. fileData.insert(fileData.end(), d.begin(), d.end());
  1495. }
  1496. }
  1497. auto status = reader->Finish();
  1498. if (!status.ok()) {
  1499. spdlog::error("Client::downloadFile failed: {}", status.error_message());
  1500. throw std::runtime_error(status.error_message());
  1501. }
  1502. return fileData;
  1503. }
  1504. std::optional<Client::FileRecord> getFileInfo(const std::string& id) {
  1505. smartbotic::databasepb::GetFileInfoRequest request;
  1506. request.set_id(id);
  1507. request.set_project(config_.project);
  1508. smartbotic::databasepb::FileInfo response;
  1509. grpc::ClientContext context;
  1510. setDeadline(context);
  1511. auto status = stub_->GetFileInfo(&context, request, &response);
  1512. if (!status.ok()) {
  1513. if (status.error_code() == grpc::StatusCode::NOT_FOUND) return std::nullopt;
  1514. spdlog::error("Client::getFileInfo failed: {}", status.error_message());
  1515. throw std::runtime_error(status.error_message());
  1516. }
  1517. Client::FileRecord record;
  1518. record.id = response.id();
  1519. record.name = response.name();
  1520. record.mime_type = response.mime_type();
  1521. record.size = response.size();
  1522. record.file_type = response.file_type();
  1523. record.related_id = response.related_id();
  1524. record.checksum = response.checksum();
  1525. record.is_public = response.is_public();
  1526. record.ref_count = response.ref_count();
  1527. record.created_at = response.created_at();
  1528. record.project = response.project();
  1529. for (const auto& [key, value] : response.metadata()) {
  1530. record.metadata[key] = value;
  1531. }
  1532. return record;
  1533. }
  1534. bool deleteFile(const std::string& id) {
  1535. smartbotic::databasepb::DeleteFileRequest request;
  1536. request.set_id(id);
  1537. request.set_project(config_.project);
  1538. smartbotic::databasepb::DeleteFileResponse response;
  1539. grpc::ClientContext context;
  1540. setDeadline(context);
  1541. auto status = stub_->DeleteFile(&context, request, &response);
  1542. if (!status.ok()) {
  1543. spdlog::error("Client::deleteFile failed: {}", status.error_message());
  1544. return false;
  1545. }
  1546. return response.deleted();
  1547. }
  1548. Client::FileListResult listFiles(const std::string& file_type,
  1549. const std::string& related_id,
  1550. uint32_t limit, uint32_t offset,
  1551. const std::string& checksum,
  1552. const std::string& name) {
  1553. smartbotic::databasepb::ListFilesRequest request;
  1554. if (!file_type.empty()) request.set_file_type(file_type);
  1555. if (!related_id.empty()) request.set_related_id(related_id);
  1556. request.set_limit(limit);
  1557. request.set_offset(offset);
  1558. if (!checksum.empty()) request.set_checksum(checksum);
  1559. if (!name.empty()) request.set_name(name);
  1560. request.set_project(config_.project);
  1561. smartbotic::databasepb::ListFilesResponse response;
  1562. grpc::ClientContext context;
  1563. setDeadline(context);
  1564. auto status = stub_->ListFiles(&context, request, &response);
  1565. if (!status.ok()) {
  1566. spdlog::error("Client::listFiles failed: {}", status.error_message());
  1567. throw std::runtime_error(status.error_message());
  1568. }
  1569. Client::FileListResult result;
  1570. result.total_count = response.total_count();
  1571. result.has_more = response.has_more();
  1572. for (const auto& f : response.files()) {
  1573. Client::FileRecord record;
  1574. record.id = f.id();
  1575. record.name = f.name();
  1576. record.mime_type = f.mime_type();
  1577. record.size = f.size();
  1578. record.file_type = f.file_type();
  1579. record.related_id = f.related_id();
  1580. record.checksum = f.checksum();
  1581. record.is_public = f.is_public();
  1582. record.created_at = f.created_at();
  1583. record.project = f.project();
  1584. for (const auto& [key, value] : f.metadata()) {
  1585. record.metadata[key] = value;
  1586. }
  1587. result.files.push_back(std::move(record));
  1588. }
  1589. return result;
  1590. }
  1591. private:
  1592. // Per-context setup: deadline + (v2.4) bearer-token auth metadata.
  1593. //
  1594. // Most RPC call sites in this library funnel through this helper —
  1595. // unary (insert, get, update, ...) plus the unary file ops that need
  1596. // the default deadline. Long-running streams (Subscribe) and the
  1597. // file-streaming RPCs (UploadFile/DownloadFile) compute their own
  1598. // deadlines and call `attachAuth()` separately so the auth metadata
  1599. // still lands on every outbound RPC. Centralising both concerns here
  1600. // means we don't need a gRPC ClientInterceptorFactoryInterface — the
  1601. // interceptor route is also offered by gRPC++ but it sits in
  1602. // grpc::experimental:: and would need to be re-validated on every
  1603. // grpc minor; the per-context style keeps Stage E independent of
  1604. // gRPC's TLS/interceptor experimental flux.
  1605. void setDeadline(grpc::ClientContext& context) {
  1606. context.set_deadline(
  1607. std::chrono::system_clock::now() + std::chrono::milliseconds(config_.timeoutMs)
  1608. );
  1609. attachAuth(context);
  1610. }
  1611. // v2.4 — attach the bearer-token authorization metadata if configured.
  1612. // Use this on contexts that set their own deadline (streams,
  1613. // long-running ops) so the auth metadata still lands.
  1614. void attachAuth(grpc::ClientContext& context) {
  1615. if (!config_.auth_token.empty()) {
  1616. context.AddMetadata("authorization", "Bearer " + config_.auth_token);
  1617. }
  1618. }
  1619. Config config_;
  1620. std::shared_ptr<grpc::Channel> channel_;
  1621. std::unique_ptr<smartbotic::databasepb::DatabaseService::Stub> stub_;
  1622. std::atomic<bool> connected_{false};
  1623. };
  1624. // ===== Client Public Interface Implementation =====
  1625. Client::Client(Config config)
  1626. : impl_(std::make_unique<Impl>(std::move(config))) {}
  1627. Client::~Client() = default;
  1628. Client::Client(Client&&) noexcept = default;
  1629. Client& Client::operator=(Client&&) noexcept = default;
  1630. bool Client::connect() {
  1631. return impl_->connect();
  1632. }
  1633. bool Client::isConnected() const {
  1634. return impl_->isConnected();
  1635. }
  1636. std::string Client::insert(const std::string& collection, const nlohmann::json& data,
  1637. const std::string& id, uint32_t ttlSeconds,
  1638. const std::string& actor) {
  1639. return impl_->insert(collection, data, id, ttlSeconds, actor);
  1640. }
  1641. std::optional<nlohmann::json> Client::get(const std::string& collection, const std::string& id) {
  1642. return impl_->get(collection, id);
  1643. }
  1644. bool Client::update(const std::string& collection, const std::string& id, const nlohmann::json& data,
  1645. const std::string& actor) {
  1646. return impl_->update(collection, id, data, actor);
  1647. }
  1648. bool Client::updateIfVersion(const std::string& collection, const std::string& id,
  1649. const nlohmann::json& data, uint64_t expectedVersion,
  1650. const std::string& actor) {
  1651. return impl_->updateIfVersion(collection, id, data, expectedVersion, actor);
  1652. }
  1653. uint64_t Client::patch(const std::string& collection, const std::string& id,
  1654. const nlohmann::json& fields, const std::string& actor) {
  1655. return impl_->patch(collection, id, fields, actor);
  1656. }
  1657. std::pair<std::string, bool> Client::upsert(const std::string& collection, const nlohmann::json& data,
  1658. const std::string& id, uint32_t ttlSeconds,
  1659. const std::string& actor) {
  1660. return impl_->upsert(collection, data, id, ttlSeconds, actor);
  1661. }
  1662. bool Client::remove(const std::string& collection, const std::string& id) {
  1663. return impl_->remove(collection, id);
  1664. }
  1665. bool Client::exists(const std::string& collection, const std::string& id) {
  1666. return impl_->exists(collection, id);
  1667. }
  1668. Client::VersionHistoryResult Client::getVersionHistory(const std::string& collection,
  1669. const std::string& id, uint32_t limit, uint32_t offset) {
  1670. return impl_->getVersionHistory(collection, id, limit, offset);
  1671. }
  1672. std::optional<Client::VersionEntry> Client::getDocumentVersion(const std::string& collection,
  1673. const std::string& id, uint64_t version) {
  1674. return impl_->getDocumentVersion(collection, id, version);
  1675. }
  1676. uint64_t Client::restoreVersion(const std::string& collection, const std::string& id,
  1677. uint64_t version, const std::string& actor) {
  1678. return impl_->restoreVersion(collection, id, version, actor);
  1679. }
  1680. std::pair<uint64_t, uint64_t> Client::restoreToDate(const std::string& collection,
  1681. const std::string& id, uint64_t timestamp, const std::string& actor) {
  1682. return impl_->restoreToDate(collection, id, timestamp, actor);
  1683. }
  1684. std::vector<nlohmann::json> Client::find(const std::string& collection,
  1685. const QueryOptions& options) {
  1686. return impl_->find(collection, options);
  1687. }
  1688. std::vector<nlohmann::json> Client::find(const std::string& collection,
  1689. const QueryOptions& options,
  1690. const std::vector<std::string>& projection) {
  1691. return impl_->find(collection, options, projection);
  1692. }
  1693. Client::FindResult Client::findWithMetrics(const std::string& collection,
  1694. const QueryOptions& options) {
  1695. return impl_->findWithMetrics(collection, options);
  1696. }
  1697. uint64_t Client::count(const std::string& collection,
  1698. const std::vector<Filter>& filters) {
  1699. return impl_->count(collection, filters);
  1700. }
  1701. uint64_t Client::count(const std::string& collection) {
  1702. return impl_->count(collection, {});
  1703. }
  1704. bool Client::setAdd(const std::string& collection, const std::string& setId, const std::string& member) {
  1705. return impl_->setAdd(collection, setId, member);
  1706. }
  1707. bool Client::setRemove(const std::string& collection, const std::string& setId, const std::string& member) {
  1708. return impl_->setRemove(collection, setId, member);
  1709. }
  1710. std::vector<std::string> Client::setMembers(const std::string& collection, const std::string& setId) {
  1711. return impl_->setMembers(collection, setId);
  1712. }
  1713. bool Client::setIsMember(const std::string& collection, const std::string& setId, const std::string& member) {
  1714. return impl_->setIsMember(collection, setId, member);
  1715. }
  1716. std::vector<Client::SimilarityResult> Client::similaritySearch(
  1717. const std::string& collection, const std::vector<float>& queryVector,
  1718. uint32_t topK, float minScore) {
  1719. return impl_->similaritySearch(collection, queryVector, topK, minScore);
  1720. }
  1721. bool Client::createCollection(const std::string& name, uint32_t defaultTtlSeconds,
  1722. bool encrypted, uint32_t maxVersions, uint32_t vectorDimension) {
  1723. return impl_->createCollection(name, defaultTtlSeconds, encrypted, maxVersions, vectorDimension);
  1724. }
  1725. bool Client::dropCollection(const std::string& name) {
  1726. return impl_->dropCollection(name);
  1727. }
  1728. std::vector<std::string> Client::listCollections() {
  1729. return impl_->listCollections();
  1730. }
  1731. std::optional<Client::CollectionInfo> Client::getCollectionInfo(const std::string& name) {
  1732. return impl_->getCollectionInfo(name);
  1733. }
  1734. // ===== Project Management (v2.3 Stage F) =====
  1735. std::vector<std::string> Client::listProjects() {
  1736. return impl_->listProjects();
  1737. }
  1738. bool Client::createProject(const std::string& name) {
  1739. return impl_->createProject(name);
  1740. }
  1741. void Client::dropProject(const std::string& name) {
  1742. impl_->dropProject(name);
  1743. }
  1744. bool Client::configureCollection(const std::string& collection, const CollectionConfig& cfg) {
  1745. return impl_->configureCollection(collection, cfg);
  1746. }
  1747. Client::CollectionConfig Client::getCollectionConfig(const std::string& collection) {
  1748. return impl_->getCollectionConfig(collection);
  1749. }
  1750. bool Client::hasCollectionConfig(const std::string& collection) {
  1751. return impl_->hasCollectionConfig(collection);
  1752. }
  1753. Client::TimestampMigrationResult Client::migrateCollectionTimestamps(
  1754. const std::string& collection,
  1755. const std::string& fromPrecision,
  1756. const std::string& toPrecision
  1757. ) {
  1758. return impl_->migrateCollectionTimestamps(collection, fromPrecision, toPrecision);
  1759. }
  1760. bool Client::createView(const std::string& name,
  1761. const std::string& collection,
  1762. const std::vector<std::string>& include,
  1763. const std::vector<std::string>& exclude,
  1764. const std::vector<Filter>& where,
  1765. const std::optional<Sort>& defaultSort) {
  1766. return impl_->createView(name, collection, include, exclude, where, defaultSort);
  1767. }
  1768. bool Client::dropView(const std::string& name) {
  1769. return impl_->dropView(name);
  1770. }
  1771. std::vector<Client::ViewDefinition> Client::listViews() {
  1772. return impl_->listViews();
  1773. }
  1774. std::optional<Client::ViewDefinition> Client::getViewInfo(const std::string& name) {
  1775. return impl_->getViewInfo(name);
  1776. }
  1777. std::shared_ptr<void> Client::subscribe(const std::vector<std::string>& collections, EventCallback callback) {
  1778. return impl_->subscribe(collections, std::move(callback));
  1779. }
  1780. bool Client::healthCheck() {
  1781. return impl_->healthCheck();
  1782. }
  1783. std::optional<Client::HealthInfo> Client::getHealthInfo() {
  1784. return impl_->getHealthInfo();
  1785. }
  1786. std::optional<Client::StatsInfo> Client::getStats() {
  1787. return impl_->getStats();
  1788. }
  1789. Client::MemoryStats Client::getMemoryStats() {
  1790. return impl_->getMemoryStats();
  1791. }
  1792. bool Client::setReadOnly(bool readOnly) {
  1793. return impl_->setReadOnly(readOnly);
  1794. }
  1795. Client::ReadOnlyStatus Client::getReadOnlyStatus() {
  1796. return impl_->getReadOnlyStatus();
  1797. }
  1798. Client::FileUploadResult Client::uploadFile(const std::vector<uint8_t>& data, const FileUploadMeta& meta) {
  1799. return impl_->uploadFile(data, meta);
  1800. }
  1801. Client::FileUploadResult Client::uploadFile(const std::vector<uint8_t>& data,
  1802. const FileUploadMeta& meta,
  1803. uint32_t ttlSeconds) {
  1804. return impl_->uploadFile(data, meta, ttlSeconds, /*explicitTtl=*/true);
  1805. }
  1806. std::vector<uint8_t> Client::downloadFile(const std::string& id) {
  1807. return impl_->downloadFile(id);
  1808. }
  1809. std::optional<Client::FileRecord> Client::getFileInfo(const std::string& id) {
  1810. return impl_->getFileInfo(id);
  1811. }
  1812. bool Client::deleteFile(const std::string& id) {
  1813. return impl_->deleteFile(id);
  1814. }
  1815. Client::FileListResult Client::listFiles(const std::string& file_type,
  1816. const std::string& related_id,
  1817. uint32_t limit, uint32_t offset,
  1818. const std::string& checksum,
  1819. const std::string& name) {
  1820. return impl_->listFiles(file_type, related_id, limit, offset, checksum, name);
  1821. }
  1822. } // namespace smartbotic::database