webserver_service.cpp 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485
  1. #include "webserver_service.hpp"
  2. #include "api/auth_controller.hpp"
  3. #include "api/user_controller.hpp"
  4. #include "api/workflow_controller.hpp"
  5. #include "api/workflow_group_controller.hpp"
  6. #include "api/execution_controller.hpp"
  7. #include "api/node_controller.hpp"
  8. #include "api/runner_controller.hpp"
  9. #include "api/webhook_controller.hpp"
  10. #include "api/database_controller.hpp"
  11. #include "api/file_controller.hpp"
  12. #include "api/proxy_controller.hpp"
  13. #include "api/credential_controller.hpp"
  14. #include "nodes/node_store.hpp"
  15. #include "grpc/node_sync_service.hpp"
  16. #include "grpc/credential_service.hpp"
  17. #include "credentials/credential_store.hpp"
  18. #include "scheduler/workflow_scheduler.hpp"
  19. #include "common/time_utils.hpp"
  20. #include "common/config_defaults.hpp"
  21. #include "logging/logger.hpp"
  22. #include <grpcpp/grpcpp.h>
  23. #include "proto/runner.grpc.pb.h"
  24. namespace smartbotic::webserver {
  25. WebServerService::WebServerService(const WebServerServiceConfig& config)
  26. : config_(config) {
  27. // Initialize storage client
  28. storage::StorageClientConfig storage_config;
  29. storage_config.address = config_.database_address;
  30. storage_config.project = config_.database_project;
  31. storage_ = std::make_unique<storage::StorageClient>(storage_config);
  32. // Initialize JWT
  33. jwt_ = std::make_unique<auth::JwtUtils>(config_.jwt_config);
  34. // Initialize auth store
  35. auth_store_ = std::make_unique<auth::AuthStore>(*storage_, *jwt_);
  36. // Initialize auth middleware
  37. auth_middleware_ = std::make_unique<auth::AuthMiddleware>(*jwt_, *auth_store_);
  38. // Initialize runner registry
  39. runner_registry_ = std::make_unique<runners::RunnerRegistry>(*storage_, config_.runner_config);
  40. // Initialize load balancer
  41. load_balancer_ = std::make_unique<runners::LoadBalancer>(*runner_registry_,
  42. config_.load_balancer_config);
  43. // Initialize node store
  44. node_store_ = std::make_unique<nodes::NodeStore>(*storage_);
  45. // Migrate/sync nodes from disk - updates existing nodes if code has changed
  46. auto migrate_result = node_store_->migrateFromFiles("./nodes");
  47. if (migrate_result.ok()) {
  48. LOG_INFO("Node migration: {} new nodes imported", migrate_result.value().size());
  49. } else {
  50. LOG_WARN("Node migration failed: {}", migrate_result.error().message());
  51. }
  52. // Initialize credential store
  53. credentials::CredentialStoreConfig cred_config;
  54. cred_config.master_key = config_.credentials_config.master_key;
  55. cred_config.pbkdf2_iterations = config_.credentials_config.pbkdf2_iterations;
  56. credential_store_ = std::make_unique<credentials::CredentialStore>(*storage_, cred_config);
  57. credential_store_->initialize();
  58. // Initialize workflow scheduler
  59. scheduler_ = std::make_unique<WorkflowScheduler>();
  60. scheduler_->setExecuteCallback([this](const std::string& workflow_id,
  61. const std::string& trigger_node_id,
  62. const std::string& trigger_type) {
  63. executeScheduledWorkflow(workflow_id, trigger_node_id, trigger_type);
  64. });
  65. // Initialize HTTP server
  66. HttpServerConfig http_config;
  67. http_config.port = config_.http_port;
  68. http_config.static_files_path = config_.static_files_path;
  69. http_server_ = std::make_unique<HttpServer>(http_config);
  70. // Initialize WebSocket server on separate port
  71. WebSocketServerConfig ws_config;
  72. ws_config.port = config_.http_port + 1; // WebSocket on next port (8091)
  73. ws_server_ = std::make_unique<WebSocketServer>(ws_config, *jwt_);
  74. // Initialize NodeSync gRPC server for runners
  75. node_sync_server_ = std::make_unique<grpc::NodeSyncServer>(*node_store_, config_.node_sync_port);
  76. // Initialize Credential gRPC server for runners
  77. credential_server_ = std::make_unique<grpc::CredentialServer>(*credential_store_, config_.credential_service_port);
  78. }
  79. WebServerService::~WebServerService() {
  80. stop();
  81. }
  82. WebServerServiceConfig WebServerService::loadConfig(const std::filesystem::path& path) {
  83. WebServerServiceConfig config;
  84. auto result = config::Config::fromFile(path);
  85. if (result.ok()) {
  86. auto& cfg = result.value();
  87. config.http_port = cfg.getOr<int>("http_port", 8080);
  88. config.node_sync_port = cfg.getOr<int>("node_sync_port", 9002);
  89. config.credential_service_port = cfg.getOr<int>("credential_service_port", 9003);
  90. config.static_files_path = cfg.getOr<std::string>("static_files_path", "./webui/dist");
  91. config.database_address = cfg.getOr<std::string>("database_address", "localhost:9004");
  92. config.database_project = cfg.getOr<std::string>("database_project", "smartbotic-automation");
  93. // Runner config
  94. config.runner_config.heartbeat_timeout_sec =
  95. cfg.getOr<int>("runners.heartbeat_timeout_sec", 30);
  96. config.runner_config.offline_removal_sec =
  97. cfg.getOr<int>("runners.offline_removal_sec", 60);
  98. // Load balancer config
  99. auto strategy_str = cfg.getOr<std::string>("runners.load_balancing", "least-connections");
  100. config.load_balancer_config.strategy =
  101. runners::loadBalancingStrategyFromString(strategy_str);
  102. // JWT config
  103. config.jwt_config.secret = cfg.getOr<std::string>("auth.jwt_secret",
  104. "dev-secret-change-in-production");
  105. config.jwt_config.access_token_lifetime_sec =
  106. cfg.getOr<int>("auth.access_token_lifetime_sec", 900);
  107. // Credentials config
  108. config.credentials_config.master_key = cfg.getOr<std::string>("credentials.master_key",
  109. "dev-key-change-in-production");
  110. config.credentials_config.pbkdf2_iterations =
  111. cfg.getOr<int>("credentials.pbkdf2_iterations", 100000);
  112. }
  113. return config;
  114. }
  115. void WebServerService::start() {
  116. LOG_INFO("Starting WebServer service...");
  117. // Ensure admin user exists
  118. auth_store_->ensureAdminUser();
  119. // Setup API routes
  120. setupRoutes();
  121. // Start runner registry cleanup
  122. runner_registry_->start();
  123. // Start workflow scheduler
  124. scheduler_->start();
  125. // Ensure the executions summary view exists before serving requests
  126. ensureExecutionsSummaryView();
  127. // Load active workflows from database and register with scheduler
  128. loadScheduledWorkflows();
  129. // Start NodeSync gRPC server for runners
  130. node_sync_server_->start();
  131. // Start Credential gRPC server for runners
  132. credential_server_->start();
  133. // Start WebSocket server
  134. ws_server_->start();
  135. // Start HTTP server (blocking in background thread)
  136. http_server_->start();
  137. LOG_INFO("WebServer service started on port {}", config_.http_port);
  138. }
  139. void WebServerService::stop() {
  140. LOG_INFO("Stopping WebServer service...");
  141. http_server_->stop();
  142. ws_server_->stop();
  143. credential_server_->stop();
  144. node_sync_server_->stop();
  145. scheduler_->stop();
  146. runner_registry_->stop();
  147. LOG_INFO("WebServer service stopped");
  148. }
  149. void WebServerService::setupRoutes() {
  150. auto& server = http_server_->server();
  151. // Health check
  152. server.Get("/health", [](const httplib::Request& req, httplib::Response& res) {
  153. res.set_content(R"({"status":"ok"})", "application/json");
  154. });
  155. // API version
  156. server.Get("/api/v1", [](const httplib::Request& req, httplib::Response& res) {
  157. res.set_content(R"({"name":"SmartBotic API","version":"1.0.0"})", "application/json");
  158. });
  159. // Register controllers - stored as members to ensure they outlive httplib callbacks
  160. auth_ctrl_ = std::make_unique<api::AuthController>(*auth_store_, *auth_middleware_);
  161. auth_ctrl_->registerRoutes(server);
  162. user_ctrl_ = std::make_unique<api::UserController>(*auth_store_, *auth_middleware_);
  163. user_ctrl_->registerRoutes(server);
  164. workflow_ctrl_ = std::make_unique<api::WorkflowController>(
  165. *storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_,
  166. *scheduler_, *node_store_);
  167. workflow_ctrl_->registerRoutes(server);
  168. workflow_group_ctrl_ = std::make_unique<api::WorkflowGroupController>(
  169. *storage_, *auth_middleware_, *ws_server_);
  170. workflow_group_ctrl_->registerRoutes(server);
  171. execution_ctrl_ = std::make_unique<api::ExecutionController>(
  172. *storage_, *auth_middleware_, *ws_server_, *scheduler_, *load_balancer_,
  173. [this](const std::string& workflow_id, const std::string& execution_id,
  174. const std::string& error) {
  175. runErrorWorkflow(workflow_id, execution_id, error);
  176. });
  177. file_ctrl_ = std::make_unique<api::FileController>(*storage_, *auth_middleware_);
  178. execution_ctrl_->registerRoutes(server);
  179. file_ctrl_->registerRoutes(server);
  180. proxy_ctrl_ = std::make_unique<api::ProxyController>(*auth_middleware_);
  181. proxy_ctrl_->registerRoutes(server);
  182. node_ctrl_ = std::make_unique<api::NodeController>(
  183. *node_store_, *auth_middleware_, &node_sync_server_->service());
  184. node_ctrl_->registerRoutes(server);
  185. runner_ctrl_ = std::make_unique<api::RunnerController>(*runner_registry_, *auth_middleware_);
  186. runner_ctrl_->registerRoutes(server);
  187. webhook_ctrl_ = std::make_unique<api::WebhookController>(
  188. *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_);
  189. webhook_ctrl_->registerRoutes(server);
  190. database_ctrl_ = std::make_unique<api::DatabaseController>(*storage_, *auth_middleware_);
  191. database_ctrl_->registerRoutes(server);
  192. credential_ctrl_ = std::make_unique<api::CredentialController>(*credential_store_, *auth_middleware_);
  193. credential_ctrl_->registerRoutes(server);
  194. LOG_INFO("API routes registered");
  195. }
  196. void WebServerService::ensureExecutionsSummaryView() {
  197. // View names are namespaced per project from smartbotic-database 2.4.2, so a
  198. // view created by this client is also queryable by it. Earlier versions
  199. // registered the name unqualified while queries were project-prefixed, which
  200. // made the view unreachable.
  201. // Recreate rather than reuse. A view is metadata over a collection, so
  202. // rebuilding it costs nothing, and keeping an existing one means the view
  203. // silently outlives the field list below: after a database recovery the old
  204. // view served documents that were no longer in the collection at all.
  205. for (const auto& view : storage_->listViews()) {
  206. if (view.name == kExecutionsSummaryView) {
  207. LOG_INFO("Rebuilding the executions summary view");
  208. storage_->dropView(kExecutionsSummaryView);
  209. break;
  210. }
  211. }
  212. // stopReason belongs here for the same reason error does: it is the message
  213. // explaining how a run ended, and the listing is where someone looks for it.
  214. // A stopped run has status "completed" and an empty error, so without this
  215. // field the list cannot tell a run that finished its work from one that
  216. // deliberately ended early.
  217. auto result = storage_->createView(kExecutionsSummaryView, "executions",
  218. {"workflowId", "workflowName", "status", "triggerType",
  219. "startedAt", "finishedAt", "error", "runnerId",
  220. "stopped", "stopReason", "stoppedNodeId"});
  221. if (result.failed()) {
  222. LOG_ERROR("Could not create the executions summary view: {}. The executions "
  223. "listing will fail until this is resolved.", result.error().message());
  224. return;
  225. }
  226. LOG_INFO("Created server-side view '{}' over executions", kExecutionsSummaryView);
  227. }
  228. void WebServerService::loadScheduledWorkflows() {
  229. LOG_INFO("Loading scheduled workflows from database...");
  230. // Query all active workflows
  231. storage::QueryOptions options;
  232. options.filters.push_back({"active", true});
  233. options.page_size = 1000; // Load up to 1000 workflows
  234. auto result = storage_->query("workflows", options);
  235. if (result.failed()) {
  236. LOG_ERROR("Failed to load workflows: {}", result.error().message());
  237. return;
  238. }
  239. int registered_count = 0;
  240. for (const auto& workflow : result.value().documents) {
  241. std::string workflow_id = workflow.value("_id", "");
  242. std::string workflow_name = workflow.value("name", "");
  243. auto nodes = workflow.value("nodes", nlohmann::json::array());
  244. // Find scheduled trigger nodes
  245. for (const auto& node : nodes) {
  246. std::string node_id = node.value("id", "");
  247. std::string node_type = node.value("type", "");
  248. // Get node definition to check if it's a scheduled trigger
  249. auto node_result = node_store_->get(node_type);
  250. if (node_result.failed()) {
  251. continue;
  252. }
  253. const auto& node_def = node_result.value();
  254. if (!node_def.is_trigger || !node_def.is_scheduled) {
  255. continue;
  256. }
  257. // Get interval from node config, filling in any schema defaults
  258. // the stored config is missing (e.g. an untouched form field).
  259. auto config = smartbotic::common::applyConfigDefaults(
  260. node.value("config", nlohmann::json::object()), node_def.config_schema);
  261. int interval = config.value("pollInterval", 0);
  262. if (interval > 0) {
  263. auto policy = overlapPolicyFromString(
  264. config.value("overlapPolicy", std::string("skip")));
  265. scheduler_->registerWorkflow(
  266. workflow_id,
  267. workflow_name,
  268. node_id,
  269. node_type,
  270. interval,
  271. policy,
  272. config.value("maxConcurrent", 1),
  273. config.value("maxRunMinutes", 0)
  274. );
  275. registered_count++;
  276. LOG_DEBUG("Registered workflow '{}' ({}) for scheduled execution every {} minutes",
  277. workflow_name, workflow_id, interval);
  278. }
  279. }
  280. }
  281. LOG_INFO("Loaded {} scheduled workflows", registered_count);
  282. }
  283. void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
  284. const std::string& failed_execution_id,
  285. const std::string& error_message) {
  286. {
  287. std::lock_guard<std::mutex> lock(handled_failures_mutex_);
  288. if (!handled_failures_.insert(failed_execution_id).second) {
  289. return; // already handled this failure
  290. }
  291. if (handled_failures_.size() > 512) {
  292. handled_failures_.erase(handled_failures_.begin());
  293. }
  294. }
  295. auto failed = storage_->get("workflows", failed_workflow_id);
  296. if (failed.failed()) {
  297. return;
  298. }
  299. const auto settings = failed.value().value("settings", nlohmann::json::object());
  300. const std::string handler_id = settings.value("errorWorkflowId", std::string());
  301. if (handler_id.empty()) {
  302. return;
  303. }
  304. // A handler that fails must not summon itself, which would run forever.
  305. if (handler_id == failed_workflow_id) {
  306. LOG_WARN("Workflow {} names itself as its error workflow; not running it",
  307. failed_workflow_id);
  308. return;
  309. }
  310. auto handler = storage_->get("workflows", handler_id);
  311. if (handler.failed()) {
  312. LOG_WARN("Workflow {} names error workflow {}, which no longer exists",
  313. failed_workflow_id, handler_id);
  314. return;
  315. }
  316. auto runner = load_balancer_->selectRunner();
  317. if (!runner) {
  318. LOG_ERROR("No runners available to run error workflow {}", handler_id);
  319. return;
  320. }
  321. // The handler is told what failed rather than having to look it up, so it can
  322. // notify or record without needing read access to the executions collection.
  323. nlohmann::json trigger_data;
  324. trigger_data["errorWorkflow"] = true;
  325. trigger_data["failedWorkflowId"] = failed_workflow_id;
  326. trigger_data["failedWorkflowName"] = failed.value().value("name", std::string());
  327. trigger_data["failedExecutionId"] = failed_execution_id;
  328. trigger_data["error"] = error_message;
  329. trigger_data["failedAt"] = common::TimeUtils::nowMs();
  330. auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
  331. auto stub = proto::RunnerService::NewStub(channel);
  332. proto::ExecuteWorkflowRequest request;
  333. request.set_workflow_id(handler_id);
  334. request.set_trigger_type("error-workflow");
  335. request.set_trigger_data(trigger_data.dump());
  336. request.set_wait_for_completion(false);
  337. proto::ExecuteWorkflowResponse response;
  338. ::grpc::ClientContext context;
  339. context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
  340. auto status = stub->ExecuteWorkflow(&context, request, &response);
  341. if (!status.ok()) {
  342. LOG_ERROR("Error workflow {} could not be started: {}", handler_id, status.error_message());
  343. return;
  344. }
  345. LOG_INFO("Error workflow {} started as {} after {} failed",
  346. handler_id, response.execution_id(), failed_workflow_id);
  347. scheduler_->notifyExecutionStarted(handler_id, response.execution_id());
  348. }
  349. void WebServerService::executeScheduledWorkflow(const std::string& workflow_id,
  350. const std::string& trigger_node_id,
  351. const std::string& trigger_type) {
  352. // Select a runner
  353. auto runner = load_balancer_->selectRunner();
  354. if (!runner) {
  355. LOG_ERROR("No runners available for scheduled workflow {}", workflow_id);
  356. return;
  357. }
  358. // Create gRPC channel and stub
  359. auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
  360. auto stub = proto::RunnerService::NewStub(channel);
  361. // Prepare request
  362. proto::ExecuteWorkflowRequest request;
  363. request.set_workflow_id(workflow_id);
  364. request.set_trigger_type(trigger_type);
  365. nlohmann::json trigger_data;
  366. trigger_data["triggerNodeId"] = trigger_node_id;
  367. trigger_data["scheduledExecution"] = true;
  368. request.set_trigger_data(trigger_data.dump());
  369. request.set_wait_for_completion(false);
  370. // Execute
  371. proto::ExecuteWorkflowResponse response;
  372. ::grpc::ClientContext context;
  373. context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
  374. auto status = stub->ExecuteWorkflow(&context, request, &response);
  375. if (!status.ok()) {
  376. LOG_ERROR("Failed to execute scheduled workflow {}: {}", workflow_id, status.error_message());
  377. return;
  378. }
  379. LOG_INFO("Scheduled workflow {} execution started: {} on runner {}",
  380. workflow_id, response.execution_id(), runner->id);
  381. // Dispatch is fire-and-forget, so the scheduler only learns about the run here.
  382. scheduler_->notifyExecutionStarted(workflow_id, response.execution_id());
  383. // Broadcast execution started
  384. ws_server_->broadcast("executions." + response.execution_id() + ".started", {
  385. {"executionId", response.execution_id()},
  386. {"workflowId", workflow_id},
  387. {"runnerId", runner->id},
  388. {"triggeredBy", trigger_type},
  389. {"scheduled", true}
  390. });
  391. }
  392. } // namespace smartbotic::webserver