#include "smartbotic/database/client.hpp" #include #include #include #include #include #include #include #include namespace smartbotic::database { namespace { // v1.7.0 T11 — transient transport errors that are safe to retry for idempotent // writes (insert with explicit ID, updateIfVersion, upsert, remove, patch). // All other codes (INVALID_ARGUMENT, FAILED_PRECONDITION, NOT_FOUND, ...) are // final — no retry. bool retryableStatus(const grpc::Status& s) { switch (s.error_code()) { case grpc::StatusCode::DEADLINE_EXCEEDED: // server queued behind eviction/WAL case grpc::StatusCode::RESOURCE_EXHAUSTED: // admission control or gRPC concurrency cap case grpc::StatusCode::UNAVAILABLE: // transient connection issue return true; default: return false; } } // Exponential backoff with jitter. attempt=0 => baseMs; attempt=1 => baseMs*2; ... uint32_t computeBackoffMs(uint32_t baseMs, uint32_t attempt, uint32_t capMs, double jitterPct) { uint64_t exp = static_cast(baseMs) << attempt; if (exp > capMs) exp = capMs; // ±jitterPct jitter; std::rand() is fine here — not a security sensitive RNG. double jitter = 1.0 + ((double(std::rand()) / RAND_MAX) * 2.0 - 1.0) * jitterPct; if (jitter < 0.1) jitter = 0.1; return static_cast(exp * jitter); } } // anonymous namespace // ===== PIMPL Implementation ===== class Client::Impl { public: explicit Impl(Config config) : config_(std::move(config)) {} ~Impl() { disconnect(); } bool connect() { try { auto channelArgs = grpc::ChannelArguments(); // Use longer keepalive intervals to avoid "too_many_pings" errors from server channelArgs.SetInt(GRPC_ARG_KEEPALIVE_TIME_MS, 60000); // 60 seconds channelArgs.SetInt(GRPC_ARG_KEEPALIVE_TIMEOUT_MS, 20000); // 20 seconds channelArgs.SetInt(GRPC_ARG_KEEPALIVE_PERMIT_WITHOUT_CALLS, 0); // Only ping when there are active calls // Set max message size to match server (100MB for file uploads) channelArgs.SetMaxReceiveMessageSize(100 * 1024 * 1024); channelArgs.SetMaxSendMessageSize(100 * 1024 * 1024); channel_ = grpc::CreateCustomChannel( config_.address, grpc::InsecureChannelCredentials(), channelArgs ); stub_ = smartbotic::databasepb::DatabaseService::NewStub(channel_); connected_ = true; spdlog::info("Database client connected to {}", config_.address); return true; } catch (const std::exception& e) { spdlog::error("Database client connection failed: {}", e.what()); return false; } } void disconnect() { connected_ = false; stub_.reset(); channel_.reset(); } bool isConnected() const { if (!connected_ || !channel_) { return false; } // Check if channel is in a usable state (not failed or shutdown) auto state = channel_->GetState(false); return state == GRPC_CHANNEL_READY || state == GRPC_CHANNEL_IDLE || state == GRPC_CHANNEL_CONNECTING; } // ===== Document Operations ===== std::string insert(const std::string& collection, const nlohmann::json& data, const std::string& id, uint32_t ttlSeconds, const std::string& actor) { smartbotic::databasepb::InsertRequest request; request.set_collection(collection); request.set_data(data.dump()); if (!id.empty()) { request.set_id(id); } if (ttlSeconds > 0) { request.set_ttl_seconds(ttlSeconds); } if (!actor.empty()) { request.set_actor(actor); } smartbotic::databasepb::InsertResponse response; grpc::Status status; for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) { grpc::ClientContext context; setDeadline(context); response.Clear(); status = stub_->Insert(&context, request, &response); if (status.ok() || !retryableStatus(status)) { break; } if (attempt < config_.writeRetries) { uint32_t backoffMs = computeBackoffMs( config_.writeRetryBackoffMs, attempt, config_.writeRetryMaxBackoffMs, config_.writeRetryJitter); spdlog::warn("Client::insert {}; retrying in {}ms (attempt {}/{})", status.error_message(), backoffMs, attempt + 1, config_.writeRetries); std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs)); } } if (!status.ok()) { spdlog::error("Client::insert failed after retries: {}", status.error_message()); throw std::runtime_error(status.error_message()); } return response.id(); } std::optional get(const std::string& collection, const std::string& id) { smartbotic::databasepb::GetRequest request; request.set_collection(collection); request.set_id(id); smartbotic::databasepb::GetResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->Get(&context, request, &response); if (!status.ok()) { spdlog::error("Client::get failed: {}", status.error_message()); return std::nullopt; } if (!response.found()) { return std::nullopt; } auto json = nlohmann::json::parse(response.document().data()); json["_id"] = response.document().id(); json["_version"] = response.document().version(); json["_created_at"] = response.document().created_at(); json["_updated_at"] = response.document().updated_at(); json["_created_by"] = response.document().created_by(); json["_updated_by"] = response.document().updated_by(); return json; } bool update(const std::string& collection, const std::string& id, const nlohmann::json& data, const std::string& actor) { // Strip metadata fields from data — callers pass the JSON from get() which // includes _version, _id, etc. We use _version for optimistic locking. auto cleanData = data; cleanData.erase("_id"); cleanData.erase("_version"); cleanData.erase("_created_at"); cleanData.erase("_updated_at"); cleanData.erase("_created_by"); cleanData.erase("_updated_by"); // Optimistic locking with automatic retry: // 1. Read current version // 2. Call updateIfVersion // 3. On version conflict (another writer), re-read and retry for (uint32_t attempt = 0; attempt <= config_.maxRetries; ++attempt) { // Get current version auto current = get(collection, id); if (!current) { return false; // document doesn't exist } uint64_t version = (*current)["_version"].get(); // Attempt version-checked update if (updateIfVersion(collection, id, cleanData, version, actor)) { return true; } // Version conflict — another client wrote between our get() and update if (attempt < config_.maxRetries) { spdlog::debug("Client::update version conflict on {}/{}, retry {}/{}", collection, id, attempt + 1, config_.maxRetries); } } spdlog::warn("Client::update failed after {} retries due to version conflicts on {}/{}", config_.maxRetries, collection, id); return false; } bool updateIfVersion(const std::string& collection, const std::string& id, const nlohmann::json& data, uint64_t expectedVersion, const std::string& actor) { smartbotic::databasepb::UpdateRequest request; request.set_collection(collection); request.set_id(id); request.set_data(data.dump()); request.set_expected_version(expectedVersion); if (!actor.empty()) { request.set_actor(actor); } smartbotic::databasepb::UpdateResponse response; grpc::Status status; for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) { grpc::ClientContext context; setDeadline(context); response.Clear(); status = stub_->Update(&context, request, &response); if (status.ok() || !retryableStatus(status)) { break; } if (attempt < config_.writeRetries) { uint32_t backoffMs = computeBackoffMs( config_.writeRetryBackoffMs, attempt, config_.writeRetryMaxBackoffMs, config_.writeRetryJitter); spdlog::warn("Client::updateIfVersion {}; retrying in {}ms (attempt {}/{})", status.error_message(), backoffMs, attempt + 1, config_.writeRetries); std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs)); } } if (!status.ok()) { spdlog::error("Client::updateIfVersion failed after retries: {}", status.error_message()); return false; } return response.success(); } uint64_t patch(const std::string& collection, const std::string& id, const nlohmann::json& fields, const std::string& actor) { smartbotic::databasepb::PatchDocumentRequest request; request.set_collection(collection); request.set_id(id); request.set_patch_json(fields.dump()); if (!actor.empty()) { request.set_actor(actor); } smartbotic::databasepb::PatchDocumentResponse response; grpc::Status status; for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) { grpc::ClientContext context; setDeadline(context); response.Clear(); status = stub_->PatchDocument(&context, request, &response); if (status.ok() || !retryableStatus(status)) { break; } if (attempt < config_.writeRetries) { uint32_t backoffMs = computeBackoffMs( config_.writeRetryBackoffMs, attempt, config_.writeRetryMaxBackoffMs, config_.writeRetryJitter); spdlog::warn("Client::patch {}; retrying in {}ms (attempt {}/{})", status.error_message(), backoffMs, attempt + 1, config_.writeRetries); std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs)); } } if (!status.ok()) { spdlog::error("Client::patch failed after retries: {}", status.error_message()); return 0; } if (!response.success()) { spdlog::error("Client::patch failed: {}", response.error()); return 0; } return response.new_version(); } std::pair upsert(const std::string& collection, const nlohmann::json& data, const std::string& id, uint32_t ttlSeconds, const std::string& actor) { smartbotic::databasepb::UpsertRequest request; request.set_collection(collection); request.set_data(data.dump()); if (!id.empty()) { request.set_id(id); } if (ttlSeconds > 0) { request.set_ttl_seconds(ttlSeconds); } if (!actor.empty()) { request.set_actor(actor); } smartbotic::databasepb::UpsertResponse response; grpc::Status status; for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) { grpc::ClientContext context; setDeadline(context); response.Clear(); status = stub_->Upsert(&context, request, &response); if (status.ok() || !retryableStatus(status)) { break; } if (attempt < config_.writeRetries) { uint32_t backoffMs = computeBackoffMs( config_.writeRetryBackoffMs, attempt, config_.writeRetryMaxBackoffMs, config_.writeRetryJitter); spdlog::warn("Client::upsert {}; retrying in {}ms (attempt {}/{})", status.error_message(), backoffMs, attempt + 1, config_.writeRetries); std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs)); } } if (!status.ok()) { spdlog::error("Client::upsert failed after retries: {}", status.error_message()); throw std::runtime_error(status.error_message()); } return {response.id(), response.inserted()}; } bool remove(const std::string& collection, const std::string& id) { smartbotic::databasepb::DeleteRequest request; request.set_collection(collection); request.set_id(id); smartbotic::databasepb::DeleteResponse response; grpc::Status status; for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) { grpc::ClientContext context; setDeadline(context); response.Clear(); status = stub_->Delete(&context, request, &response); if (status.ok() || !retryableStatus(status)) { break; } if (attempt < config_.writeRetries) { uint32_t backoffMs = computeBackoffMs( config_.writeRetryBackoffMs, attempt, config_.writeRetryMaxBackoffMs, config_.writeRetryJitter); spdlog::warn("Client::remove {}; retrying in {}ms (attempt {}/{})", status.error_message(), backoffMs, attempt + 1, config_.writeRetries); std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs)); } } if (!status.ok()) { spdlog::error("Client::remove failed after retries: {}", status.error_message()); return false; } return response.deleted(); } bool exists(const std::string& collection, const std::string& id) { smartbotic::databasepb::ExistsRequest request; request.set_collection(collection); request.set_id(id); smartbotic::databasepb::ExistsResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->Exists(&context, request, &response); if (!status.ok()) { spdlog::error("Client::exists failed: {}", status.error_message()); return false; } return response.exists(); } // ===== Query Operations ===== std::vector find(const std::string& collection, const Client::QueryOptions& options) { smartbotic::databasepb::FindRequest request; request.set_collection(collection); // Set filters for (const auto& filter : options.filters) { auto* pb = request.add_filters(); pb->set_field(filter.field); pb->set_value(filter.value.dump()); pb->set_op(static_cast(filter.op)); } // Set sorting if (!options.sortField.empty()) { auto* sort = request.mutable_sort(); sort->set_field(options.sortField); sort->set_descending(options.sortDescending); } // Set pagination request.set_limit(options.limit); request.set_offset(options.offset); smartbotic::databasepb::FindResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->Find(&context, request, &response); if (!status.ok()) { spdlog::error("Client::find failed: {}", status.error_message()); return {}; } std::vector results; results.reserve(response.documents_size()); for (const auto& doc : response.documents()) { auto json = nlohmann::json::parse(doc.data()); json["_id"] = doc.id(); json["_version"] = doc.version(); json["_created_at"] = doc.created_at(); json["_updated_at"] = doc.updated_at(); json["_created_by"] = doc.created_by(); json["_updated_by"] = doc.updated_by(); results.push_back(json); } return results; } Client::FindResult findWithMetrics(const std::string& collection, const Client::QueryOptions& options) { smartbotic::databasepb::FindRequest request; request.set_collection(collection); // Set filters for (const auto& filter : options.filters) { auto* pb = request.add_filters(); pb->set_field(filter.field); pb->set_value(filter.value.dump()); pb->set_op(static_cast(filter.op)); } // Set sorting if (!options.sortField.empty()) { auto* sort = request.mutable_sort(); sort->set_field(options.sortField); sort->set_descending(options.sortDescending); } // Set pagination request.set_limit(options.limit); request.set_offset(options.offset); smartbotic::databasepb::FindResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->Find(&context, request, &response); if (!status.ok()) { spdlog::error("Client::findWithMetrics failed: {}", status.error_message()); return {}; } Client::FindResult result; result.totalCount = response.total_count(); result.hasMore = response.has_more(); result.usedWalFallback = response.used_wal_fallback(); result.memoryMatchCount = response.memory_match_count(); result.walMatchCount = response.wal_match_count(); result.memorySearchMicros = response.memory_search_micros(); result.walSearchMicros = response.wal_search_micros(); result.documents.reserve(response.documents_size()); for (const auto& doc : response.documents()) { auto json = nlohmann::json::parse(doc.data()); json["_id"] = doc.id(); json["_version"] = doc.version(); json["_created_at"] = doc.created_at(); json["_updated_at"] = doc.updated_at(); json["_created_by"] = doc.created_by(); json["_updated_by"] = doc.updated_by(); result.documents.push_back(json); } return result; } uint64_t count(const std::string& collection, const std::vector& filters) { smartbotic::databasepb::CountRequest request; request.set_collection(collection); for (const auto& filter : filters) { auto* pb = request.add_filters(); pb->set_field(filter.field); pb->set_value(filter.value.dump()); pb->set_op(static_cast(filter.op)); } smartbotic::databasepb::CountResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->Count(&context, request, &response); if (!status.ok()) { spdlog::error("Client::count failed: {}", status.error_message()); return 0; } return response.count(); } // ===== Set Operations ===== bool setAdd(const std::string& collection, const std::string& setId, const std::string& member) { smartbotic::databasepb::SetAddRequest request; request.set_collection(collection); request.set_set_id(setId); request.set_member(member); smartbotic::databasepb::SetAddResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->SetAdd(&context, request, &response); if (!status.ok()) { spdlog::error("Client::setAdd failed: {}", status.error_message()); return false; } return response.added(); } bool setRemove(const std::string& collection, const std::string& setId, const std::string& member) { smartbotic::databasepb::SetRemoveRequest request; request.set_collection(collection); request.set_set_id(setId); request.set_member(member); smartbotic::databasepb::SetRemoveResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->SetRemove(&context, request, &response); if (!status.ok()) { spdlog::error("Client::setRemove failed: {}", status.error_message()); return false; } return response.removed(); } std::vector setMembers(const std::string& collection, const std::string& setId) { smartbotic::databasepb::SetMembersRequest request; request.set_collection(collection); request.set_set_id(setId); smartbotic::databasepb::SetMembersResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->SetMembers(&context, request, &response); if (!status.ok()) { spdlog::error("Client::setMembers failed: {}", status.error_message()); return {}; } return {response.members().begin(), response.members().end()}; } bool setIsMember(const std::string& collection, const std::string& setId, const std::string& member) { smartbotic::databasepb::SetIsMemberRequest request; request.set_collection(collection); request.set_set_id(setId); request.set_member(member); smartbotic::databasepb::SetIsMemberResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->SetIsMember(&context, request, &response); if (!status.ok()) { spdlog::error("Client::setIsMember failed: {}", status.error_message()); return false; } return response.is_member(); } // ===== Collection Management ===== std::vector similaritySearch( const std::string& collection, const std::vector& queryVector, uint32_t topK, float minScore) { smartbotic::databasepb::SimilaritySearchRequest request; request.set_collection(collection); for (float v : queryVector) { request.add_query_vector(v); } request.set_top_k(topK); request.set_min_score(minScore); smartbotic::databasepb::SimilaritySearchResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->SimilaritySearch(&context, request, &response); if (!status.ok()) { spdlog::error("Client::similaritySearch failed: {}", status.error_message()); throw std::runtime_error(status.error_message()); } std::vector results; results.reserve(response.results_size()); for (const auto& r : response.results()) { Client::SimilarityResult entry; entry.id = r.id(); entry.score = r.score(); if (!r.data().empty()) { entry.data = nlohmann::json::parse(r.data()); } results.push_back(std::move(entry)); } return results; } bool createCollection(const std::string& name, uint32_t defaultTtlSeconds, bool encrypted, uint32_t maxVersions, uint32_t vectorDimension) { smartbotic::databasepb::CreateCollectionRequest request; request.set_name(name); auto* options = request.mutable_options(); if (defaultTtlSeconds > 0) { options->set_default_ttl_seconds(defaultTtlSeconds); } options->set_encrypted(encrypted); if (maxVersions > 0) { options->set_max_versions(maxVersions); } if (vectorDimension > 0) { options->set_vector_dimension(vectorDimension); } smartbotic::databasepb::CreateCollectionResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->CreateCollection(&context, request, &response); if (!status.ok()) { spdlog::error("Client::createCollection failed: {}", status.error_message()); return false; } return response.created(); } bool dropCollection(const std::string& name) { smartbotic::databasepb::DropCollectionRequest request; request.set_name(name); smartbotic::databasepb::DropCollectionResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->DropCollection(&context, request, &response); if (!status.ok()) { spdlog::error("Client::dropCollection failed: {}", status.error_message()); return false; } return response.dropped(); } std::vector listCollections() { smartbotic::databasepb::ListCollectionsRequest request; smartbotic::databasepb::ListCollectionsResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->ListCollections(&context, request, &response); if (!status.ok()) { spdlog::error("Client::listCollections failed: {}", status.error_message()); return {}; } return {response.names().begin(), response.names().end()}; } std::optional getCollectionInfo(const std::string& name) { smartbotic::databasepb::GetCollectionInfoRequest request; request.set_name(name); smartbotic::databasepb::GetCollectionInfoResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->GetCollectionInfo(&context, request, &response); if (!status.ok()) { spdlog::error("Client::getCollectionInfo failed: {}", status.error_message()); return std::nullopt; } if (!response.found()) { return std::nullopt; } const auto& info = response.info(); Client::CollectionInfo result; result.name = info.name(); result.documentCount = info.document_count(); result.sizeBytes = info.size_bytes(); result.defaultTtlSeconds = info.options().default_ttl_seconds(); result.encrypted = info.options().encrypted(); result.maxVersions = info.options().max_versions(); result.createdAt = info.created_at(); result.updatedAt = info.updated_at(); return result; } // ===== Collection Configuration ===== bool configureCollection(const std::string& collection, const Client::CollectionConfig& cfg) { smartbotic::databasepb::ConfigureCollectionRequest request; request.set_collection(collection); request.mutable_config()->set_timestamp_precision(cfg.timestampPrecision); smartbotic::databasepb::ConfigureCollectionResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->ConfigureCollection(&context, request, &response); if (!status.ok()) { spdlog::error("Client::configureCollection failed: {}", status.error_message()); return false; } if (!response.success()) { spdlog::error("Client::configureCollection rejected: {}", response.error()); return false; } return true; } Client::CollectionConfig getCollectionConfig(const std::string& collection) { smartbotic::databasepb::GetCollectionConfigRequest request; request.set_collection(collection); smartbotic::databasepb::GetCollectionConfigResponse response; grpc::ClientContext context; setDeadline(context); Client::CollectionConfig out; auto status = stub_->GetCollectionConfig(&context, request, &response); if (!status.ok()) { spdlog::error("Client::getCollectionConfig failed: {}", status.error_message()); return out; } out.timestampPrecision = response.config().timestamp_precision(); if (out.timestampPrecision.empty()) out.timestampPrecision = "ms"; return out; } bool hasCollectionConfig(const std::string& collection) { smartbotic::databasepb::GetCollectionConfigRequest request; request.set_collection(collection); smartbotic::databasepb::GetCollectionConfigResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->GetCollectionConfig(&context, request, &response); if (!status.ok()) return false; return response.found(); } Client::TimestampMigrationResult migrateCollectionTimestamps( const std::string& collection, const std::string& fromPrecision, const std::string& toPrecision ) { smartbotic::databasepb::MigrateCollectionTimestampsRequest request; request.set_collection(collection); request.set_from_precision(fromPrecision); request.set_to_precision(toPrecision); smartbotic::databasepb::MigrateCollectionTimestampsResponse response; grpc::ClientContext context; // Longer deadline — can iterate many rows auto deadline = std::chrono::system_clock::now() + std::chrono::minutes(10); context.set_deadline(deadline); Client::TimestampMigrationResult out; auto status = stub_->MigrateCollectionTimestamps(&context, request, &response); if (!status.ok()) { out.success = false; out.error = status.error_message(); return out; } out.success = response.success(); out.error = response.error(); out.rowsMigrated = response.rows_migrated(); out.rowsSkipped = response.rows_skipped(); return out; } // ===== View Management ===== bool createView(const std::string& name, const std::string& collection, const std::vector& include, const std::vector& exclude, const std::vector& where, const std::optional& defaultSort) { smartbotic::databasepb::CreateViewRequest request; request.set_name(name); request.set_collection(collection); for (const auto& p : include) request.add_include(p); for (const auto& p : exclude) request.add_exclude(p); for (const auto& f : where) { auto* pb = request.add_where(); pb->set_field(f.field); pb->set_op(static_cast(f.op)); pb->set_value(f.value.dump()); } if (defaultSort) { auto* s = request.mutable_default_sort(); s->set_field(defaultSort->field); s->set_descending(defaultSort->descending); } smartbotic::databasepb::CreateViewResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->CreateView(&context, request, &response); if (!status.ok()) { spdlog::error("Client::createView failed: {}", status.error_message()); return false; } if (!response.success()) { spdlog::error("Client::createView rejected: {}", response.error()); return false; } return true; } bool dropView(const std::string& name) { smartbotic::databasepb::DropViewRequest request; request.set_name(name); smartbotic::databasepb::DropViewResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->DropView(&context, request, &response); if (!status.ok()) { spdlog::error("Client::dropView failed: {}", status.error_message()); return false; } return response.success(); } std::vector listViews() { smartbotic::databasepb::ListViewsRequest request; smartbotic::databasepb::ListViewsResponse response; grpc::ClientContext context; setDeadline(context); std::vector out; auto status = stub_->ListViews(&context, request, &response); if (!status.ok()) { spdlog::error("Client::listViews failed: {}", status.error_message()); return out; } for (const auto& pb : response.views()) { Client::ViewDefinition v; v.name = pb.name(); v.collection = pb.collection(); for (const auto& p : pb.include()) v.include.push_back(p); for (const auto& p : pb.exclude()) v.exclude.push_back(p); for (const auto& pbf : pb.where()) { Client::Filter f; f.field = pbf.field(); f.op = static_cast(pbf.op()); try { f.value = nlohmann::json::parse(pbf.value()); } catch (...) { f.value = pbf.value(); } v.where.push_back(f); } if (pb.has_default_sort()) { Client::Sort s; s.field = pb.default_sort().field(); s.descending = pb.default_sort().descending(); v.defaultSort = s; } v.createdAt = pb.created_at(); v.updatedAt = pb.updated_at(); out.push_back(std::move(v)); } return out; } std::optional getViewInfo(const std::string& name) { smartbotic::databasepb::GetViewInfoRequest request; request.set_name(name); smartbotic::databasepb::GetViewInfoResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->GetViewInfo(&context, request, &response); if (!status.ok() || !response.found()) { return std::nullopt; } Client::ViewDefinition v; v.name = response.view().name(); v.collection = response.view().collection(); for (const auto& p : response.view().include()) v.include.push_back(p); for (const auto& p : response.view().exclude()) v.exclude.push_back(p); for (const auto& pbf : response.view().where()) { Client::Filter f; f.field = pbf.field(); f.op = static_cast(pbf.op()); try { f.value = nlohmann::json::parse(pbf.value()); } catch (...) { f.value = pbf.value(); } v.where.push_back(f); } if (response.view().has_default_sort()) { Client::Sort s; s.field = response.view().default_sort().field(); s.descending = response.view().default_sort().descending(); v.defaultSort = s; } v.createdAt = response.view().created_at(); v.updatedAt = response.view().updated_at(); return v; } // ===== Event Subscription ===== class SubscriptionHandle { public: SubscriptionHandle(std::shared_ptr ctx, std::unique_ptr> reader, std::thread readerThread) : context_(std::move(ctx)) , reader_(std::move(reader)) , readerThread_(std::move(readerThread)) , active_(true) {} ~SubscriptionHandle() { cancel(); } void cancel() { if (active_.exchange(false)) { context_->TryCancel(); if (readerThread_.joinable()) { readerThread_.join(); } } } private: std::shared_ptr context_; std::unique_ptr> reader_; std::thread readerThread_; std::atomic active_; }; std::shared_ptr subscribe(const std::vector& collections, Client::EventCallback callback) { auto context = std::make_shared(); smartbotic::databasepb::SubscribeRequest request; for (const auto& coll : collections) { request.add_collections(coll); } request.set_include_data(true); auto reader = stub_->Subscribe(context.get(), request); // Create reader thread auto readerThread = std::thread([reader = reader.get(), callback = std::move(callback)]() { smartbotic::databasepb::DatabaseEvent event; while (reader->Read(&event)) { std::optional data; if (!event.data().empty()) { data = nlohmann::json::parse(event.data()); } // Convert proto event type to string std::string eventType; switch (event.type()) { case smartbotic::databasepb::EVENT_INSERT: eventType = "insert"; break; case smartbotic::databasepb::EVENT_UPDATE: eventType = "update"; break; case smartbotic::databasepb::EVENT_DELETE: eventType = "delete"; break; case smartbotic::databasepb::EVENT_EXPIRE: eventType = "expire"; break; case smartbotic::databasepb::EVENT_INVALIDATE: eventType = "invalidate"; break; default: eventType = "unknown"; break; } callback( event.collection(), event.document_id(), eventType, data ); } }); return std::make_shared( std::move(context), std::move(reader), std::move(readerThread) ); } // ===== Version History ===== Client::VersionHistoryResult getVersionHistory(const std::string& collection, const std::string& id, uint32_t limit, uint32_t offset) { smartbotic::databasepb::GetVersionHistoryRequest request; request.set_collection(collection); request.set_id(id); if (limit > 0) request.set_limit(limit); if (offset > 0) request.set_offset(offset); smartbotic::databasepb::GetVersionHistoryResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->GetVersionHistory(&context, request, &response); if (!status.ok()) { spdlog::error("Client::getVersionHistory failed: {}", status.error_message()); return {}; } Client::VersionHistoryResult result; result.currentVersion = response.current_version(); result.totalCount = response.total_count(); result.documentDeleted = response.document_deleted(); result.versions.reserve(response.versions_size()); for (const auto& ver : response.versions()) { Client::VersionEntry entry; entry.version = ver.version(); entry.timestamp = ver.timestamp(); entry.updatedBy = ver.updated_by(); if (!ver.data().empty()) { entry.data = nlohmann::json::parse(ver.data()); } result.versions.push_back(std::move(entry)); } return result; } std::optional getDocumentVersion(const std::string& collection, const std::string& id, uint64_t version) { smartbotic::databasepb::GetDocumentVersionRequest request; request.set_collection(collection); request.set_id(id); request.set_version(version); smartbotic::databasepb::GetDocumentVersionResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->GetDocumentVersion(&context, request, &response); if (!status.ok()) { spdlog::error("Client::getDocumentVersion failed: {}", status.error_message()); return std::nullopt; } if (!response.found()) { return std::nullopt; } Client::VersionEntry entry; entry.version = response.version_entry().version(); entry.timestamp = response.version_entry().timestamp(); entry.updatedBy = response.version_entry().updated_by(); if (!response.version_entry().data().empty()) { entry.data = nlohmann::json::parse(response.version_entry().data()); } return entry; } uint64_t restoreVersion(const std::string& collection, const std::string& id, uint64_t version, const std::string& actor) { smartbotic::databasepb::RestoreVersionRequest request; request.set_collection(collection); request.set_id(id); request.set_version(version); if (!actor.empty()) request.set_actor(actor); smartbotic::databasepb::RestoreVersionResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->RestoreVersion(&context, request, &response); if (!status.ok()) { spdlog::error("Client::restoreVersion failed: {}", status.error_message()); return 0; } if (!response.success()) { spdlog::error("Client::restoreVersion: {}", response.error()); return 0; } return response.new_version(); } std::pair restoreToDate(const std::string& collection, const std::string& id, uint64_t timestamp, const std::string& actor) { smartbotic::databasepb::RestoreToDateRequest request; request.set_collection(collection); request.set_id(id); request.set_timestamp(timestamp); if (!actor.empty()) request.set_actor(actor); smartbotic::databasepb::RestoreToDateResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->RestoreToDate(&context, request, &response); if (!status.ok()) { spdlog::error("Client::restoreToDate failed: {}", status.error_message()); return {0, 0}; } if (!response.success()) { spdlog::error("Client::restoreToDate: {}", response.error()); return {0, 0}; } return {response.restored_version(), response.new_version()}; } // ===== Health ===== bool healthCheck() { smartbotic::databasepb::HealthCheckRequest request; smartbotic::databasepb::HealthCheckResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->HealthCheck(&context, request, &response); if (!status.ok()) { return false; } return response.healthy(); } std::optional getHealthInfo() { smartbotic::databasepb::HealthCheckRequest request; smartbotic::databasepb::HealthCheckResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->HealthCheck(&context, request, &response); if (!status.ok()) { return std::nullopt; } Client::HealthInfo info; info.healthy = response.healthy(); info.uptimeMs = response.uptime_seconds() * 1000; info.documentCount = response.document_count(); info.memoryUsedBytes = response.memory_used_bytes(); info.walSizeBytes = response.wal_size_bytes(); return info; } std::optional getStats() { smartbotic::databasepb::GetStatsRequest request; smartbotic::databasepb::GetStatsResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->GetStats(&context, request, &response); if (!status.ok()) { return std::nullopt; } Client::StatsInfo info; info.totalDocuments = response.total_documents(); info.totalCollections = response.total_collections(); info.memoryUsedBytes = response.memory_used_bytes(); info.walSequence = response.wal_sequence(); info.walSizeBytes = response.wal_size_bytes(); info.snapshotCount = response.snapshot_count(); info.lastSnapshotSequence = response.last_snapshot_sequence(); info.insertCount = response.insert_count(); info.updateCount = response.update_count(); info.deleteCount = response.delete_count(); info.queryCount = response.query_count(); // Memory eviction stats info.evictedDocuments = response.evicted_documents(); info.totalEvictions = response.total_evictions(); info.recoveryCount = response.recovery_count(); // Memory configuration info.maxMemoryBytes = response.max_memory_bytes(); info.evictionThresholdPercent = response.eviction_threshold_percent(); info.evictionTargetPercent = response.eviction_target_percent(); // Operation timing (microseconds) - for performance monitoring info.getCount = response.get_count(); info.getTotalMicros = response.get_total_micros(); info.getMaxMicros = response.get_max_micros(); info.insertTotalMicros = response.insert_total_micros(); info.insertMaxMicros = response.insert_max_micros(); info.updateTotalMicros = response.update_total_micros(); info.updateMaxMicros = response.update_max_micros(); info.queryTotalMicros = response.query_total_micros(); info.queryMaxMicros = response.query_max_micros(); return info; } Client::MemoryStats getMemoryStats() { smartbotic::databasepb::GetMemoryStatsRequest request; smartbotic::databasepb::GetMemoryStatsResponse response; grpc::ClientContext context; setDeadline(context); Client::MemoryStats out; auto status = stub_->GetMemoryStats(&context, request, &response); if (!status.ok()) { spdlog::error("Client::getMemoryStats failed: {}", status.error_message()); return out; } out.totalMemoryBytes = response.total_memory_bytes(); out.maxMemoryBytes = response.max_memory_bytes(); out.pressurePercent = response.pressure_percent(); out.pressureLevel = response.pressure_level(); out.lastEvictionTimestamp = response.last_eviction_timestamp(); out.lastEvictionDocs = response.last_eviction_docs(); out.lastEvictionBytesFreed = response.last_eviction_bytes_freed(); out.collections.reserve(response.collections_size()); for (const auto& pbc : response.collections()) { Client::MemoryCollectionStats c; c.collection = pbc.collection(); c.documentCount = pbc.document_count(); c.estimatedBytes = pbc.estimated_bytes(); c.evictedStubCount = pbc.evicted_stub_count(); switch (pbc.priority()) { case smartbotic::databasepb::MEMORY_PRIORITY_LOW: c.priority = "low"; break; case smartbotic::databasepb::MEMORY_PRIORITY_NORMAL: c.priority = "normal"; break; case smartbotic::databasepb::MEMORY_PRIORITY_HIGH: c.priority = "high"; break; default: c.priority = "normal"; break; } out.collections.push_back(std::move(c)); } return out; } // ===== Read-Only Control ===== bool setReadOnly(bool readOnly) { smartbotic::databasepb::SetReadOnlyRequest request; request.set_read_only(readOnly); smartbotic::databasepb::SetReadOnlyResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->SetReadOnly(&context, request, &response); if (!status.ok()) { spdlog::error("Client::setReadOnly failed: {}", status.error_message()); return false; } return response.success(); } Client::ReadOnlyStatus getReadOnlyStatus() { smartbotic::databasepb::GetReadOnlyStatusRequest request; smartbotic::databasepb::GetReadOnlyStatusResponse response; grpc::ClientContext context; setDeadline(context); Client::ReadOnlyStatus out; auto status = stub_->GetReadOnlyStatus(&context, request, &response); if (!status.ok()) { spdlog::error("Client::getReadOnlyStatus failed: {}", status.error_message()); return out; } out.readOnly = response.read_only(); out.reason = response.reason(); out.recoveryOutcome = response.recovery_outcome(); out.expectedSnapshot = response.expected_snapshot(); out.snapshotUsed = response.snapshot_used(); out.failureReason = response.failure_reason(); out.walEntriesReplayed = response.wal_entries_replayed(); out.snapshotsAttempted = response.snapshots_attempted(); return out; } // ===== File Operations ===== Client::FileUploadResult uploadFile(const std::vector& data, const Client::FileUploadMeta& meta) { grpc::ClientContext context; auto timeout_ms = config_.timeoutMs + (data.size() / (1024 * 1024)) * 1000; context.set_deadline( std::chrono::system_clock::now() + std::chrono::milliseconds(timeout_ms)); smartbotic::databasepb::UploadFileResponse response; auto writer = stub_->UploadFile(&context, &response); // First chunk: metadata smartbotic::databasepb::FileChunk metaChunk; auto* m = metaChunk.mutable_metadata(); m->set_name(meta.name); m->set_mime_type(meta.mime_type); m->set_file_type(meta.file_type); m->set_related_id(meta.related_id); m->set_is_public(meta.is_public); for (const auto& [key, value] : meta.metadata) { (*m->mutable_metadata())[key] = value; } writer->Write(metaChunk); // Data chunks (64KB each) constexpr size_t CHUNK_SIZE = 64 * 1024; for (size_t offset = 0; offset < data.size(); offset += CHUNK_SIZE) { smartbotic::databasepb::FileChunk dataChunk; size_t chunkSize = std::min(CHUNK_SIZE, data.size() - offset); dataChunk.set_data(data.data() + offset, chunkSize); if (!writer->Write(dataChunk)) break; } writer->WritesDone(); auto status = writer->Finish(); if (!status.ok()) { spdlog::error("Client::uploadFile failed: {}", status.error_message()); throw std::runtime_error(status.error_message()); } return {response.id(), response.size(), response.checksum(), response.deduplicated()}; } std::vector downloadFile(const std::string& id) { smartbotic::databasepb::DownloadFileRequest request; request.set_id(id); grpc::ClientContext context; context.set_deadline( std::chrono::system_clock::now() + std::chrono::milliseconds(config_.timeoutMs * 10)); auto reader = stub_->DownloadFile(&context, request); std::vector fileData; smartbotic::databasepb::FileChunk chunk; while (reader->Read(&chunk)) { if (chunk.has_data()) { const auto& d = chunk.data(); fileData.insert(fileData.end(), d.begin(), d.end()); } } auto status = reader->Finish(); if (!status.ok()) { spdlog::error("Client::downloadFile failed: {}", status.error_message()); throw std::runtime_error(status.error_message()); } return fileData; } std::optional getFileInfo(const std::string& id) { smartbotic::databasepb::GetFileInfoRequest request; request.set_id(id); smartbotic::databasepb::FileInfo response; grpc::ClientContext context; setDeadline(context); auto status = stub_->GetFileInfo(&context, request, &response); if (!status.ok()) { if (status.error_code() == grpc::StatusCode::NOT_FOUND) return std::nullopt; spdlog::error("Client::getFileInfo failed: {}", status.error_message()); throw std::runtime_error(status.error_message()); } Client::FileRecord record; record.id = response.id(); record.name = response.name(); record.mime_type = response.mime_type(); record.size = response.size(); record.file_type = response.file_type(); record.related_id = response.related_id(); record.checksum = response.checksum(); record.is_public = response.is_public(); record.ref_count = response.ref_count(); record.created_at = response.created_at(); for (const auto& [key, value] : response.metadata()) { record.metadata[key] = value; } return record; } bool deleteFile(const std::string& id) { smartbotic::databasepb::DeleteFileRequest request; request.set_id(id); smartbotic::databasepb::DeleteFileResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->DeleteFile(&context, request, &response); if (!status.ok()) { spdlog::error("Client::deleteFile failed: {}", status.error_message()); return false; } return response.deleted(); } Client::FileListResult listFiles(const std::string& file_type, const std::string& related_id, uint32_t limit, uint32_t offset, const std::string& checksum, const std::string& name) { smartbotic::databasepb::ListFilesRequest request; if (!file_type.empty()) request.set_file_type(file_type); if (!related_id.empty()) request.set_related_id(related_id); request.set_limit(limit); request.set_offset(offset); if (!checksum.empty()) request.set_checksum(checksum); if (!name.empty()) request.set_name(name); smartbotic::databasepb::ListFilesResponse response; grpc::ClientContext context; setDeadline(context); auto status = stub_->ListFiles(&context, request, &response); if (!status.ok()) { spdlog::error("Client::listFiles failed: {}", status.error_message()); throw std::runtime_error(status.error_message()); } Client::FileListResult result; result.total_count = response.total_count(); result.has_more = response.has_more(); for (const auto& f : response.files()) { Client::FileRecord record; record.id = f.id(); record.name = f.name(); record.mime_type = f.mime_type(); record.size = f.size(); record.file_type = f.file_type(); record.related_id = f.related_id(); record.checksum = f.checksum(); record.is_public = f.is_public(); record.created_at = f.created_at(); for (const auto& [key, value] : f.metadata()) { record.metadata[key] = value; } result.files.push_back(std::move(record)); } return result; } private: void setDeadline(grpc::ClientContext& context) { context.set_deadline( std::chrono::system_clock::now() + std::chrono::milliseconds(config_.timeoutMs) ); } Config config_; std::shared_ptr channel_; std::unique_ptr stub_; std::atomic connected_{false}; }; // ===== Client Public Interface Implementation ===== Client::Client(Config config) : impl_(std::make_unique(std::move(config))) {} Client::~Client() = default; Client::Client(Client&&) noexcept = default; Client& Client::operator=(Client&&) noexcept = default; bool Client::connect() { return impl_->connect(); } bool Client::isConnected() const { return impl_->isConnected(); } std::string Client::insert(const std::string& collection, const nlohmann::json& data, const std::string& id, uint32_t ttlSeconds, const std::string& actor) { return impl_->insert(collection, data, id, ttlSeconds, actor); } std::optional Client::get(const std::string& collection, const std::string& id) { return impl_->get(collection, id); } bool Client::update(const std::string& collection, const std::string& id, const nlohmann::json& data, const std::string& actor) { return impl_->update(collection, id, data, actor); } bool Client::updateIfVersion(const std::string& collection, const std::string& id, const nlohmann::json& data, uint64_t expectedVersion, const std::string& actor) { return impl_->updateIfVersion(collection, id, data, expectedVersion, actor); } uint64_t Client::patch(const std::string& collection, const std::string& id, const nlohmann::json& fields, const std::string& actor) { return impl_->patch(collection, id, fields, actor); } std::pair Client::upsert(const std::string& collection, const nlohmann::json& data, const std::string& id, uint32_t ttlSeconds, const std::string& actor) { return impl_->upsert(collection, data, id, ttlSeconds, actor); } bool Client::remove(const std::string& collection, const std::string& id) { return impl_->remove(collection, id); } bool Client::exists(const std::string& collection, const std::string& id) { return impl_->exists(collection, id); } Client::VersionHistoryResult Client::getVersionHistory(const std::string& collection, const std::string& id, uint32_t limit, uint32_t offset) { return impl_->getVersionHistory(collection, id, limit, offset); } std::optional Client::getDocumentVersion(const std::string& collection, const std::string& id, uint64_t version) { return impl_->getDocumentVersion(collection, id, version); } uint64_t Client::restoreVersion(const std::string& collection, const std::string& id, uint64_t version, const std::string& actor) { return impl_->restoreVersion(collection, id, version, actor); } std::pair Client::restoreToDate(const std::string& collection, const std::string& id, uint64_t timestamp, const std::string& actor) { return impl_->restoreToDate(collection, id, timestamp, actor); } std::vector Client::find(const std::string& collection, const QueryOptions& options) { return impl_->find(collection, options); } std::vector Client::find(const std::string& collection) { return impl_->find(collection, QueryOptions{}); } Client::FindResult Client::findWithMetrics(const std::string& collection, const QueryOptions& options) { return impl_->findWithMetrics(collection, options); } uint64_t Client::count(const std::string& collection, const std::vector& filters) { return impl_->count(collection, filters); } uint64_t Client::count(const std::string& collection) { return impl_->count(collection, {}); } bool Client::setAdd(const std::string& collection, const std::string& setId, const std::string& member) { return impl_->setAdd(collection, setId, member); } bool Client::setRemove(const std::string& collection, const std::string& setId, const std::string& member) { return impl_->setRemove(collection, setId, member); } std::vector Client::setMembers(const std::string& collection, const std::string& setId) { return impl_->setMembers(collection, setId); } bool Client::setIsMember(const std::string& collection, const std::string& setId, const std::string& member) { return impl_->setIsMember(collection, setId, member); } std::vector Client::similaritySearch( const std::string& collection, const std::vector& queryVector, uint32_t topK, float minScore) { return impl_->similaritySearch(collection, queryVector, topK, minScore); } bool Client::createCollection(const std::string& name, uint32_t defaultTtlSeconds, bool encrypted, uint32_t maxVersions, uint32_t vectorDimension) { return impl_->createCollection(name, defaultTtlSeconds, encrypted, maxVersions, vectorDimension); } bool Client::dropCollection(const std::string& name) { return impl_->dropCollection(name); } std::vector Client::listCollections() { return impl_->listCollections(); } std::optional Client::getCollectionInfo(const std::string& name) { return impl_->getCollectionInfo(name); } bool Client::configureCollection(const std::string& collection, const CollectionConfig& cfg) { return impl_->configureCollection(collection, cfg); } Client::CollectionConfig Client::getCollectionConfig(const std::string& collection) { return impl_->getCollectionConfig(collection); } bool Client::hasCollectionConfig(const std::string& collection) { return impl_->hasCollectionConfig(collection); } Client::TimestampMigrationResult Client::migrateCollectionTimestamps( const std::string& collection, const std::string& fromPrecision, const std::string& toPrecision ) { return impl_->migrateCollectionTimestamps(collection, fromPrecision, toPrecision); } bool Client::createView(const std::string& name, const std::string& collection, const std::vector& include, const std::vector& exclude, const std::vector& where, const std::optional& defaultSort) { return impl_->createView(name, collection, include, exclude, where, defaultSort); } bool Client::dropView(const std::string& name) { return impl_->dropView(name); } std::vector Client::listViews() { return impl_->listViews(); } std::optional Client::getViewInfo(const std::string& name) { return impl_->getViewInfo(name); } std::shared_ptr Client::subscribe(const std::vector& collections, EventCallback callback) { return impl_->subscribe(collections, std::move(callback)); } bool Client::healthCheck() { return impl_->healthCheck(); } std::optional Client::getHealthInfo() { return impl_->getHealthInfo(); } std::optional Client::getStats() { return impl_->getStats(); } Client::MemoryStats Client::getMemoryStats() { return impl_->getMemoryStats(); } bool Client::setReadOnly(bool readOnly) { return impl_->setReadOnly(readOnly); } Client::ReadOnlyStatus Client::getReadOnlyStatus() { return impl_->getReadOnlyStatus(); } Client::FileUploadResult Client::uploadFile(const std::vector& data, const FileUploadMeta& meta) { return impl_->uploadFile(data, meta); } std::vector Client::downloadFile(const std::string& id) { return impl_->downloadFile(id); } std::optional Client::getFileInfo(const std::string& id) { return impl_->getFileInfo(id); } bool Client::deleteFile(const std::string& id) { return impl_->deleteFile(id); } Client::FileListResult Client::listFiles(const std::string& file_type, const std::string& related_id, uint32_t limit, uint32_t offset, const std::string& checksum, const std::string& name) { return impl_->listFiles(file_type, related_id, limit, offset, checksum, name); } } // namespace smartbotic::database