| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602 |
- #include "database_service.hpp"
- #include "logging/logger.hpp"
- #include "common/time_utils.hpp"
- #include <grpcpp/health_check_service_interface.h>
- namespace smartbotic::database {
- using namespace common;
- // DatabaseServiceImpl implementation
- DatabaseServiceImpl::DatabaseServiceImpl(MemoryStore& store)
- : store_(store) {}
- grpc::Status DatabaseServiceImpl::Get(grpc::ServerContext* context,
- const proto::GetRequest* request,
- proto::Document* response) {
- auto result = store_.get(request->collection(), request->id());
- if (result.failed()) {
- return grpc::Status(grpc::StatusCode::NOT_FOUND, result.error().message());
- }
- documentToProto(result.value(), response);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::Query(grpc::ServerContext* context,
- const proto::QueryRequest* request,
- proto::QueryResponse* response) {
- struct Query q(request->collection());
- // Build filters
- for (const auto& f : request->filters()) {
- Filter filter;
- filter.field = f.field();
- filter.op = protoToFilterOp(f.op());
- filter.value = nlohmann::json::parse(f.value());
- q.filters.push_back(std::move(filter));
- }
- // Build sorts
- for (const auto& s : request->sorts()) {
- Sort sort;
- sort.field = s.field();
- sort.direction = s.direction() == proto::SORT_DIRECTION_DESC ?
- SortDirection::Desc : SortDirection::Asc;
- q.sorts.push_back(std::move(sort));
- }
- // Pagination
- if (request->has_pagination()) {
- q.offset = (request->pagination().page() - 1) * request->pagination().page_size();
- q.limit = request->pagination().page_size();
- }
- // Projection
- for (const auto& field : request->fields()) {
- q.fields.push_back(field);
- }
- auto result = store_.query(q);
- for (const auto& doc : result.documents) {
- documentToProto(doc, response->add_documents());
- }
- auto* pagination = response->mutable_pagination();
- pagination->set_total_count(result.total_count);
- pagination->set_has_more(result.has_more);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::Insert(grpc::ServerContext* context,
- const proto::InsertRequest* request,
- proto::MutationResponse* response) {
- Document doc;
- doc.id = request->id();
- doc.collection = request->collection();
- doc.data = nlohmann::json::parse(request->data());
- if (request->ttl_ms() > 0) {
- doc.expires_at = TimeUtils::nowMs() + request->ttl_ms();
- }
- auto result = store_.insert(request->collection(), std::move(doc));
- if (result.failed()) {
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(static_cast<int32_t>(result.error().code()));
- error->set_message(result.error().message());
- return grpc::Status::OK;
- }
- response->set_id(result.value().id);
- response->set_version(result.value().version);
- response->set_success(true);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::Update(grpc::ServerContext* context,
- const proto::UpdateRequest* request,
- proto::MutationResponse* response) {
- auto data = nlohmann::json::parse(request->data());
- auto result = store_.update(request->collection(), request->id(), data,
- request->expected_version(), request->partial());
- if (result.failed()) {
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(static_cast<int32_t>(result.error().code()));
- error->set_message(result.error().message());
- return grpc::Status::OK;
- }
- response->set_id(result.value().id);
- response->set_version(result.value().version);
- response->set_success(true);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::Delete(grpc::ServerContext* context,
- const proto::DeleteRequest* request,
- proto::MutationResponse* response) {
- auto result = store_.remove(request->collection(), request->id(),
- request->expected_version());
- if (result.failed()) {
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(static_cast<int32_t>(result.error().code()));
- error->set_message(result.error().message());
- return grpc::Status::OK;
- }
- response->set_id(request->id());
- response->set_success(true);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::BatchInsert(grpc::ServerContext* context,
- const proto::BatchInsertRequest* request,
- proto::BatchMutationResponse* response) {
- int32_t success_count = 0;
- int32_t failure_count = 0;
- for (const auto& req : request->documents()) {
- Document doc;
- doc.id = req.id();
- doc.collection = request->collection();
- doc.data = nlohmann::json::parse(req.data());
- if (req.ttl_ms() > 0) {
- doc.expires_at = TimeUtils::nowMs() + req.ttl_ms();
- }
- auto result = store_.insert(request->collection(), std::move(doc));
- auto* mutation_result = response->add_results();
- if (result.ok()) {
- mutation_result->set_id(result.value().id);
- mutation_result->set_version(result.value().version);
- mutation_result->set_success(true);
- ++success_count;
- } else {
- mutation_result->set_success(false);
- auto* error = mutation_result->mutable_error();
- error->set_code(static_cast<int32_t>(result.error().code()));
- error->set_message(result.error().message());
- ++failure_count;
- }
- }
- response->set_success_count(success_count);
- response->set_failure_count(failure_count);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::BatchDelete(grpc::ServerContext* context,
- const proto::BatchDeleteRequest* request,
- proto::BatchMutationResponse* response) {
- int32_t success_count = 0;
- int32_t failure_count = 0;
- for (const auto& id : request->ids()) {
- auto result = store_.remove(request->collection(), id, 0);
- auto* mutation_result = response->add_results();
- mutation_result->set_id(id);
- if (result.ok()) {
- mutation_result->set_success(true);
- ++success_count;
- } else {
- mutation_result->set_success(false);
- auto* error = mutation_result->mutable_error();
- error->set_code(static_cast<int32_t>(result.error().code()));
- error->set_message(result.error().message());
- ++failure_count;
- }
- }
- response->set_success_count(success_count);
- response->set_failure_count(failure_count);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::CreateCollection(grpc::ServerContext* context,
- const proto::CreateCollectionRequest* request,
- proto::Empty* response) {
- CollectionConfig config(request->name());
- config.default_ttl_ms = request->default_ttl_ms();
- if (!request->schema().empty()) {
- config.schema = nlohmann::json::parse(request->schema());
- }
- // Handle versioning configuration
- if (request->has_versioning()) {
- VersioningConfig versioning;
- versioning.enabled = request->versioning().enabled();
- versioning.max_versions = request->versioning().max_versions();
- versioning.version_ttl_ms = request->versioning().version_ttl_ms();
- versioning.keep_on_delete = request->versioning().keep_on_delete();
- config.versioning = versioning;
- }
- auto result = store_.createCollection(config);
- if (result.failed()) {
- return grpc::Status(grpc::StatusCode::ALREADY_EXISTS, result.error().message());
- }
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::DropCollection(grpc::ServerContext* context,
- const proto::DropCollectionRequest* request,
- proto::Empty* response) {
- auto result = store_.dropCollection(request->name());
- if (result.failed()) {
- return grpc::Status(grpc::StatusCode::NOT_FOUND, result.error().message());
- }
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::ListCollections(grpc::ServerContext* context,
- const proto::ListCollectionsRequest* request,
- proto::ListCollectionsResponse* response) {
- auto collections = store_.listCollections();
- for (const auto& name : collections) {
- response->add_collections(name);
- }
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::GetVersion(grpc::ServerContext* context,
- const proto::GetVersionRequest* request,
- proto::Document* response) {
- auto result = store_.getVersion(request->collection(), request->id(), request->version());
- if (result.failed()) {
- return grpc::Status(grpc::StatusCode::NOT_FOUND, result.error().message());
- }
- const auto& ver = result.value();
- response->set_id(ver.document_id);
- response->set_collection(request->collection());
- response->set_data(ver.data.dump());
- response->set_version(ver.version);
- response->set_created_at(ver.created_at);
- response->set_updated_at(ver.created_at);
- response->set_expires_at(ver.expires_at);
- return grpc::Status::OK;
- }
- grpc::Status DatabaseServiceImpl::ListVersions(grpc::ServerContext* context,
- const proto::ListVersionsRequest* request,
- proto::ListVersionsResponse* response) {
- int32_t limit = request->limit() > 0 ? request->limit() : 100;
- int32_t offset = request->offset();
- auto result = store_.listVersions(request->collection(), request->id(), limit, offset);
- for (const auto& doc : result.documents) {
- auto ver = DocumentVersion::fromJson(doc.data);
- auto* proto_ver = response->add_versions();
- documentVersionToProto(ver, proto_ver);
- }
- response->set_total_count(result.total_count);
- response->set_has_more(result.has_more);
- return grpc::Status::OK;
- }
- void DatabaseServiceImpl::documentToProto(const Document& doc, proto::Document* proto) {
- proto->set_id(doc.id);
- proto->set_collection(doc.collection);
- proto->set_data(doc.data.dump());
- proto->set_version(doc.version);
- proto->set_created_at(doc.created_at);
- proto->set_updated_at(doc.updated_at);
- proto->set_expires_at(doc.expires_at);
- }
- void DatabaseServiceImpl::documentVersionToProto(const DocumentVersion& ver, proto::DocumentVersion* proto) {
- proto->set_id(ver.id);
- proto->set_document_id(ver.document_id);
- proto->set_version(ver.version);
- proto->set_data(ver.data.dump());
- proto->set_created_at(ver.created_at);
- proto->set_expires_at(ver.expires_at);
- }
- FilterOp DatabaseServiceImpl::protoToFilterOp(proto::FilterOp op) {
- switch (op) {
- case proto::FILTER_OP_EQ: return FilterOp::Eq;
- case proto::FILTER_OP_NE: return FilterOp::Ne;
- case proto::FILTER_OP_GT: return FilterOp::Gt;
- case proto::FILTER_OP_GTE: return FilterOp::Gte;
- case proto::FILTER_OP_LT: return FilterOp::Lt;
- case proto::FILTER_OP_LTE: return FilterOp::Lte;
- case proto::FILTER_OP_IN: return FilterOp::In;
- case proto::FILTER_OP_NIN: return FilterOp::Nin;
- case proto::FILTER_OP_CONTAINS: return FilterOp::Contains;
- case proto::FILTER_OP_REGEX: return FilterOp::Regex;
- case proto::FILTER_OP_EXISTS: return FilterOp::Exists;
- default: return FilterOp::Eq;
- }
- }
- // DatabaseService implementation
- DatabaseService::DatabaseService(const Config& config)
- : config_(config) {
- // Setup WAL
- WAL::Config wal_config;
- wal_config.directory = std::filesystem::path(config_.data_directory) / "wal";
- wal_config.sync_interval_ms = config_.wal.sync_interval_ms;
- wal_config.enabled = true;
- wal_ = std::make_unique<WAL>(wal_config);
- // Setup snapshot manager
- SnapshotManager::Config snapshot_config;
- snapshot_config.directory = std::filesystem::path(config_.data_directory) / "snapshots";
- snapshot_config.interval_sec = config_.snapshot.interval_sec;
- snapshot_config.enabled = true;
- snapshot_manager_ = std::make_unique<SnapshotManager>(snapshot_config);
- // Setup mutation callback for WAL
- store_.setMutationCallback([this](const WalEntry& entry) {
- wal_->append(entry);
- });
- }
- DatabaseService::~DatabaseService() {
- stop();
- }
- DatabaseService::Config DatabaseService::loadConfig(const std::filesystem::path& path) {
- Config config;
- auto result = config::Config::fromFile(path);
- if (result.ok()) {
- auto& cfg = result.value();
- config.grpc_port = cfg.getOr<int>("grpc_port", 9001);
- config.data_directory = cfg.getOr<std::string>("data_directory", "./data/database");
- config.wal.sync_interval_ms = cfg.getOr<int>("persistence.wal_sync_interval_ms", 100);
- config.snapshot.interval_sec = cfg.getOr<int>("persistence.snapshot_interval_sec", 3600);
- }
- return config;
- }
- void DatabaseService::start() {
- if (running_) {
- return;
- }
- LOG_INFO("Starting database service...");
- // Create data directory
- std::filesystem::create_directories(config_.data_directory);
- // Recovery from persistence
- recoveryFromPersistence();
- // Create predefined collections
- createPredefinedCollections();
- // Start WAL
- wal_->start();
- // Start gRPC server
- service_impl_ = std::make_unique<DatabaseServiceImpl>(store_);
- // Disable health check to avoid async completion queue issues
- // grpc::EnableDefaultHealthCheckService(true);
- grpc::ServerBuilder builder;
- builder.AddListeningPort("0.0.0.0:" + std::to_string(config_.grpc_port),
- grpc::InsecureServerCredentials());
- builder.RegisterService(service_impl_.get());
- server_ = builder.BuildAndStart();
- LOG_INFO("Database gRPC server listening on port {}", config_.grpc_port);
- running_ = true;
- // Start background threads
- snapshot_thread_ = std::thread(&DatabaseService::snapshotLoop, this);
- ttl_thread_ = std::thread(&DatabaseService::ttlLoop, this);
- }
- void DatabaseService::stop() {
- if (!running_) {
- return;
- }
- LOG_INFO("Stopping database service...");
- // Signal shutdown to background threads
- {
- std::lock_guard<std::mutex> lock(shutdown_mutex_);
- running_ = false;
- }
- shutdown_cv_.notify_all();
- // Stop background threads (they will wake up immediately now)
- if (snapshot_thread_.joinable()) {
- snapshot_thread_.join();
- }
- if (ttl_thread_.joinable()) {
- ttl_thread_.join();
- }
- // Final snapshot before shutdown
- LOG_INFO("Creating final snapshot before shutdown...");
- createSnapshot();
- // Stop WAL
- wal_->stop();
- // Stop gRPC server
- if (server_) {
- server_->Shutdown();
- }
- LOG_INFO("Database service stopped");
- }
- void DatabaseService::createPredefinedCollections() {
- // Users collection
- store_.createCollection(CollectionConfig("users"));
- // Workflows collection
- store_.createCollection(CollectionConfig("workflows"));
- // Workflow groups collection
- store_.createCollection(CollectionConfig("workflow_groups"));
- // Executions collection with 7-day TTL
- store_.createCollection(
- CollectionConfig("executions").withTTL(7 * 24 * 60 * 60 * 1000));
- // Sessions collection with 1-day TTL
- store_.createCollection(
- CollectionConfig("sessions").withTTL(24 * 60 * 60 * 1000));
- // API keys collection
- store_.createCollection(CollectionConfig("api_keys"));
- // Credentials collection
- store_.createCollection(CollectionConfig("credentials"));
- // Runners collection
- store_.createCollection(CollectionConfig("runners"));
- // Nodes collection - stores node definitions
- store_.createCollection(CollectionConfig("nodes"));
- LOG_INFO("Created predefined collections");
- }
- void DatabaseService::createSnapshot() {
- auto snapshot = store_.createSnapshot();
- auto result = snapshot_manager_->save(snapshot);
- if (result.ok()) {
- // Truncate WAL after successful snapshot
- wal_->truncate(snapshot.value("sequence", 0));
- }
- }
- void DatabaseService::recoveryFromPersistence() {
- LOG_INFO("Starting recovery from persistence...");
- // Try to load latest snapshot
- auto snapshot_result = snapshot_manager_->loadLatest();
- if (snapshot_result.ok()) {
- store_.loadSnapshot(snapshot_result.value());
- LOG_INFO("Loaded snapshot");
- }
- // Replay WAL entries after snapshot
- int64_t last_sequence = snapshot_manager_->getLatestSequence();
- wal_->replay([this, last_sequence](const WalEntry& entry) {
- if (entry.sequence <= last_sequence) {
- return; // Skip entries already in snapshot
- }
- // Apply WAL entry
- switch (entry.type) {
- case WalEntryType::Insert: {
- auto doc = Document::fromJson(entry.data);
- store_.getCollection(entry.collection)->loadFromSnapshot({doc});
- break;
- }
- case WalEntryType::Update: {
- auto doc = Document::fromJson(entry.data);
- store_.getCollection(entry.collection)->loadFromSnapshot({doc});
- break;
- }
- case WalEntryType::Delete: {
- auto* col = store_.getCollection(entry.collection);
- if (col) {
- col->remove(entry.document_id, 0);
- }
- break;
- }
- case WalEntryType::CreateCollection: {
- CollectionConfig config(entry.collection);
- if (entry.data.contains("default_ttl_ms")) {
- config.default_ttl_ms = entry.data["default_ttl_ms"];
- }
- if (entry.data.contains("versioning")) {
- config.versioning = VersioningConfig::fromJson(entry.data["versioning"]);
- }
- store_.createCollection(config);
- break;
- }
- case WalEntryType::DropCollection: {
- store_.dropCollection(entry.collection);
- break;
- }
- case WalEntryType::StoreVersion: {
- // Version entries are stored in-memory, no special recovery needed
- // They are reconstructed from document updates during normal operation
- break;
- }
- case WalEntryType::DeleteVersion: {
- // Handled similarly to StoreVersion
- break;
- }
- }
- });
- LOG_INFO("Recovery complete");
- }
- void DatabaseService::snapshotLoop() {
- while (running_) {
- // Wait for shutdown signal or timeout
- {
- std::unique_lock<std::mutex> lock(shutdown_mutex_);
- if (shutdown_cv_.wait_for(lock,
- std::chrono::seconds(config_.snapshot.interval_sec),
- [this] { return !running_.load(); })) {
- // Shutdown signaled, exit loop
- break;
- }
- }
- if (!running_) break;
- // Check if WAL is large enough to trigger snapshot
- if (wal_->getSize() >= config_.snapshot.wal_size_trigger) {
- createSnapshot();
- }
- }
- }
- void DatabaseService::ttlLoop() {
- constexpr int ttl_check_interval_sec = 60; // Check every minute
- while (running_) {
- // Wait for shutdown signal or timeout
- {
- std::unique_lock<std::mutex> lock(shutdown_mutex_);
- if (shutdown_cv_.wait_for(lock,
- std::chrono::seconds(ttl_check_interval_sec),
- [this] { return !running_.load(); })) {
- // Shutdown signaled, exit loop
- break;
- }
- }
- if (!running_) break;
- store_.expireAllDocuments();
- store_.expireAllVersions();
- }
- }
- } // namespace smartbotic::database
|