#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 "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_); // 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: {} new nodes imported", migrate_result.value().size()); } 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(); 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); }); // Initialize HTTP server HttpServerConfig http_config; http_config.port = config_.http_port; http_config.static_files_path = config_.static_files_path; 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"); // 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(); // 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); user_ctrl_ = std::make_unique(*auth_store_, *auth_middleware_); user_ctrl_->registerRoutes(server); workflow_ctrl_ = std::make_unique( *storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_, *scheduler_, *node_store_); workflow_ctrl_->registerRoutes(server); workflow_group_ctrl_ = std::make_unique( *storage_, *auth_middleware_, *ws_server_); workflow_group_ctrl_->registerRoutes(server); execution_ctrl_ = std::make_unique( *storage_, *auth_middleware_, *ws_server_, *scheduler_, *load_balancer_, [this](const std::string& workflow_id, const std::string& execution_id, const std::string& error) { runErrorWorkflow(workflow_id, execution_id, 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_, &node_sync_server_->service()); node_ctrl_->registerRoutes(server); runner_ctrl_ = std::make_unique(*runner_registry_, *auth_middleware_); runner_ctrl_->registerRoutes(server); webhook_ctrl_ = std::make_unique( *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_); webhook_ctrl_->registerRoutes(server); database_ctrl_ = std::make_unique(*storage_, *auth_middleware_); database_ctrl_->registerRoutes(server); credential_ctrl_ = std::make_unique(*credential_store_, *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 || !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::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"); 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) { // 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); nlohmann::json trigger_data; trigger_data["triggerNodeId"] = trigger_node_id; trigger_data["scheduledExecution"] = true; 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