#include "webserver_service.hpp" #include "api/auth_controller.hpp" #include "api/user_controller.hpp" #include "api/workflow_controller.hpp" #include "api/workflow_group_controller.hpp" #include "api/execution_controller.hpp" #include "api/node_controller.hpp" #include "api/runner_controller.hpp" #include "api/webhook_controller.hpp" #include "api/database_controller.hpp" #include "api/file_controller.hpp" #include "api/proxy_controller.hpp" #include "api/credential_controller.hpp" #include "nodes/node_store.hpp" #include "grpc/node_sync_service.hpp" #include "grpc/credential_service.hpp" #include "credentials/credential_store.hpp" #include "scheduler/workflow_scheduler.hpp" #include "common/time_utils.hpp" #include #include "common/config_defaults.hpp" #include "logging/logger.hpp" #include #include "proto/runner.grpc.pb.h" namespace smartbotic::webserver { WebServerService::WebServerService(const WebServerServiceConfig& config) : config_(config) { // Initialize storage client storage::StorageClientConfig storage_config; storage_config.address = config_.database_address; storage_config.project = config_.database_project; storage_ = std::make_unique(storage_config); // Initialize JWT jwt_ = std::make_unique(config_.jwt_config); // Initialize auth store auth_store_ = std::make_unique(*storage_, *jwt_); // The indexes this installation's queries actually need. Declared here, and // idempotently, so a fresh install or a restored backup gets them without // anyone remembering to - they live in the database, not in the config. // // Only what was measured to help. An index costs write throughput, so a // declaration that buys nothing is a real cost paid for ever: // // sessions.refreshToken every token refresh scanned the whole // collection - 24 ms over 9,661 sessions, 1 ms now // executions.workflowId a lookup that missed took 1,124 ms and now takes // 1. What is left when it hits is the size of the // execution documents, about 30 ms each, which no // index can help with // workflows.projectId small today; the listing filters by it on every // page load and workflows are written by hand // users.username/email login looks up by both // // executions.startedAt was declared, dropped, and declared again. On 2.9 it // changed nothing: ranges and sorts were not served from an index, and a // sorted single row still took 422 ms over 10,000 rows. Re-measured on // 2.10, which fixed that, it makes no measurable difference either - but the // collection is 1,500 rows now rather than 10,000, having had its orphans // removed, so the two measurements are not comparable and neither is a // verdict. It is kept because the field is perfectly selective (every value // distinct), the planner declines an index it cannot use, and the collection // will grow again. for (const auto& [collection, field] : std::initializer_list>{ {"executions", "workflowId"}, {"executions", "startedAt"}, {"workflows", "projectId"}, {"users", "username"}, {"users", "email"}, {"sessions", "refreshToken"}}) { auto created = storage_->createIndex(collection, field); if (created.failed()) { // Not fatal: an older database has no index support and every query // still works, just by scanning. LOG_DEBUG("Index {}.{} not declared: {}", collection, field, created.error().message()); } else if (created.value() > 0) { LOG_INFO("Index {}.{} declared, {} rows backfilled", collection, field, created.value()); } } // Sessions issued under a longer lifetime than is configured now would // otherwise keep it until they ran out, so shortening the setting would not // take effect for as long as the old one lasted. auth_store_->enforceSessionLifetime(); // Initialize auth middleware auth_middleware_ = std::make_unique(*jwt_, *auth_store_); // Initialize runner registry runner_registry_ = std::make_unique(*storage_, config_.runner_config); // Initialize load balancer load_balancer_ = std::make_unique(*runner_registry_, config_.load_balancer_config); // Initialize node store node_store_ = std::make_unique(*storage_); // Migrate/sync nodes from disk - updates existing nodes if code has changed auto migrate_result = node_store_->migrateFromFiles("./nodes"); if (migrate_result.ok()) { LOG_INFO("Node migration: {} nodes imported/updated, {} rejected", migrate_result.value().migrated.size(), migrate_result.value().rejected.size()); for (const auto& rejection : migrate_result.value().rejected) { for (const auto& reason : rejection.reasons) { LOG_WARN("Node migration rejected {}: {}", rejection.node_id, reason); } } } else { LOG_WARN("Node migration failed: {}", migrate_result.error().message()); } // Initialize credential store credentials::CredentialStoreConfig cred_config; cred_config.master_key = config_.credentials_config.master_key; cred_config.pbkdf2_iterations = config_.credentials_config.pbkdf2_iterations; credential_store_ = std::make_unique(*storage_, cred_config); credential_store_->initialize(); // Initialize workflow scheduler scheduler_ = std::make_unique(); // The scheduler holds a slot per run and frees it when the completion event // arrives. Events get lost - a database outage mid-run is enough - so it can // also ask whether a run has finished. Answering from the execution record // rather than from memory turns a 75-minute stall into one tick. scheduler_->setExecutionFinishedCheck([this](const std::string& execution_id) { auto record = storage_->get("executions", execution_id); if (record.failed()) { // Unreadable is not finished. Saying otherwise here would start a // second run of a workflow that is still going, which is worse than // waiting for the deadline the scheduler already has. return false; } const std::string status = record.value().value("status", ""); return status == "completed" || status == "failed" || status == "cancelled"; }); scheduler_->setExecuteCallback([this](const std::string& workflow_id, const std::string& trigger_node_id, const std::string& trigger_type) { executeScheduledWorkflow(workflow_id, trigger_node_id, trigger_type); }); // Triggers that are told rather than asking. The database says when a // collection changed, so these workflows do not poll for it. db_watcher_ = std::make_unique(*storage_); db_watcher_->setExecuteCallback([this](const std::string& workflow_id, const std::string& trigger_node_id, const std::string& trigger_type, const nlohmann::json& event) { executeScheduledWorkflow(workflow_id, trigger_node_id, trigger_type, event); }); // Initialize HTTP server HttpServerConfig http_config; http_config.port = config_.http_port; http_config.static_files_path = config_.static_files_path; http_config.max_upload_mb = config_.max_upload_mb; http_server_ = std::make_unique(http_config); // Initialize WebSocket server on separate port WebSocketServerConfig ws_config; ws_config.port = config_.http_port + 1; // WebSocket on next port (8091) ws_server_ = std::make_unique(ws_config, *jwt_); // Initialize NodeSync gRPC server for runners node_sync_server_ = std::make_unique(*node_store_, config_.node_sync_port); // Initialize Credential gRPC server for runners credential_server_ = std::make_unique(*credential_store_, config_.credential_service_port); } WebServerService::~WebServerService() { stop(); } WebServerServiceConfig WebServerService::loadConfig(const std::filesystem::path& path) { WebServerServiceConfig config; auto result = config::Config::fromFile(path); if (result.ok()) { auto& cfg = result.value(); config.http_port = cfg.getOr("http_port", 8080); config.node_sync_port = cfg.getOr("node_sync_port", 9002); config.credential_service_port = cfg.getOr("credential_service_port", 9003); config.static_files_path = cfg.getOr("static_files_path", "./webui/dist"); config.database_address = cfg.getOr("database_address", "localhost:9004"); config.database_project = cfg.getOr("database_project", "smartbotic-automation"); // HttpServer already enforced a body cap before this - a hardcoded // 16 MB literal at the point it constructs httplib::Server - so this // is not closing an open door, it is making that cap configurable // and raising it. The default goes from 16 MB to 32 MB because an // upload node needs to carry up to a 16 MB image field, and that // field cannot fit inside a 16 MB total body once multipart framing // is added on top. config.max_upload_mb = cfg.getOr("server.max_upload_mb", 32); // Runner config config.runner_config.heartbeat_timeout_sec = cfg.getOr("runners.heartbeat_timeout_sec", 30); config.runner_config.offline_removal_sec = cfg.getOr("runners.offline_removal_sec", 60); // Load balancer config auto strategy_str = cfg.getOr("runners.load_balancing", "least-connections"); config.load_balancer_config.strategy = runners::loadBalancingStrategyFromString(strategy_str); // JWT config config.jwt_config.secret = cfg.getOr("auth.jwt_secret", "dev-secret-change-in-production"); config.jwt_config.access_token_lifetime_sec = cfg.getOr("auth.access_token_lifetime_sec", 900); // Credentials config config.credentials_config.master_key = cfg.getOr("credentials.master_key", "dev-key-change-in-production"); config.credentials_config.pbkdf2_iterations = cfg.getOr("credentials.pbkdf2_iterations", 100000); } return config; } void WebServerService::start() { LOG_INFO("Starting WebServer service..."); // Ensure admin user exists auth_store_->ensureAdminUser(); // Setup API routes setupRoutes(); // Start runner registry cleanup runner_registry_->start(); // Start workflow scheduler scheduler_->start(); // Keep a history of workflow documents. Without it the version number // still climbs on every write but nothing older is retained, so there is no // earlier version to publish, compare against or go back to. // // Turned on here rather than at creation because the collection already // exists on every install that predates this, and createCollection refuses // a collection that is already there. { storage::CollectionConfig cfg; cfg.versioning_enabled = true; auto configured = storage_->configureCollection("workflows", cfg); if (configured.failed()) { LOG_WARN("Could not enable version history on workflows: {}", configured.error().message()); } else { LOG_INFO("Version history enabled on workflows"); } } // Everybody gets a personal project, and anything stored before projects // existed is filed into one - a record with no project is one only an // instance admin can see. if (project_ctrl_) { project_ctrl_->migrateExistingRecords(); } // Somebody has to be the owner - the account that cannot be demoted or // deleted, so an installation can never end up with nobody able to // administer it. If nobody is, the first admin becomes it. { auto users = auth_store_->listUsers(1, 1000); if (users.ok()) { bool have_owner = false; for (const auto& user : users.value()) { if (user.role == "owner") { have_owner = true; break; } } if (!have_owner) { for (const auto& user : users.value()) { if (user.role != "admin") continue; auto promoted = storage_->update("users", user.id, {{"role", "owner"}}, 0, true); if (promoted.ok()) { LOG_INFO("{} is now the owner of this installation", user.username); } else { LOG_WARN("Could not make {} the owner: {}", user.username, promoted.error().message()); } break; } } } } // Ensure the executions summary view exists before serving requests ensureExecutionsSummaryView(); // Load active workflows from database and register with scheduler loadScheduledWorkflows(); // Start NodeSync gRPC server for runners node_sync_server_->start(); // Start Credential gRPC server for runners credential_server_->start(); // Start WebSocket server ws_server_->start(); // Start HTTP server (blocking in background thread) http_server_->start(); LOG_INFO("WebServer service started on port {}", config_.http_port); } void WebServerService::stop() { LOG_INFO("Stopping WebServer service..."); http_server_->stop(); ws_server_->stop(); credential_server_->stop(); node_sync_server_->stop(); scheduler_->stop(); runner_registry_->stop(); LOG_INFO("WebServer service stopped"); } void WebServerService::setupRoutes() { auto& server = http_server_->server(); // Health check server.Get("/health", [](const httplib::Request& req, httplib::Response& res) { res.set_content(R"({"status":"ok"})", "application/json"); }); // API version server.Get("/api/v1", [](const httplib::Request& req, httplib::Response& res) { res.set_content(R"({"name":"SmartBotic API","version":"1.0.0"})", "application/json"); }); // Register controllers - stored as members to ensure they outlive httplib callbacks auth_ctrl_ = std::make_unique(*auth_store_, *auth_middleware_); auth_ctrl_->registerRoutes(server); // Who may do what. Built before any controller that asks it - dereferencing // this while it was still null bound a reference to nothing, and the first // request that used it took the server down with a segfault rather than // failing anywhere near the mistake. access_ = std::make_unique(*storage_); user_ctrl_ = std::make_unique( *auth_store_, *auth_middleware_, *access_, *storage_); user_ctrl_->registerRoutes(server); // Before every controller that holds a reference to it. Dereferencing an // empty unique_ptr here does not fail here - it fails later, inside the // first request that uses it, as a segfault nowhere near the mistake. That // has happened once on this branch already. retention_ = std::make_unique(*storage_); project_ctrl_ = std::make_unique( *storage_, *auth_middleware_, *access_, *auth_store_, *retention_); project_ctrl_->registerRoutes(server); workflow_ctrl_ = std::make_unique( *storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_, *scheduler_, *db_watcher_, *access_, *node_store_, *retention_); workflow_ctrl_->registerRoutes(server); workflow_group_ctrl_ = std::make_unique( *storage_, *access_, *auth_middleware_, *ws_server_); workflow_group_ctrl_->registerRoutes(server); execution_ctrl_ = std::make_unique( *storage_, *auth_middleware_, *access_, *ws_server_, *scheduler_, *load_balancer_, [this](const std::string& workflow_id, const std::string& execution_id, bool failed, const std::string& error) { noteExecutionOutcome(workflow_id, execution_id, failed, error); }); file_ctrl_ = std::make_unique(*storage_, *auth_middleware_); execution_ctrl_->registerRoutes(server); file_ctrl_->registerRoutes(server); proxy_ctrl_ = std::make_unique(*auth_middleware_); proxy_ctrl_->registerRoutes(server); node_ctrl_ = std::make_unique( *node_store_, *auth_middleware_, *load_balancer_, &node_sync_server_->service()); node_ctrl_->registerRoutes(server); runner_ctrl_ = std::make_unique( *runner_registry_, *auth_middleware_, [this](const std::string& runner_id, const std::string& address) { reconcileOrphanedExecutions(runner_id, address); }); runner_ctrl_->registerRoutes(server); webhook_ctrl_ = std::make_unique( *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_, *scheduler_); webhook_ctrl_->registerRoutes(server); database_ctrl_ = std::make_unique(*storage_, *auth_middleware_); database_ctrl_->registerRoutes(server); credential_ctrl_ = std::make_unique( *credential_store_, *access_, *storage_, *auth_middleware_); credential_ctrl_->registerRoutes(server); LOG_INFO("API routes registered"); } void WebServerService::ensureExecutionsSummaryView() { // View names are namespaced per project from smartbotic-database 2.4.2, so a // view created by this client is also queryable by it. Earlier versions // registered the name unqualified while queries were project-prefixed, which // made the view unreachable. // Recreate rather than reuse. A view is metadata over a collection, so // rebuilding it costs nothing, and keeping an existing one means the view // silently outlives the field list below: after a database recovery the old // view served documents that were no longer in the collection at all. for (const auto& view : storage_->listViews()) { if (view.name == kExecutionsSummaryView) { LOG_INFO("Rebuilding the executions summary view"); storage_->dropView(kExecutionsSummaryView); break; } } // stopReason belongs here for the same reason error does: it is the message // explaining how a run ended, and the listing is where someone looks for it. // A stopped run has status "completed" and an empty error, so without this // field the list cannot tell a run that finished its work from one that // deliberately ended early. auto result = storage_->createView(kExecutionsSummaryView, "executions", {"workflowId", "workflowName", "status", "triggerType", "startedAt", "finishedAt", "error", "runnerId", "stopped", "stopReason", "stoppedNodeId"}); if (result.failed()) { LOG_ERROR("Could not create the executions summary view: {}. The executions " "listing will fail until this is resolved.", result.error().message()); return; } LOG_INFO("Created server-side view '{}' over executions", kExecutionsSummaryView); } void WebServerService::loadScheduledWorkflows() { LOG_INFO("Loading scheduled workflows from database..."); // Query all active workflows storage::QueryOptions options; options.filters.push_back({"active", true}); options.page_size = 1000; // Load up to 1000 workflows auto result = storage_->query("workflows", options); if (result.failed()) { LOG_ERROR("Failed to load workflows: {}", result.error().message()); return; } int registered_count = 0; for (const auto& workflow : result.value().documents) { std::string workflow_id = workflow.value("_id", ""); std::string workflow_name = workflow.value("name", ""); auto nodes = workflow.value("nodes", nlohmann::json::array()); // Find scheduled trigger nodes for (const auto& node : nodes) { std::string node_id = node.value("id", ""); std::string node_type = node.value("type", ""); // Get node definition to check if it's a scheduled trigger auto node_result = node_store_->get(node_type); if (node_result.failed()) { continue; } const auto& node_def = node_result.value(); if (!node_def.is_trigger) { continue; } // Told, not asked: the database reports the change, so there is no // interval. Registered here as well as on activate, or a restart // would leave every event-driven workflow deaf until somebody // toggled it. if (node_type == "database-change") { auto config = smartbotic::common::applyConfigDefaults( node.value("config", nlohmann::json::object()), node_def.config_schema); const std::string collection = config.value("collection", ""); if (collection.empty()) continue; std::vector event_types; if (config.contains("eventTypes") && config["eventTypes"].is_array()) { for (const auto& t : config["eventTypes"]) { if (t.is_string()) event_types.push_back(t.get()); } } db_watcher_->watch(workflow_id, workflow_name, node_id, collection, event_types); registered_count++; continue; } if (!node_def.is_scheduled) { continue; } // Get interval from node config, filling in any schema defaults // the stored config is missing (e.g. an untouched form field). auto config = smartbotic::common::applyConfigDefaults( node.value("config", nlohmann::json::object()), node_def.config_schema); int interval = config.value("pollInterval", 0); if (interval > 0) { auto policy = overlapPolicyFromString( config.value("overlapPolicy", std::string("skip"))); scheduler_->registerWorkflow( workflow_id, workflow_name, node_id, node_type, interval, policy, config.value("maxConcurrent", 1), config.value("maxRunMinutes", 0) ); registered_count++; LOG_DEBUG("Registered workflow '{}' ({}) for scheduled execution every {} minutes", workflow_name, workflow_id, interval); } } } LOG_INFO("Loaded {} scheduled workflows", registered_count); } void WebServerService::noteExecutionOutcome(const std::string& workflow_id, const std::string& execution_id, bool failed, const std::string& error) { if (failed) { runErrorWorkflow(workflow_id, execution_id, error); } auto stored = storage_->get("workflows", workflow_id); if (stored.failed()) { return; } const auto workflow = stored.value(); const auto settings = workflow.value("settings", nlohmann::json::object()); const int limit = settings.value("deactivateAfterFailures", 0); const int streak = workflow.value("consecutiveFailures", 0); if (!failed) { // Only written when there is something to clear, so an ordinary run does // not cost a write. if (streak != 0) { storage_->update("workflows", workflow_id, {{"consecutiveFailures", 0}}, 0, true); } return; } const int next = streak + 1; nlohmann::json patch = {{"consecutiveFailures", next}}; // Counting is worth doing even when nothing is switched off - it is the // number someone looks at when asking how long this has been going wrong. if (limit <= 0 || next < limit || workflow.value("active", false) != true) { storage_->update("workflows", workflow_id, patch, 0, true); return; } // Switched off, and said out loud. A workflow that simply appeared // "Inactive" one morning with no reason recorded is worse than one that // kept failing, because nobody can tell which of the two happened. patch["active"] = false; patch["deactivatedReason"] = "Switched off after " + std::to_string(next) + " failures in a row. The last one: " + (error.empty() ? "no reason given" : error); patch["deactivatedAt"] = common::TimeUtils::nowMs(); patch["updatedAt"] = common::TimeUtils::nowMs(); storage_->update("workflows", workflow_id, patch, 0, true); // Off the schedule too, or it would keep firing while reading as inactive. scheduler_->unregisterWorkflow(workflow_id); LOG_WARN("Workflow {} deactivated after {} consecutive failures: {}", workflow_id, next, error); ws_server_->broadcast("workflows.deactivated", { {"id", workflow_id}, {"reason", patch["deactivatedReason"]}, {"consecutiveFailures", next} }); } void WebServerService::reconcileOrphanedExecutions(const std::string& runner_id, const std::string& address) { // Ask the runner what it is actually running rather than assuming. A runner // that re-registers while working - a duplicate call, a flapping network - // must not have its live executions closed underneath it, and only the // runner knows which those are. std::unordered_set still_running; { auto channel = ::grpc::CreateChannel(address, ::grpc::InsecureChannelCredentials()); auto stub = proto::RunnerService::NewStub(channel); proto::ListActiveExecutionsRequest request; proto::ListActiveExecutionsResponse response; ::grpc::ClientContext context; context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(10)); auto status = stub->ListActiveExecutions(&context, request, &response); if (!status.ok()) { // Without an answer there is no way to tell an orphan from a live // run, and closing a live one is far worse than leaving a stale // record for someone to press Stop on. LOG_WARN("Runner {} could not say what it is running ({}), so nothing was reconciled", runner_id, status.error_message()); return; } for (const auto& id : response.execution_ids()) { still_running.insert(id); } } storage::QueryOptions options; options.filters.push_back({"runnerId", runner_id}); options.filters.push_back({"status", "running"}); options.page = 1; options.page_size = 500; auto found = storage_->query("executions", options); if (found.failed()) { LOG_WARN("Could not look for orphaned executions on runner {}: {}", runner_id, found.error().message()); return; } int closed = 0; for (const auto& record : found.value().documents) { const std::string id = record.value("_id", ""); if (id.empty() || still_running.contains(id)) { continue; } // Deliberately not touching Waiting. That is a run parked on a person, // not on a runner, and it is meant to outlive one. const nlohmann::json patch = { {"status", "cancelled"}, {"error", "The runner restarted while this was running, so nothing was left to finish it"}, {"finishedAt", common::TimeUtils::nowMs()} }; if (storage_->update("executions", id, patch, 0, true).ok()) { ++closed; ws_server_->broadcast("executions." + id + ".cancelled", {{"executionId", id}}); } } if (closed > 0) { LOG_WARN("Closed {} execution(s) that runner {} was recorded as running but is not", closed, runner_id); } } void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id, const std::string& failed_execution_id, const std::string& error_message) { { std::lock_guard lock(handled_failures_mutex_); if (!handled_failures_.insert(failed_execution_id).second) { return; // already handled this failure } if (handled_failures_.size() > 512) { handled_failures_.erase(handled_failures_.begin()); } } auto failed = storage_->get("workflows", failed_workflow_id); if (failed.failed()) { return; } const auto settings = failed.value().value("settings", nlohmann::json::object()); const std::string handler_id = settings.value("errorWorkflowId", std::string()); if (handler_id.empty()) { return; } // A handler that fails must not summon itself, which would run forever. if (handler_id == failed_workflow_id) { LOG_WARN("Workflow {} names itself as its error workflow; not running it", failed_workflow_id); return; } auto handler = storage_->get("workflows", handler_id); if (handler.failed()) { LOG_WARN("Workflow {} names error workflow {}, which no longer exists", failed_workflow_id, handler_id); return; } auto runner = load_balancer_->selectRunner(); if (!runner) { LOG_ERROR("No runners available to run error workflow {}", handler_id); return; } // The handler is told what failed rather than having to look it up, so it can // notify or record without needing read access to the executions collection. nlohmann::json trigger_data; trigger_data["errorWorkflow"] = true; trigger_data["failedWorkflowId"] = failed_workflow_id; trigger_data["failedWorkflowName"] = failed.value().value("name", std::string()); trigger_data["failedExecutionId"] = failed_execution_id; trigger_data["error"] = error_message; trigger_data["failedAt"] = common::TimeUtils::nowMs(); auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials()); auto stub = proto::RunnerService::NewStub(channel); proto::ExecuteWorkflowRequest request; request.set_workflow_id(handler_id); request.set_trigger_type("error-workflow"); // Started by the system, so it runs what was published. request.set_use_published(true); request.set_trigger_data(trigger_data.dump()); request.set_wait_for_completion(false); proto::ExecuteWorkflowResponse response; ::grpc::ClientContext context; context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30)); auto status = stub->ExecuteWorkflow(&context, request, &response); if (!status.ok()) { LOG_ERROR("Error workflow {} could not be started: {}", handler_id, status.error_message()); return; } LOG_INFO("Error workflow {} started as {} after {} failed", handler_id, response.execution_id(), failed_workflow_id); scheduler_->notifyExecutionStarted(handler_id, response.execution_id()); } void WebServerService::executeScheduledWorkflow(const std::string& workflow_id, const std::string& trigger_node_id, const std::string& trigger_type, const nlohmann::json& extra_trigger_data) { // Select a runner auto runner = load_balancer_->selectRunner(); if (!runner) { LOG_ERROR("No runners available for scheduled workflow {}", workflow_id); return; } // Create gRPC channel and stub auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials()); auto stub = proto::RunnerService::NewStub(channel); // Prepare request proto::ExecuteWorkflowRequest request; request.set_workflow_id(workflow_id); request.set_trigger_type(trigger_type); // A schedule firing is not somebody testing an edit: it runs the published // version, so a half-finished change on somebody's canvas never goes live // just because it was saved. request.set_use_published(true); nlohmann::json trigger_data; trigger_data["triggerNodeId"] = trigger_node_id; trigger_data["scheduledExecution"] = true; // What actually happened, for the triggers that are told rather than the // ones that ask - a database change carries the document with it. if (extra_trigger_data.is_object()) { for (const auto& [key, value] : extra_trigger_data.items()) { trigger_data[key] = value; } } request.set_trigger_data(trigger_data.dump()); request.set_wait_for_completion(false); // Execute proto::ExecuteWorkflowResponse response; ::grpc::ClientContext context; context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30)); auto status = stub->ExecuteWorkflow(&context, request, &response); if (!status.ok()) { LOG_ERROR("Failed to execute scheduled workflow {}: {}", workflow_id, status.error_message()); return; } LOG_INFO("Scheduled workflow {} execution started: {} on runner {}", workflow_id, response.execution_id(), runner->id); // Dispatch is fire-and-forget, so the scheduler only learns about the run here. scheduler_->notifyExecutionStarted(workflow_id, response.execution_id()); // Broadcast execution started ws_server_->broadcast("executions." + response.execution_id() + ".started", { {"executionId", response.execution_id()}, {"workflowId", workflow_id}, {"runnerId", runner->id}, {"triggeredBy", trigger_type}, {"scheduled", true} }); } } // namespace smartbotic::webserver