|
@@ -0,0 +1,1200 @@
|
|
|
|
|
+#include "smartbotic/webserver/script_service.hpp"
|
|
|
|
|
+
|
|
|
|
|
+#include <algorithm>
|
|
|
|
|
+#include <chrono>
|
|
|
|
|
+#include <iomanip>
|
|
|
|
|
+#include <sstream>
|
|
|
|
|
+#include <random>
|
|
|
|
|
+
|
|
|
|
|
+#include <spdlog/spdlog.h>
|
|
|
|
|
+
|
|
|
|
|
+#include "smartbotic/webserver/authorization_service.hpp"
|
|
|
|
|
+#include "smartbotic/webserver/document_service.hpp"
|
|
|
|
|
+#include "smartbotic/webserver/event_manager.hpp"
|
|
|
|
|
+#include "smartbotic/webserver/permissions.hpp"
|
|
|
|
|
+#include "smartbotic/runner/js_engine.hpp"
|
|
|
|
|
+
|
|
|
|
|
+namespace smartbotic::webserver {
|
|
|
|
|
+
|
|
|
|
|
+namespace {
|
|
|
|
|
+
|
|
|
|
|
+/// Get current timestamp as ISO 8601 string
|
|
|
|
|
+auto GetTimestamp() -> std::string {
|
|
|
|
|
+ auto now = std::chrono::system_clock::now();
|
|
|
|
|
+ auto time = std::chrono::system_clock::to_time_t(now);
|
|
|
|
|
+ auto ms = std::chrono::duration_cast<std::chrono::milliseconds>(
|
|
|
|
|
+ now.time_since_epoch()) % 1000;
|
|
|
|
|
+
|
|
|
|
|
+ std::tm tm{};
|
|
|
|
|
+ gmtime_r(&time, &tm);
|
|
|
|
|
+
|
|
|
|
|
+ std::ostringstream oss;
|
|
|
|
|
+ oss << std::put_time(&tm, "%Y-%m-%dT%H:%M:%S");
|
|
|
|
|
+ oss << "." << std::setfill('0') << std::setw(3) << ms.count() << "Z";
|
|
|
|
|
+ return oss.str();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+/// Generate a UUID v4
|
|
|
|
|
+auto GenerateUuidV4() -> std::string {
|
|
|
|
|
+ static std::random_device rd;
|
|
|
|
|
+ static std::mt19937 gen(rd());
|
|
|
|
|
+ static std::uniform_int_distribution<> dis(0, 15);
|
|
|
|
|
+ static std::uniform_int_distribution<> dis2(8, 11);
|
|
|
|
|
+
|
|
|
|
|
+ std::stringstream ss;
|
|
|
|
|
+ ss << std::hex;
|
|
|
|
|
+ for (int i = 0; i < 8; i++) { ss << dis(gen); }
|
|
|
|
|
+ ss << "-";
|
|
|
|
|
+ for (int i = 0; i < 4; i++) { ss << dis(gen); }
|
|
|
|
|
+ ss << "-4"; // Version 4
|
|
|
|
|
+ for (int i = 0; i < 3; i++) { ss << dis(gen); }
|
|
|
|
|
+ ss << "-";
|
|
|
|
|
+ ss << dis2(gen); // Variant
|
|
|
|
|
+ for (int i = 0; i < 3; i++) { ss << dis(gen); }
|
|
|
|
|
+ ss << "-";
|
|
|
|
|
+ for (int i = 0; i < 12; i++) { ss << dis(gen); }
|
|
|
|
|
+ return ss.str();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+} // namespace
|
|
|
|
|
+
|
|
|
|
|
+ScriptService::ScriptService(DatabaseClient& db_client) : db_client_(db_client) {}
|
|
|
|
|
+
|
|
|
|
|
+ScriptService::~ScriptService() {
|
|
|
|
|
+ StopScheduler();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::SetEventManager(EventManager* event_manager) {
|
|
|
|
|
+ event_manager_ = event_manager;
|
|
|
|
|
+ if (event_manager_) {
|
|
|
|
|
+ SubscribeToEvents();
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::SetAuthorizationService(AuthorizationService* auth_service) {
|
|
|
|
|
+ auth_service_ = auth_service;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// Move operations deleted because std::atomic and std::mutex are not movable
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::Initialize() -> bool {
|
|
|
|
|
+ auto* collection_service = db_client_.GetCollectionService();
|
|
|
|
|
+ if (collection_service == nullptr) {
|
|
|
|
|
+ spdlog::error("ScriptService: Collection service not available");
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Initialize _scripts collection
|
|
|
|
|
+ {
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::GetCollectionMetadataRequest req;
|
|
|
|
|
+ ::smartbotic::database::CollectionMetadata resp;
|
|
|
|
|
+ req.set_name(kScriptsCollection);
|
|
|
|
|
+ auto status = collection_service->GetCollectionMetadata(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok() && status.error_code() == grpc::StatusCode::NOT_FOUND) {
|
|
|
|
|
+ grpc::ClientContext create_ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateCollectionRequest create_req;
|
|
|
|
|
+ ::smartbotic::database::CollectionMetadata create_resp;
|
|
|
|
|
+ create_req.set_name(kScriptsCollection);
|
|
|
|
|
+ auto create_status = collection_service->CreateCollection(&create_ctx, create_req, &create_resp);
|
|
|
|
|
+ if (!create_status.ok()) {
|
|
|
|
|
+ spdlog::error("ScriptService: Failed to create {}: {}", kScriptsCollection, create_status.error_message());
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ spdlog::info("ScriptService: Created collection {}", kScriptsCollection);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Initialize _script_executions collection
|
|
|
|
|
+ {
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::GetCollectionMetadataRequest req;
|
|
|
|
|
+ ::smartbotic::database::CollectionMetadata resp;
|
|
|
|
|
+ req.set_name(kExecutionsCollection);
|
|
|
|
|
+ auto status = collection_service->GetCollectionMetadata(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok() && status.error_code() == grpc::StatusCode::NOT_FOUND) {
|
|
|
|
|
+ grpc::ClientContext create_ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateCollectionRequest create_req;
|
|
|
|
|
+ ::smartbotic::database::CollectionMetadata create_resp;
|
|
|
|
|
+ create_req.set_name(kExecutionsCollection);
|
|
|
|
|
+ auto create_status = collection_service->CreateCollection(&create_ctx, create_req, &create_resp);
|
|
|
|
|
+ if (!create_status.ok()) {
|
|
|
|
|
+ spdlog::error("ScriptService: Failed to create {}: {}", kExecutionsCollection, create_status.error_message());
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ spdlog::info("ScriptService: Created collection {}", kExecutionsCollection);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Initialize _script_storage collection
|
|
|
|
|
+ {
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::GetCollectionMetadataRequest req;
|
|
|
|
|
+ ::smartbotic::database::CollectionMetadata resp;
|
|
|
|
|
+ req.set_name(kStorageCollection);
|
|
|
|
|
+ auto status = collection_service->GetCollectionMetadata(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok() && status.error_code() == grpc::StatusCode::NOT_FOUND) {
|
|
|
|
|
+ grpc::ClientContext create_ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateCollectionRequest create_req;
|
|
|
|
|
+ ::smartbotic::database::CollectionMetadata create_resp;
|
|
|
|
|
+ create_req.set_name(kStorageCollection);
|
|
|
|
|
+ auto create_status = collection_service->CreateCollection(&create_ctx, create_req, &create_resp);
|
|
|
|
|
+ if (!create_status.ok()) {
|
|
|
|
|
+ spdlog::error("ScriptService: Failed to create {}: {}", kStorageCollection, create_status.error_message());
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ spdlog::info("ScriptService: Created collection {}", kStorageCollection);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("ScriptService: Initialized successfully");
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::StartScheduler() {
|
|
|
|
|
+ if (scheduler_running_.load()) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ scheduler_running_.store(true);
|
|
|
|
|
+ scheduler_thread_ = std::thread([this]() { SchedulerLoop(); });
|
|
|
|
|
+ spdlog::info("ScriptService: Cron scheduler started");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::StopScheduler() {
|
|
|
|
|
+ if (!scheduler_running_.load()) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ scheduler_running_.store(false);
|
|
|
|
|
+ if (scheduler_thread_.joinable()) {
|
|
|
|
|
+ scheduler_thread_.join();
|
|
|
|
|
+ }
|
|
|
|
|
+ spdlog::info("ScriptService: Cron scheduler stopped");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::SubscribeToEvents() {
|
|
|
|
|
+ if (event_manager_ == nullptr) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ event_listener_id_ = event_manager_->AddListener([this](const EntityEvent& event) {
|
|
|
|
|
+ OnEntityEvent(event);
|
|
|
|
|
+ });
|
|
|
|
|
+ spdlog::debug("ScriptService: Subscribed to entity events");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// === CRUD Operations ===
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::CreateScript(const CreateScriptRequest& request) -> ScriptResult {
|
|
|
|
|
+ ScriptInfo script;
|
|
|
|
|
+ script.id = GenerateUuidV4();
|
|
|
|
|
+ script.workspace_id = request.workspace_id;
|
|
|
|
|
+ script.name = request.name;
|
|
|
|
|
+ script.description = request.description;
|
|
|
|
|
+ script.code = request.code;
|
|
|
|
|
+ script.status = ScriptStatus::Draft;
|
|
|
|
|
+ script.config = request.config;
|
|
|
|
|
+ script.triggers = request.triggers;
|
|
|
|
|
+ script.created_at = GetTimestamp();
|
|
|
|
|
+ script.updated_at = script.created_at;
|
|
|
|
|
+ script.created_by = request.created_by;
|
|
|
|
|
+
|
|
|
|
|
+ // Assign IDs to triggers if not set
|
|
|
|
|
+ for (auto& trigger : script.triggers) {
|
|
|
|
|
+ if (trigger.id.empty()) {
|
|
|
|
|
+ trigger.id = GenerateUuidV4();
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Validate config
|
|
|
|
|
+ if (script.config.timeout_ms > 30000) {
|
|
|
|
|
+ script.config.timeout_ms = 30000;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (script.config.max_memory_mb > 64) {
|
|
|
|
|
+ script.config.max_memory_mb = 64;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto* doc_service = db_client_.GetDocumentService();
|
|
|
|
|
+ if (doc_service == nullptr) {
|
|
|
|
|
+ return {.success = false, .error = "Document service not available", .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+
|
|
|
|
|
+ req.set_collection(kScriptsCollection);
|
|
|
|
|
+ req.set_id(script.id);
|
|
|
|
|
+ *req.mutable_data() = DocumentService::JsonToMapValue(ScriptToJson(script));
|
|
|
|
|
+
|
|
|
|
|
+ auto status = doc_service->CreateDocument(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ spdlog::error("ScriptService: Failed to create script: {}", status.error_message());
|
|
|
|
|
+ return {.success = false, .error = status.error_message(), .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("ScriptService: Created script {} in workspace {}", script.id, script.workspace_id);
|
|
|
|
|
+
|
|
|
|
|
+ // Emit event
|
|
|
|
|
+ if (event_manager_) {
|
|
|
|
|
+ event_manager_->EmitScriptEvent(EventAction::Create, script.id, script.workspace_id,
|
|
|
|
|
+ ScriptToJson(script), script.created_by);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .script = script};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::GetScript(const std::string& workspace_id, const std::string& id,
|
|
|
|
|
+ bool include_deleted) -> ScriptResult {
|
|
|
|
|
+ auto* doc_service = db_client_.GetDocumentService();
|
|
|
|
|
+ if (doc_service == nullptr) {
|
|
|
|
|
+ return {.success = false, .error = "Document service not available", .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::GetDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+
|
|
|
|
|
+ req.set_collection(kScriptsCollection);
|
|
|
|
|
+ req.set_id(id);
|
|
|
|
|
+
|
|
|
|
|
+ auto status = doc_service->GetDocument(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ if (status.error_code() == grpc::StatusCode::NOT_FOUND) {
|
|
|
|
|
+ return {.success = false, .error = "Script not found", .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+ return {.success = false, .error = status.error_message(), .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto json = DocumentService::MapValueToJson(resp.data());
|
|
|
|
|
+ json["id"] = resp.id();
|
|
|
|
|
+
|
|
|
|
|
+ auto script = JsonToScript(json);
|
|
|
|
|
+
|
|
|
|
|
+ // Check workspace match
|
|
|
|
|
+ if (script.workspace_id != workspace_id) {
|
|
|
|
|
+ return {.success = false, .error = "Script not found", .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Check if deleted
|
|
|
|
|
+ if (script.IsDeleted() && !include_deleted) {
|
|
|
|
|
+ return {.success = false, .error = "Script not found", .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .script = script};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::ListScripts(const std::string& workspace_id,
|
|
|
|
|
+ bool include_deleted) -> ScriptListResult {
|
|
|
|
|
+ auto* query_service = db_client_.GetQueryService();
|
|
|
|
|
+ if (query_service == nullptr) {
|
|
|
|
|
+ return {.success = false, .error = "Query service not available", .scripts = {}, .total_count = 0};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::QueryRequest req;
|
|
|
|
|
+ ::smartbotic::database::QueryResponse resp;
|
|
|
|
|
+
|
|
|
|
|
+ req.set_collection(kScriptsCollection);
|
|
|
|
|
+ req.set_limit(1000);
|
|
|
|
|
+
|
|
|
|
|
+ // Build composite filter for workspace_id and deleted_at
|
|
|
|
|
+ auto* filter = req.mutable_filter();
|
|
|
|
|
+ auto* composite = filter->mutable_composite();
|
|
|
|
|
+ composite->set_operator_(::smartbotic::database::COMPOSITE_OPERATOR_AND);
|
|
|
|
|
+
|
|
|
|
|
+ auto* ws_filter = composite->add_filters()->mutable_field();
|
|
|
|
|
+ ws_filter->set_field("workspace_id");
|
|
|
|
|
+ ws_filter->set_operator_(::smartbotic::database::FILTER_OPERATOR_EQUAL);
|
|
|
|
|
+ ws_filter->mutable_value()->set_string_value(workspace_id);
|
|
|
|
|
+
|
|
|
|
|
+ if (!include_deleted) {
|
|
|
|
|
+ auto* deleted_filter = composite->add_filters()->mutable_field();
|
|
|
|
|
+ deleted_filter->set_field("deleted_at");
|
|
|
|
|
+ deleted_filter->set_operator_(::smartbotic::database::FILTER_OPERATOR_EQUAL);
|
|
|
|
|
+ deleted_filter->mutable_value()->set_string_value("");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto status = query_service->Query(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ return {.success = false, .error = status.error_message(), .scripts = {}, .total_count = 0};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ std::vector<ScriptInfo> scripts;
|
|
|
|
|
+ for (const auto& doc : resp.documents()) {
|
|
|
|
|
+ auto json = DocumentService::MapValueToJson(doc.data());
|
|
|
|
|
+ json["id"] = doc.id();
|
|
|
|
|
+ scripts.push_back(JsonToScript(json));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Sort by updated_at descending
|
|
|
|
|
+ std::sort(scripts.begin(), scripts.end(), [](const auto& a, const auto& b) {
|
|
|
|
|
+ return a.updated_at > b.updated_at;
|
|
|
|
|
+ });
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .scripts = scripts, .total_count = static_cast<int64_t>(scripts.size())};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::UpdateScript(const UpdateScriptRequest& request) -> ScriptResult {
|
|
|
|
|
+ auto existing = GetScript(request.workspace_id, request.id, false);
|
|
|
|
|
+ if (!existing.success || !existing.script) {
|
|
|
|
|
+ return {.success = false, .error = existing.error.empty() ? "Script not found" : existing.error, .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto script = *existing.script;
|
|
|
|
|
+
|
|
|
|
|
+ // Apply updates
|
|
|
|
|
+ if (request.name) {
|
|
|
|
|
+ script.name = *request.name;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (request.description) {
|
|
|
|
|
+ script.description = *request.description;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (request.code) {
|
|
|
|
|
+ script.code = *request.code;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (request.status) {
|
|
|
|
|
+ script.status = *request.status;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (request.config) {
|
|
|
|
|
+ script.config = *request.config;
|
|
|
|
|
+ // Validate limits
|
|
|
|
|
+ if (script.config.timeout_ms > 30000) {
|
|
|
|
|
+ script.config.timeout_ms = 30000;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (script.config.max_memory_mb > 64) {
|
|
|
|
|
+ script.config.max_memory_mb = 64;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ if (request.triggers) {
|
|
|
|
|
+ script.triggers = *request.triggers;
|
|
|
|
|
+ // Assign IDs to new triggers
|
|
|
|
|
+ for (auto& trigger : script.triggers) {
|
|
|
|
|
+ if (trigger.id.empty()) {
|
|
|
|
|
+ trigger.id = GenerateUuidV4();
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ script.updated_at = GetTimestamp();
|
|
|
|
|
+
|
|
|
|
|
+ auto* doc_service = db_client_.GetDocumentService();
|
|
|
|
|
+ if (doc_service == nullptr) {
|
|
|
|
|
+ return {.success = false, .error = "Document service not available", .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+
|
|
|
|
|
+ req.set_collection(kScriptsCollection);
|
|
|
|
|
+ req.set_id(script.id);
|
|
|
|
|
+ *req.mutable_data() = DocumentService::JsonToMapValue(ScriptToJson(script));
|
|
|
|
|
+
|
|
|
|
|
+ auto status = doc_service->UpdateDocument(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ return {.success = false, .error = status.error_message(), .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("ScriptService: Updated script {} in workspace {}", script.id, script.workspace_id);
|
|
|
|
|
+
|
|
|
|
|
+ // Update cron schedule if triggers changed
|
|
|
|
|
+ if (request.triggers || request.status) {
|
|
|
|
|
+ UpdateCronSchedule(script);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Emit event
|
|
|
|
|
+ if (event_manager_) {
|
|
|
|
|
+ event_manager_->EmitScriptEvent(EventAction::Update, script.id, script.workspace_id,
|
|
|
|
|
+ ScriptToJson(script), script.created_by);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .script = script};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::DeleteScript(const std::string& workspace_id, const std::string& id) -> ScriptResult {
|
|
|
|
|
+ auto existing = GetScript(workspace_id, id, false);
|
|
|
|
|
+ if (!existing.success || !existing.script) {
|
|
|
|
|
+ return {.success = false, .error = existing.error.empty() ? "Script not found" : existing.error, .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto script = *existing.script;
|
|
|
|
|
+ script.deleted_at = GetTimestamp();
|
|
|
|
|
+ script.updated_at = script.deleted_at;
|
|
|
|
|
+
|
|
|
|
|
+ auto* doc_service = db_client_.GetDocumentService();
|
|
|
|
|
+ if (doc_service == nullptr) {
|
|
|
|
|
+ return {.success = false, .error = "Document service not available", .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+
|
|
|
|
|
+ req.set_collection(kScriptsCollection);
|
|
|
|
|
+ req.set_id(script.id);
|
|
|
|
|
+ *req.mutable_data() = DocumentService::JsonToMapValue(ScriptToJson(script));
|
|
|
|
|
+
|
|
|
|
|
+ auto status = doc_service->UpdateDocument(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ return {.success = false, .error = status.error_message(), .script = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Remove from cron schedule
|
|
|
|
|
+ RemoveCronSchedule(script.id);
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("ScriptService: Deleted script {} in workspace {}", script.id, script.workspace_id);
|
|
|
|
|
+
|
|
|
|
|
+ // Emit event
|
|
|
|
|
+ if (event_manager_) {
|
|
|
|
|
+ event_manager_->EmitScriptEvent(EventAction::Delete, script.id, script.workspace_id,
|
|
|
|
|
+ ScriptToJson(script), script.created_by);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .script = script};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// === Execution Operations ===
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::ExecuteScript(const ExecuteScriptRequest& request) -> ScriptExecutionResult {
|
|
|
|
|
+ // Get the script
|
|
|
|
|
+ auto script_result = GetScript(request.workspace_id, request.script_id, false);
|
|
|
|
|
+ if (!script_result.success || !script_result.script) {
|
|
|
|
|
+ return {.success = false, .error = "Script not found", .execution = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ const auto& script = *script_result.script;
|
|
|
|
|
+
|
|
|
|
|
+ // Check if script is active (unless manual execution)
|
|
|
|
|
+ if (request.trigger_type != "manual" && script.status != ScriptStatus::Active) {
|
|
|
|
|
+ return {.success = false, .error = "Script is not active", .execution = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Create execution record
|
|
|
|
|
+ ScriptExecution execution;
|
|
|
|
|
+ execution.id = GenerateUuidV4();
|
|
|
|
|
+ execution.script_id = script.id;
|
|
|
|
|
+ execution.workspace_id = script.workspace_id;
|
|
|
|
|
+ execution.trigger_type = request.trigger_type;
|
|
|
|
|
+ execution.trigger_event = request.trigger_event;
|
|
|
|
|
+ execution.status = ExecutionStatus::Running;
|
|
|
|
|
+ execution.started_at = GetTimestamp();
|
|
|
|
|
+ execution.executed_by = request.executed_by.empty() ? script.created_by : request.executed_by;
|
|
|
|
|
+
|
|
|
|
|
+ // Save execution record
|
|
|
|
|
+ auto* doc_service = db_client_.GetDocumentService();
|
|
|
|
|
+ if (doc_service != nullptr) {
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+ req.set_collection(kExecutionsCollection);
|
|
|
|
|
+ req.set_id(execution.id);
|
|
|
|
|
+ *req.mutable_data() = DocumentService::JsonToMapValue(ExecutionToJson(execution));
|
|
|
|
|
+ doc_service->CreateDocument(&ctx, req, &resp);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Emit execution start event
|
|
|
|
|
+ if (event_manager_) {
|
|
|
|
|
+ event_manager_->EmitScriptExecutionEvent(EventAction::Create, execution.id,
|
|
|
|
|
+ execution.workspace_id, ExecutionToJson(execution),
|
|
|
|
|
+ execution.executed_by);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Set up JavaScript engine
|
|
|
|
|
+ smartbotic::runner::JsEngineConfig js_config;
|
|
|
|
|
+ js_config.max_memory = script.config.max_memory_mb * 1024 * 1024;
|
|
|
|
|
+ js_config.max_execution_time_ms = script.config.timeout_ms;
|
|
|
|
|
+ js_config.max_db_queries = 100;
|
|
|
|
|
+ js_config.allow_eval = false;
|
|
|
|
|
+
|
|
|
|
|
+ smartbotic::runner::JsEngine engine(js_config);
|
|
|
|
|
+
|
|
|
|
|
+ // Set up context
|
|
|
|
|
+ smartbotic::runner::JsContext js_ctx;
|
|
|
|
|
+ js_ctx.workspace_id = script.workspace_id;
|
|
|
|
|
+ js_ctx.user_id = execution.executed_by;
|
|
|
|
|
+
|
|
|
|
|
+ // Add event data to context if available
|
|
|
|
|
+ if (!request.trigger_event.is_null()) {
|
|
|
|
|
+ js_ctx.custom["event"] = request.trigger_event.dump();
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ engine.SetContext(js_ctx);
|
|
|
|
|
+
|
|
|
|
|
+ // Set up database callbacks with permission checks
|
|
|
|
|
+ const std::string creator_id = script.created_by;
|
|
|
|
|
+ const std::string ws_id = script.workspace_id;
|
|
|
|
|
+ const std::vector<std::string> allowed_collections = script.config.allowed_collections;
|
|
|
|
|
+
|
|
|
|
|
+ engine.SetDbCallbacks(
|
|
|
|
|
+ // Query callback
|
|
|
|
|
+ [this, creator_id, ws_id, allowed_collections](const std::string& collection,
|
|
|
|
|
+ const std::string& filter_json,
|
|
|
|
|
+ int limit, int offset) -> std::string {
|
|
|
|
|
+ // Check if collection is allowed
|
|
|
|
|
+ if (!allowed_collections.empty()) {
|
|
|
|
|
+ auto it = std::find(allowed_collections.begin(), allowed_collections.end(), collection);
|
|
|
|
|
+ if (it == allowed_collections.end()) {
|
|
|
|
|
+ spdlog::warn("ScriptService: Access denied to collection {}", collection);
|
|
|
|
|
+ return "[]";
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Check creator's permission
|
|
|
|
|
+ if (!CreatorCanAccessCollection(creator_id, ws_id, collection, "read_all")) {
|
|
|
|
|
+ spdlog::warn("ScriptService: Creator {} lacks read permission for {}", creator_id, collection);
|
|
|
|
|
+ return "[]";
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto* query_service = db_client_.GetQueryService();
|
|
|
|
|
+ if (query_service == nullptr) { return "[]"; }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::QueryRequest req;
|
|
|
|
|
+ ::smartbotic::database::QueryResponse resp;
|
|
|
|
|
+ req.set_collection(collection);
|
|
|
|
|
+ req.set_limit(limit);
|
|
|
|
|
+ req.set_offset(offset);
|
|
|
|
|
+
|
|
|
|
|
+ // Parse filter JSON and build structured filter
|
|
|
|
|
+ if (!filter_json.empty() && filter_json != "{}") {
|
|
|
|
|
+ auto filter_obj = nlohmann::json::parse(filter_json, nullptr, false);
|
|
|
|
|
+ if (!filter_obj.is_discarded() && filter_obj.is_object() && !filter_obj.empty()) {
|
|
|
|
|
+ auto* filter = req.mutable_filter();
|
|
|
|
|
+ auto* composite = filter->mutable_composite();
|
|
|
|
|
+ composite->set_operator_(::smartbotic::database::COMPOSITE_OPERATOR_AND);
|
|
|
|
|
+
|
|
|
|
|
+ for (auto& [key, value] : filter_obj.items()) {
|
|
|
|
|
+ auto* field_filter = composite->add_filters()->mutable_field();
|
|
|
|
|
+ field_filter->set_field(key);
|
|
|
|
|
+ field_filter->set_operator_(::smartbotic::database::FILTER_OPERATOR_EQUAL);
|
|
|
|
|
+ if (value.is_string()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_string_value(value.get<std::string>());
|
|
|
|
|
+ } else if (value.is_number_integer()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_int_value(value.get<int64_t>());
|
|
|
|
|
+ } else if (value.is_boolean()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_bool_value(value.get<bool>());
|
|
|
|
|
+ } else if (value.is_number_float()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_double_value(value.get<double>());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto status = query_service->Query(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) { return "[]"; }
|
|
|
|
|
+
|
|
|
|
|
+ nlohmann::json results = nlohmann::json::array();
|
|
|
|
|
+ for (const auto& doc : resp.documents()) {
|
|
|
|
|
+ auto json = DocumentService::MapValueToJson(doc.data());
|
|
|
|
|
+ json["_id"] = doc.id();
|
|
|
|
|
+ results.push_back(json);
|
|
|
|
|
+ }
|
|
|
|
|
+ return results.dump();
|
|
|
|
|
+ },
|
|
|
|
|
+ // Get callback
|
|
|
|
|
+ [this, creator_id, ws_id, allowed_collections](const std::string& collection,
|
|
|
|
|
+ const std::string& id) -> std::string {
|
|
|
|
|
+ if (!allowed_collections.empty()) {
|
|
|
|
|
+ auto it = std::find(allowed_collections.begin(), allowed_collections.end(), collection);
|
|
|
|
|
+ if (it == allowed_collections.end()) {
|
|
|
|
|
+ return "";
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!CreatorCanAccessCollection(creator_id, ws_id, collection, "read_all")) {
|
|
|
|
|
+ return "";
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto* doc_service = db_client_.GetDocumentService();
|
|
|
|
|
+ if (doc_service == nullptr) { return ""; }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::GetDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+ req.set_collection(collection);
|
|
|
|
|
+ req.set_id(id);
|
|
|
|
|
+
|
|
|
|
|
+ auto status = doc_service->GetDocument(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) { return ""; }
|
|
|
|
|
+
|
|
|
|
|
+ auto json = DocumentService::MapValueToJson(resp.data());
|
|
|
|
|
+ json["_id"] = resp.id();
|
|
|
|
|
+ return json.dump();
|
|
|
|
|
+ },
|
|
|
|
|
+ // Count callback
|
|
|
|
|
+ [this, creator_id, ws_id, allowed_collections](const std::string& collection,
|
|
|
|
|
+ const std::string& filter_json) -> int64_t {
|
|
|
|
|
+ if (!allowed_collections.empty()) {
|
|
|
|
|
+ auto it = std::find(allowed_collections.begin(), allowed_collections.end(), collection);
|
|
|
|
|
+ if (it == allowed_collections.end()) {
|
|
|
|
|
+ return 0;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!CreatorCanAccessCollection(creator_id, ws_id, collection, "read_all")) {
|
|
|
|
|
+ return 0;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto* query_service = db_client_.GetQueryService();
|
|
|
|
|
+ if (query_service == nullptr) { return 0; }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::CountRequest req;
|
|
|
|
|
+ ::smartbotic::database::CountResponse resp;
|
|
|
|
|
+ req.set_collection(collection);
|
|
|
|
|
+
|
|
|
|
|
+ // Parse filter JSON and build structured filter
|
|
|
|
|
+ if (!filter_json.empty() && filter_json != "{}") {
|
|
|
|
|
+ auto filter_obj = nlohmann::json::parse(filter_json, nullptr, false);
|
|
|
|
|
+ if (!filter_obj.is_discarded() && filter_obj.is_object() && !filter_obj.empty()) {
|
|
|
|
|
+ auto* filter = req.mutable_filter();
|
|
|
|
|
+ auto* composite = filter->mutable_composite();
|
|
|
|
|
+ composite->set_operator_(::smartbotic::database::COMPOSITE_OPERATOR_AND);
|
|
|
|
|
+
|
|
|
|
|
+ for (auto& [key, value] : filter_obj.items()) {
|
|
|
|
|
+ auto* field_filter = composite->add_filters()->mutable_field();
|
|
|
|
|
+ field_filter->set_field(key);
|
|
|
|
|
+ field_filter->set_operator_(::smartbotic::database::FILTER_OPERATOR_EQUAL);
|
|
|
|
|
+ if (value.is_string()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_string_value(value.get<std::string>());
|
|
|
|
|
+ } else if (value.is_number_integer()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_int_value(value.get<int64_t>());
|
|
|
|
|
+ } else if (value.is_boolean()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_bool_value(value.get<bool>());
|
|
|
|
|
+ } else if (value.is_number_float()) {
|
|
|
|
|
+ field_filter->mutable_value()->set_double_value(value.get<double>());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto status = query_service->Count(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) { return 0; }
|
|
|
|
|
+ return resp.count();
|
|
|
|
|
+ }
|
|
|
|
|
+ );
|
|
|
|
|
+
|
|
|
|
|
+ // Execute the script
|
|
|
|
|
+ auto start_time = std::chrono::steady_clock::now();
|
|
|
|
|
+ auto result = engine.Execute(script.code);
|
|
|
|
|
+ auto end_time = std::chrono::steady_clock::now();
|
|
|
|
|
+
|
|
|
|
|
+ execution.execution_time_ms = std::chrono::duration_cast<std::chrono::milliseconds>(
|
|
|
|
|
+ end_time - start_time).count();
|
|
|
|
|
+
|
|
|
|
|
+ if (result.success) {
|
|
|
|
|
+ execution.status = ExecutionStatus::Success;
|
|
|
|
|
+ execution.result = result.result;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ if (result.error.find("timeout") != std::string::npos ||
|
|
|
|
|
+ result.error.find("interrupt") != std::string::npos) {
|
|
|
|
|
+ execution.status = ExecutionStatus::Timeout;
|
|
|
|
|
+ } else {
|
|
|
|
|
+ execution.status = ExecutionStatus::Error;
|
|
|
|
|
+ }
|
|
|
|
|
+ execution.error = result.error;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ execution.completed_at = GetTimestamp();
|
|
|
|
|
+
|
|
|
|
|
+ // Update execution record
|
|
|
|
|
+ if (doc_service != nullptr) {
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+ req.set_collection(kExecutionsCollection);
|
|
|
|
|
+ req.set_id(execution.id);
|
|
|
|
|
+ *req.mutable_data() = DocumentService::JsonToMapValue(ExecutionToJson(execution));
|
|
|
|
|
+ doc_service->UpdateDocument(&ctx, req, &resp);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Emit execution complete event
|
|
|
|
|
+ if (event_manager_) {
|
|
|
|
|
+ event_manager_->EmitScriptExecutionEvent(EventAction::Update, execution.id,
|
|
|
|
|
+ execution.workspace_id, ExecutionToJson(execution),
|
|
|
|
|
+ execution.executed_by);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("ScriptService: Executed script {} with status {}", script.id,
|
|
|
|
|
+ ExecStatusToString(execution.status));
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .execution = execution};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::GetExecution(const std::string& workspace_id,
|
|
|
|
|
+ const std::string& execution_id) -> ScriptExecutionResult {
|
|
|
|
|
+ auto* doc_service = db_client_.GetDocumentService();
|
|
|
|
|
+ if (doc_service == nullptr) {
|
|
|
|
|
+ return {.success = false, .error = "Document service not available", .execution = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::GetDocumentRequest req;
|
|
|
|
|
+ ::smartbotic::database::Document resp;
|
|
|
|
|
+
|
|
|
|
|
+ req.set_collection(kExecutionsCollection);
|
|
|
|
|
+ req.set_id(execution_id);
|
|
|
|
|
+
|
|
|
|
|
+ auto status = doc_service->GetDocument(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ if (status.error_code() == grpc::StatusCode::NOT_FOUND) {
|
|
|
|
|
+ return {.success = false, .error = "Execution not found", .execution = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+ return {.success = false, .error = status.error_message(), .execution = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto json = DocumentService::MapValueToJson(resp.data());
|
|
|
|
|
+ json["id"] = resp.id();
|
|
|
|
|
+
|
|
|
|
|
+ auto execution = JsonToExecution(json);
|
|
|
|
|
+
|
|
|
|
|
+ if (execution.workspace_id != workspace_id) {
|
|
|
|
|
+ return {.success = false, .error = "Execution not found", .execution = std::nullopt};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .execution = execution};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::ListExecutions(const std::string& workspace_id,
|
|
|
|
|
+ const std::string& script_id,
|
|
|
|
|
+ int limit, int offset) -> ExecutionListResult {
|
|
|
|
|
+ auto* query_service = db_client_.GetQueryService();
|
|
|
|
|
+ if (query_service == nullptr) {
|
|
|
|
|
+ return {.success = false, .error = "Query service not available", .executions = {}, .total_count = 0};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::QueryRequest req;
|
|
|
|
|
+ ::smartbotic::database::QueryResponse resp;
|
|
|
|
|
+
|
|
|
|
|
+ req.set_collection(kExecutionsCollection);
|
|
|
|
|
+ req.set_limit(limit);
|
|
|
|
|
+ req.set_offset(offset);
|
|
|
|
|
+
|
|
|
|
|
+ // Build composite filter for workspace_id and script_id
|
|
|
|
|
+ auto* filter = req.mutable_filter();
|
|
|
|
|
+ auto* composite = filter->mutable_composite();
|
|
|
|
|
+ composite->set_operator_(::smartbotic::database::COMPOSITE_OPERATOR_AND);
|
|
|
|
|
+
|
|
|
|
|
+ auto* ws_filter = composite->add_filters()->mutable_field();
|
|
|
|
|
+ ws_filter->set_field("workspace_id");
|
|
|
|
|
+ ws_filter->set_operator_(::smartbotic::database::FILTER_OPERATOR_EQUAL);
|
|
|
|
|
+ ws_filter->mutable_value()->set_string_value(workspace_id);
|
|
|
|
|
+
|
|
|
|
|
+ auto* script_filter = composite->add_filters()->mutable_field();
|
|
|
|
|
+ script_filter->set_field("script_id");
|
|
|
|
|
+ script_filter->set_operator_(::smartbotic::database::FILTER_OPERATOR_EQUAL);
|
|
|
|
|
+ script_filter->mutable_value()->set_string_value(script_id);
|
|
|
|
|
+
|
|
|
|
|
+ auto status = query_service->Query(&ctx, req, &resp);
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ return {.success = false, .error = status.error_message(), .executions = {}, .total_count = 0};
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ std::vector<ScriptExecution> executions;
|
|
|
|
|
+ for (const auto& doc : resp.documents()) {
|
|
|
|
|
+ auto json = DocumentService::MapValueToJson(doc.data());
|
|
|
|
|
+ json["id"] = doc.id();
|
|
|
|
|
+ executions.push_back(JsonToExecution(json));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Sort by started_at descending
|
|
|
|
|
+ std::sort(executions.begin(), executions.end(), [](const auto& a, const auto& b) {
|
|
|
|
|
+ return a.started_at > b.started_at;
|
|
|
|
|
+ });
|
|
|
|
|
+
|
|
|
|
|
+ return {.success = true, .error = "", .executions = executions, .total_count = resp.total_count()};
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// === Event Handling ===
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::OnEntityEvent(const EntityEvent& event) {
|
|
|
|
|
+ // Only handle Document events for now
|
|
|
|
|
+ if (event.entity_type != EntityType::Document) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto scripts = FindTriggeredScripts(event);
|
|
|
|
|
+ for (const auto& script : scripts) {
|
|
|
|
|
+ // Execute script asynchronously
|
|
|
|
|
+ ExecuteScriptRequest exec_req;
|
|
|
|
|
+ exec_req.workspace_id = script.workspace_id;
|
|
|
|
|
+ exec_req.script_id = script.id;
|
|
|
|
|
+ exec_req.executed_by = script.created_by;
|
|
|
|
|
+ exec_req.trigger_type = "event";
|
|
|
|
|
+ exec_req.trigger_event = {
|
|
|
|
|
+ {"entity_type", EntityTypeToString(event.entity_type)},
|
|
|
|
|
+ {"action", EventActionToString(event.action)},
|
|
|
|
|
+ {"entity_id", event.entity_id},
|
|
|
|
|
+ {"collection", event.collection_name},
|
|
|
|
|
+ {"data", event.data}
|
|
|
|
|
+ };
|
|
|
|
|
+
|
|
|
|
|
+ // Execute in background thread to not block event processing
|
|
|
|
|
+ std::thread([this, exec_req]() {
|
|
|
|
|
+ (void)ExecuteScript(exec_req); // Discard result intentionally
|
|
|
|
|
+ }).detach();
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::FindTriggeredScripts(const EntityEvent& event) -> std::vector<ScriptInfo> {
|
|
|
|
|
+ std::vector<ScriptInfo> triggered;
|
|
|
|
|
+
|
|
|
|
|
+ // List all active scripts in the workspace
|
|
|
|
|
+ auto scripts_result = ListScripts(event.workspace_id, false);
|
|
|
|
|
+ if (!scripts_result.success) {
|
|
|
|
|
+ return triggered;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ std::string event_action = EventActionToString(event.action);
|
|
|
|
|
+ std::string entity_type = EntityTypeToString(event.entity_type);
|
|
|
|
|
+
|
|
|
|
|
+ for (const auto& script : scripts_result.scripts) {
|
|
|
|
|
+ if (script.status != ScriptStatus::Active) {
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Check triggers
|
|
|
|
|
+ for (const auto& trigger : script.triggers) {
|
|
|
|
|
+ if (trigger.type != TriggerType::Event) {
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Match entity type
|
|
|
|
|
+ if (!trigger.entity_type.empty() && trigger.entity_type != entity_type) {
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Match action
|
|
|
|
|
+ if (!trigger.event_action.empty() && trigger.event_action != event_action) {
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Match collection filter
|
|
|
|
|
+ if (!trigger.collection_filter.empty() && trigger.collection_filter != event.collection_name) {
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // All conditions match
|
|
|
|
|
+ triggered.push_back(script);
|
|
|
|
|
+ break; // Don't trigger same script multiple times for same event
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return triggered;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// === Cron Scheduler ===
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::SchedulerLoop() {
|
|
|
|
|
+ while (scheduler_running_.load()) {
|
|
|
|
|
+ CheckCronTriggers();
|
|
|
|
|
+ std::this_thread::sleep_for(std::chrono::seconds(60));
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::UpdateCronSchedule(const ScriptInfo& script) {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(schedule_mutex_);
|
|
|
|
|
+
|
|
|
|
|
+ // Remove existing entries for this script
|
|
|
|
|
+ cron_schedules_.erase(
|
|
|
|
|
+ std::remove_if(cron_schedules_.begin(), cron_schedules_.end(),
|
|
|
|
|
+ [&script](const auto& entry) { return entry.script_id == script.id; }),
|
|
|
|
|
+ cron_schedules_.end());
|
|
|
|
|
+
|
|
|
|
|
+ // Only add if script is active
|
|
|
|
|
+ if (script.status != ScriptStatus::Active) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Add new cron triggers
|
|
|
|
|
+ for (const auto& trigger : script.triggers) {
|
|
|
|
|
+ if (trigger.type != TriggerType::Cron || trigger.cron_expression.empty()) {
|
|
|
|
|
+ continue;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ CronScheduleEntry entry;
|
|
|
|
|
+ entry.script_id = script.id;
|
|
|
|
|
+ entry.workspace_id = script.workspace_id;
|
|
|
|
|
+ entry.cron_expression = trigger.cron_expression;
|
|
|
|
|
+ entry.timezone = trigger.timezone.empty() ? "UTC" : trigger.timezone;
|
|
|
|
|
+ entry.trigger_id = trigger.id;
|
|
|
|
|
+ entry.next_run = CalculateNextRun(entry.cron_expression, entry.timezone);
|
|
|
|
|
+
|
|
|
|
|
+ cron_schedules_.push_back(entry);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::RemoveCronSchedule(const std::string& script_id) {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(schedule_mutex_);
|
|
|
|
|
+ cron_schedules_.erase(
|
|
|
|
|
+ std::remove_if(cron_schedules_.begin(), cron_schedules_.end(),
|
|
|
|
|
+ [&script_id](const auto& entry) { return entry.script_id == script_id; }),
|
|
|
|
|
+ cron_schedules_.end());
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void ScriptService::CheckCronTriggers() {
|
|
|
|
|
+ auto now = std::chrono::system_clock::now();
|
|
|
|
|
+ std::vector<CronScheduleEntry> due_entries;
|
|
|
|
|
+
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(schedule_mutex_);
|
|
|
|
|
+ for (auto& entry : cron_schedules_) {
|
|
|
|
|
+ if (entry.next_run <= now) {
|
|
|
|
|
+ due_entries.push_back(entry);
|
|
|
|
|
+ // Update next run time
|
|
|
|
|
+ entry.next_run = CalculateNextRun(entry.cron_expression, entry.timezone);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Execute due scripts
|
|
|
|
|
+ for (const auto& entry : due_entries) {
|
|
|
|
|
+ ExecuteScriptRequest exec_req;
|
|
|
|
|
+ exec_req.workspace_id = entry.workspace_id;
|
|
|
|
|
+ exec_req.script_id = entry.script_id;
|
|
|
|
|
+ exec_req.trigger_type = "cron";
|
|
|
|
|
+ exec_req.trigger_event = {
|
|
|
|
|
+ {"trigger_id", entry.trigger_id},
|
|
|
|
|
+ {"cron_expression", entry.cron_expression},
|
|
|
|
|
+ {"scheduled_time", std::chrono::system_clock::to_time_t(entry.next_run)}
|
|
|
|
|
+ };
|
|
|
|
|
+
|
|
|
|
|
+ // Execute asynchronously
|
|
|
|
|
+ std::thread([this, exec_req]() {
|
|
|
|
|
+ (void)ExecuteScript(exec_req); // Discard result intentionally
|
|
|
|
|
+ }).detach();
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::CalculateNextRun(const std::string& cron_expression,
|
|
|
|
|
+ const std::string& /*timezone*/) -> std::chrono::system_clock::time_point {
|
|
|
|
|
+ // Simple cron parser for basic expressions: minute hour day month weekday
|
|
|
|
|
+ // For now, just add 1 minute for simplicity - a full cron parser would be more complex
|
|
|
|
|
+ // TODO: Implement full cron expression parsing
|
|
|
|
|
+ (void)cron_expression;
|
|
|
|
|
+ return std::chrono::system_clock::now() + std::chrono::minutes(1);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::CreatorCanAccessCollection(const std::string& creator_id,
|
|
|
|
|
+ const std::string& workspace_id,
|
|
|
|
|
+ const std::string& collection,
|
|
|
|
|
+ const std::string& action) -> bool {
|
|
|
|
|
+ if (auth_service_ == nullptr) {
|
|
|
|
|
+ // No auth service - allow access
|
|
|
|
|
+ return true;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Build an AuthUser for the creator
|
|
|
|
|
+ AuthUser creator;
|
|
|
|
|
+ creator.user_id = creator_id;
|
|
|
|
|
+
|
|
|
|
|
+ // Build the permission to check
|
|
|
|
|
+ auto permission = permissions::BuildCollectionPermission(workspace_id, collection, action);
|
|
|
|
|
+
|
|
|
|
|
+ return auth_service_->HasPermission(creator, permission);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// === JSON Conversion ===
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::ScriptToJson(const ScriptInfo& script) -> nlohmann::json {
|
|
|
|
|
+ nlohmann::json triggers_json = nlohmann::json::array();
|
|
|
|
|
+ for (const auto& trigger : script.triggers) {
|
|
|
|
|
+ triggers_json.push_back(TriggerToJson(trigger));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return {
|
|
|
|
|
+ {"id", script.id},
|
|
|
|
|
+ {"workspace_id", script.workspace_id},
|
|
|
|
|
+ {"name", script.name},
|
|
|
|
|
+ {"description", script.description},
|
|
|
|
|
+ {"code", script.code},
|
|
|
|
|
+ {"status", StatusToString(script.status)},
|
|
|
|
|
+ {"config", ConfigToJson(script.config)},
|
|
|
|
|
+ {"triggers", triggers_json},
|
|
|
|
|
+ {"created_at", script.created_at},
|
|
|
|
|
+ {"updated_at", script.updated_at},
|
|
|
|
|
+ {"created_by", script.created_by},
|
|
|
|
|
+ {"deleted_at", script.deleted_at}
|
|
|
|
|
+ };
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::JsonToScript(const nlohmann::json& json) -> ScriptInfo {
|
|
|
|
|
+ ScriptInfo script;
|
|
|
|
|
+ script.id = json.value("id", "");
|
|
|
|
|
+ script.workspace_id = json.value("workspace_id", "");
|
|
|
|
|
+ script.name = json.value("name", "");
|
|
|
|
|
+ script.description = json.value("description", "");
|
|
|
|
|
+ script.code = json.value("code", "");
|
|
|
|
|
+ script.status = StringToStatus(json.value("status", "draft"));
|
|
|
|
|
+ script.created_at = json.value("created_at", "");
|
|
|
|
|
+ script.updated_at = json.value("updated_at", "");
|
|
|
|
|
+ script.created_by = json.value("created_by", "");
|
|
|
|
|
+ script.deleted_at = json.value("deleted_at", "");
|
|
|
|
|
+
|
|
|
|
|
+ if (json.contains("config") && json["config"].is_object()) {
|
|
|
|
|
+ script.config = JsonToConfig(json["config"]);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (json.contains("triggers") && json["triggers"].is_array()) {
|
|
|
|
|
+ for (const auto& t : json["triggers"]) {
|
|
|
|
|
+ script.triggers.push_back(JsonToTrigger(t));
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return script;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::ExecutionToJson(const ScriptExecution& execution) -> nlohmann::json {
|
|
|
|
|
+ return {
|
|
|
|
|
+ {"id", execution.id},
|
|
|
|
|
+ {"script_id", execution.script_id},
|
|
|
|
|
+ {"workspace_id", execution.workspace_id},
|
|
|
|
|
+ {"trigger_type", execution.trigger_type},
|
|
|
|
|
+ {"trigger_event", execution.trigger_event},
|
|
|
|
|
+ {"status", ExecStatusToString(execution.status)},
|
|
|
|
|
+ {"result", execution.result},
|
|
|
|
|
+ {"error", execution.error},
|
|
|
|
|
+ {"logs", execution.logs},
|
|
|
|
|
+ {"started_at", execution.started_at},
|
|
|
|
|
+ {"completed_at", execution.completed_at},
|
|
|
|
|
+ {"execution_time_ms", execution.execution_time_ms},
|
|
|
|
|
+ {"executed_by", execution.executed_by}
|
|
|
|
|
+ };
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::JsonToExecution(const nlohmann::json& json) -> ScriptExecution {
|
|
|
|
|
+ ScriptExecution execution;
|
|
|
|
|
+ execution.id = json.value("id", "");
|
|
|
|
|
+ execution.script_id = json.value("script_id", "");
|
|
|
|
|
+ execution.workspace_id = json.value("workspace_id", "");
|
|
|
|
|
+ execution.trigger_type = json.value("trigger_type", "");
|
|
|
|
|
+ execution.trigger_event = json.value("trigger_event", nlohmann::json::object());
|
|
|
|
|
+ execution.status = StringToExecStatus(json.value("status", "running"));
|
|
|
|
|
+ execution.result = json.value("result", "");
|
|
|
|
|
+ execution.error = json.value("error", "");
|
|
|
|
|
+ execution.started_at = json.value("started_at", "");
|
|
|
|
|
+ execution.completed_at = json.value("completed_at", "");
|
|
|
|
|
+ execution.execution_time_ms = json.value("execution_time_ms", 0);
|
|
|
|
|
+ execution.executed_by = json.value("executed_by", "");
|
|
|
|
|
+
|
|
|
|
|
+ if (json.contains("logs") && json["logs"].is_array()) {
|
|
|
|
|
+ for (const auto& log : json["logs"]) {
|
|
|
|
|
+ execution.logs.push_back(log.get<std::string>());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return execution;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::TriggerToJson(const ScriptTrigger& trigger) -> nlohmann::json {
|
|
|
|
|
+ return {
|
|
|
|
|
+ {"id", trigger.id},
|
|
|
|
|
+ {"type", TriggerTypeToString(trigger.type)},
|
|
|
|
|
+ {"entity_type", trigger.entity_type},
|
|
|
|
|
+ {"event_action", trigger.event_action},
|
|
|
|
|
+ {"collection_filter", trigger.collection_filter},
|
|
|
|
|
+ {"cron_expression", trigger.cron_expression},
|
|
|
|
|
+ {"timezone", trigger.timezone},
|
|
|
|
|
+ {"last_triggered", trigger.last_triggered},
|
|
|
|
|
+ {"trigger_count", trigger.trigger_count}
|
|
|
|
|
+ };
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::JsonToTrigger(const nlohmann::json& json) -> ScriptTrigger {
|
|
|
|
|
+ ScriptTrigger trigger;
|
|
|
|
|
+ trigger.id = json.value("id", "");
|
|
|
|
|
+ trigger.type = StringToTriggerType(json.value("type", "event"));
|
|
|
|
|
+ trigger.entity_type = json.value("entity_type", "");
|
|
|
|
|
+ trigger.event_action = json.value("event_action", "");
|
|
|
|
|
+ trigger.collection_filter = json.value("collection_filter", "");
|
|
|
|
|
+ trigger.cron_expression = json.value("cron_expression", "");
|
|
|
|
|
+ trigger.timezone = json.value("timezone", "UTC");
|
|
|
|
|
+ trigger.last_triggered = json.value("last_triggered", "");
|
|
|
|
|
+ trigger.trigger_count = json.value("trigger_count", 0);
|
|
|
|
|
+ return trigger;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::ConfigToJson(const ScriptConfig& config) -> nlohmann::json {
|
|
|
|
|
+ return {
|
|
|
|
|
+ {"timeout_ms", config.timeout_ms},
|
|
|
|
|
+ {"max_memory_mb", config.max_memory_mb},
|
|
|
|
|
+ {"allow_network", config.allow_network},
|
|
|
|
|
+ {"allow_storage", config.allow_storage},
|
|
|
|
|
+ {"allowed_collections", config.allowed_collections}
|
|
|
|
|
+ };
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::JsonToConfig(const nlohmann::json& json) -> ScriptConfig {
|
|
|
|
|
+ ScriptConfig config;
|
|
|
|
|
+ config.timeout_ms = json.value("timeout_ms", 5000);
|
|
|
|
|
+ config.max_memory_mb = json.value("max_memory_mb", 16);
|
|
|
|
|
+ config.allow_network = json.value("allow_network", false);
|
|
|
|
|
+ config.allow_storage = json.value("allow_storage", true);
|
|
|
|
|
+
|
|
|
|
|
+ if (json.contains("allowed_collections") && json["allowed_collections"].is_array()) {
|
|
|
|
|
+ for (const auto& c : json["allowed_collections"]) {
|
|
|
|
|
+ config.allowed_collections.push_back(c.get<std::string>());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return config;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::StatusToString(ScriptStatus status) -> std::string {
|
|
|
|
|
+ switch (status) {
|
|
|
|
|
+ case ScriptStatus::Active: return "active";
|
|
|
|
|
+ case ScriptStatus::Inactive: return "inactive";
|
|
|
|
|
+ case ScriptStatus::Draft: return "draft";
|
|
|
|
|
+ default: return "draft";
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::StringToStatus(const std::string& str) -> ScriptStatus {
|
|
|
|
|
+ if (str == "active") { return ScriptStatus::Active; }
|
|
|
|
|
+ if (str == "inactive") { return ScriptStatus::Inactive; }
|
|
|
|
|
+ return ScriptStatus::Draft;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::ExecStatusToString(ExecutionStatus status) -> std::string {
|
|
|
|
|
+ switch (status) {
|
|
|
|
|
+ case ExecutionStatus::Running: return "running";
|
|
|
|
|
+ case ExecutionStatus::Success: return "success";
|
|
|
|
|
+ case ExecutionStatus::Error: return "error";
|
|
|
|
|
+ case ExecutionStatus::Timeout: return "timeout";
|
|
|
|
|
+ default: return "running";
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::StringToExecStatus(const std::string& str) -> ExecutionStatus {
|
|
|
|
|
+ if (str == "success") { return ExecutionStatus::Success; }
|
|
|
|
|
+ if (str == "error") { return ExecutionStatus::Error; }
|
|
|
|
|
+ if (str == "timeout") { return ExecutionStatus::Timeout; }
|
|
|
|
|
+ return ExecutionStatus::Running;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::TriggerTypeToString(TriggerType type) -> std::string {
|
|
|
|
|
+ switch (type) {
|
|
|
|
|
+ case TriggerType::Event: return "event";
|
|
|
|
|
+ case TriggerType::Cron: return "cron";
|
|
|
|
|
+ default: return "event";
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::StringToTriggerType(const std::string& str) -> TriggerType {
|
|
|
|
|
+ if (str == "cron") { return TriggerType::Cron; }
|
|
|
|
|
+ return TriggerType::Event;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::GetCurrentTimestamp() -> std::string {
|
|
|
|
|
+ return GetTimestamp();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto ScriptService::GenerateUuid() -> std::string {
|
|
|
|
|
+ return GenerateUuidV4();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+} // namespace smartbotic::webserver
|