| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485 |
- #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 <grpcpp/grpcpp.h>
- #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::StorageClient>(storage_config);
- // Initialize JWT
- jwt_ = std::make_unique<auth::JwtUtils>(config_.jwt_config);
- // Initialize auth store
- auth_store_ = std::make_unique<auth::AuthStore>(*storage_, *jwt_);
- // Initialize auth middleware
- auth_middleware_ = std::make_unique<auth::AuthMiddleware>(*jwt_, *auth_store_);
- // Initialize runner registry
- runner_registry_ = std::make_unique<runners::RunnerRegistry>(*storage_, config_.runner_config);
- // Initialize load balancer
- load_balancer_ = std::make_unique<runners::LoadBalancer>(*runner_registry_,
- config_.load_balancer_config);
- // Initialize node store
- node_store_ = std::make_unique<nodes::NodeStore>(*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<credentials::CredentialStore>(*storage_, cred_config);
- credential_store_->initialize();
- // Initialize workflow scheduler
- scheduler_ = std::make_unique<WorkflowScheduler>();
- 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<HttpServer>(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<WebSocketServer>(ws_config, *jwt_);
- // Initialize NodeSync gRPC server for runners
- node_sync_server_ = std::make_unique<grpc::NodeSyncServer>(*node_store_, config_.node_sync_port);
- // Initialize Credential gRPC server for runners
- credential_server_ = std::make_unique<grpc::CredentialServer>(*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<int>("http_port", 8080);
- config.node_sync_port = cfg.getOr<int>("node_sync_port", 9002);
- config.credential_service_port = cfg.getOr<int>("credential_service_port", 9003);
- config.static_files_path = cfg.getOr<std::string>("static_files_path", "./webui/dist");
- config.database_address = cfg.getOr<std::string>("database_address", "localhost:9004");
- config.database_project = cfg.getOr<std::string>("database_project", "smartbotic-automation");
- // Runner config
- config.runner_config.heartbeat_timeout_sec =
- cfg.getOr<int>("runners.heartbeat_timeout_sec", 30);
- config.runner_config.offline_removal_sec =
- cfg.getOr<int>("runners.offline_removal_sec", 60);
- // Load balancer config
- auto strategy_str = cfg.getOr<std::string>("runners.load_balancing", "least-connections");
- config.load_balancer_config.strategy =
- runners::loadBalancingStrategyFromString(strategy_str);
- // JWT config
- config.jwt_config.secret = cfg.getOr<std::string>("auth.jwt_secret",
- "dev-secret-change-in-production");
- config.jwt_config.access_token_lifetime_sec =
- cfg.getOr<int>("auth.access_token_lifetime_sec", 900);
- // Credentials config
- config.credentials_config.master_key = cfg.getOr<std::string>("credentials.master_key",
- "dev-key-change-in-production");
- config.credentials_config.pbkdf2_iterations =
- cfg.getOr<int>("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<api::AuthController>(*auth_store_, *auth_middleware_);
- auth_ctrl_->registerRoutes(server);
- user_ctrl_ = std::make_unique<api::UserController>(*auth_store_, *auth_middleware_);
- user_ctrl_->registerRoutes(server);
- workflow_ctrl_ = std::make_unique<api::WorkflowController>(
- *storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_,
- *scheduler_, *node_store_);
- workflow_ctrl_->registerRoutes(server);
- workflow_group_ctrl_ = std::make_unique<api::WorkflowGroupController>(
- *storage_, *auth_middleware_, *ws_server_);
- workflow_group_ctrl_->registerRoutes(server);
- execution_ctrl_ = std::make_unique<api::ExecutionController>(
- *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<api::FileController>(*storage_, *auth_middleware_);
- execution_ctrl_->registerRoutes(server);
- file_ctrl_->registerRoutes(server);
- proxy_ctrl_ = std::make_unique<api::ProxyController>(*auth_middleware_);
- proxy_ctrl_->registerRoutes(server);
- node_ctrl_ = std::make_unique<api::NodeController>(
- *node_store_, *auth_middleware_, &node_sync_server_->service());
- node_ctrl_->registerRoutes(server);
- runner_ctrl_ = std::make_unique<api::RunnerController>(*runner_registry_, *auth_middleware_);
- runner_ctrl_->registerRoutes(server);
- webhook_ctrl_ = std::make_unique<api::WebhookController>(
- *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_);
- webhook_ctrl_->registerRoutes(server);
- database_ctrl_ = std::make_unique<api::DatabaseController>(*storage_, *auth_middleware_);
- database_ctrl_->registerRoutes(server);
- credential_ctrl_ = std::make_unique<api::CredentialController>(*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<std::mutex> 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
|