#include "database_service.hpp" #include "logging/logger.hpp" #include "common/time_utils.hpp" #include 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(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(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(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(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(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_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(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("grpc_port", 9001); config.data_directory = cfg.getOr("data_directory", "./data/database"); config.max_message_size_mb = cfg.getOr("max_message_size_mb", 64); config.wal.sync_interval_ms = cfg.getOr("persistence.wal_sync_interval_ms", 100); config.snapshot.interval_sec = cfg.getOr("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(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()); // Set max message size to handle large execution results int max_msg_size = config_.max_message_size_mb * 1024 * 1024; builder.SetMaxReceiveMessageSize(max_msg_size); builder.SetMaxSendMessageSize(max_msg_size); LOG_INFO("Max gRPC message size: {} MB", config_.max_message_size_mb); 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 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 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 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