webserver_service.cpp 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822
  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 <unordered_set>
  21. #include "common/config_defaults.hpp"
  22. #include "logging/logger.hpp"
  23. #include <grpcpp/grpcpp.h>
  24. #include "proto/runner.grpc.pb.h"
  25. namespace smartbotic::webserver {
  26. WebServerService::WebServerService(const WebServerServiceConfig& config)
  27. : config_(config) {
  28. // Initialize storage client
  29. storage::StorageClientConfig storage_config;
  30. storage_config.address = config_.database_address;
  31. storage_config.project = config_.database_project;
  32. storage_ = std::make_unique<storage::StorageClient>(storage_config);
  33. // Initialize JWT
  34. jwt_ = std::make_unique<auth::JwtUtils>(config_.jwt_config);
  35. // Initialize auth store
  36. auth_store_ = std::make_unique<auth::AuthStore>(*storage_, *jwt_);
  37. // The indexes this installation's queries actually need. Declared here, and
  38. // idempotently, so a fresh install or a restored backup gets them without
  39. // anyone remembering to - they live in the database, not in the config.
  40. //
  41. // Only what was measured to help. An index costs write throughput, so a
  42. // declaration that buys nothing is a real cost paid for ever:
  43. //
  44. // sessions.refreshToken every token refresh scanned the whole
  45. // collection - 24 ms over 9,661 sessions, 1 ms now
  46. // executions.workflowId a lookup that missed took 1,124 ms and now takes
  47. // 1. What is left when it hits is the size of the
  48. // execution documents, about 30 ms each, which no
  49. // index can help with
  50. // workflows.projectId small today; the listing filters by it on every
  51. // page load and workflows are written by hand
  52. // users.username/email login looks up by both
  53. //
  54. // executions.startedAt was declared, dropped, and declared again. On 2.9 it
  55. // changed nothing: ranges and sorts were not served from an index, and a
  56. // sorted single row still took 422 ms over 10,000 rows. Re-measured on
  57. // 2.10, which fixed that, it makes no measurable difference either - but the
  58. // collection is 1,500 rows now rather than 10,000, having had its orphans
  59. // removed, so the two measurements are not comparable and neither is a
  60. // verdict. It is kept because the field is perfectly selective (every value
  61. // distinct), the planner declines an index it cannot use, and the collection
  62. // will grow again.
  63. for (const auto& [collection, field] : std::initializer_list<std::pair<const char*, const char*>>{
  64. {"executions", "workflowId"},
  65. {"executions", "startedAt"},
  66. {"workflows", "projectId"},
  67. {"users", "username"},
  68. {"users", "email"},
  69. {"sessions", "refreshToken"}}) {
  70. auto created = storage_->createIndex(collection, field);
  71. if (created.failed()) {
  72. // Not fatal: an older database has no index support and every query
  73. // still works, just by scanning.
  74. LOG_DEBUG("Index {}.{} not declared: {}", collection, field,
  75. created.error().message());
  76. } else if (created.value() > 0) {
  77. LOG_INFO("Index {}.{} declared, {} rows backfilled", collection, field,
  78. created.value());
  79. }
  80. }
  81. // Sessions issued under a longer lifetime than is configured now would
  82. // otherwise keep it until they ran out, so shortening the setting would not
  83. // take effect for as long as the old one lasted.
  84. auth_store_->enforceSessionLifetime();
  85. // Initialize auth middleware
  86. auth_middleware_ = std::make_unique<auth::AuthMiddleware>(*jwt_, *auth_store_);
  87. // Initialize runner registry
  88. runner_registry_ = std::make_unique<runners::RunnerRegistry>(*storage_, config_.runner_config);
  89. // Initialize load balancer
  90. load_balancer_ = std::make_unique<runners::LoadBalancer>(*runner_registry_,
  91. config_.load_balancer_config);
  92. // Initialize node store
  93. node_store_ = std::make_unique<nodes::NodeStore>(*storage_);
  94. // Migrate/sync nodes from disk - updates existing nodes if code has changed
  95. auto migrate_result = node_store_->migrateFromFiles("./nodes");
  96. if (migrate_result.ok()) {
  97. LOG_INFO("Node migration: {} nodes imported/updated, {} rejected",
  98. migrate_result.value().migrated.size(), migrate_result.value().rejected.size());
  99. for (const auto& rejection : migrate_result.value().rejected) {
  100. for (const auto& reason : rejection.reasons) {
  101. LOG_WARN("Node migration rejected {}: {}", rejection.node_id, reason);
  102. }
  103. }
  104. } else {
  105. LOG_WARN("Node migration failed: {}", migrate_result.error().message());
  106. }
  107. // Initialize credential store
  108. credentials::CredentialStoreConfig cred_config;
  109. cred_config.master_key = config_.credentials_config.master_key;
  110. cred_config.pbkdf2_iterations = config_.credentials_config.pbkdf2_iterations;
  111. credential_store_ = std::make_unique<credentials::CredentialStore>(*storage_, cred_config);
  112. credential_store_->initialize();
  113. // Initialize workflow scheduler
  114. scheduler_ = std::make_unique<WorkflowScheduler>();
  115. // The scheduler holds a slot per run and frees it when the completion event
  116. // arrives. Events get lost - a database outage mid-run is enough - so it can
  117. // also ask whether a run has finished. Answering from the execution record
  118. // rather than from memory turns a 75-minute stall into one tick.
  119. scheduler_->setExecutionFinishedCheck([this](const std::string& execution_id) {
  120. auto record = storage_->get("executions", execution_id);
  121. if (record.failed()) {
  122. // Unreadable is not finished. Saying otherwise here would start a
  123. // second run of a workflow that is still going, which is worse than
  124. // waiting for the deadline the scheduler already has.
  125. return false;
  126. }
  127. const std::string status = record.value().value("status", "");
  128. return status == "completed" || status == "failed" || status == "cancelled";
  129. });
  130. scheduler_->setExecuteCallback([this](const std::string& workflow_id,
  131. const std::string& trigger_node_id,
  132. const std::string& trigger_type) {
  133. executeScheduledWorkflow(workflow_id, trigger_node_id, trigger_type);
  134. });
  135. // Triggers that are told rather than asking. The database says when a
  136. // collection changed, so these workflows do not poll for it.
  137. db_watcher_ = std::make_unique<DatabaseWatcher>(*storage_);
  138. db_watcher_->setExecuteCallback([this](const std::string& workflow_id,
  139. const std::string& trigger_node_id,
  140. const std::string& trigger_type,
  141. const nlohmann::json& event) {
  142. executeScheduledWorkflow(workflow_id, trigger_node_id, trigger_type, event);
  143. });
  144. // Initialize HTTP server
  145. HttpServerConfig http_config;
  146. http_config.port = config_.http_port;
  147. http_config.static_files_path = config_.static_files_path;
  148. http_config.max_upload_mb = config_.max_upload_mb;
  149. http_server_ = std::make_unique<HttpServer>(http_config);
  150. // Initialize WebSocket server on separate port
  151. WebSocketServerConfig ws_config;
  152. ws_config.port = config_.http_port + 1; // WebSocket on next port (8091)
  153. ws_server_ = std::make_unique<WebSocketServer>(ws_config, *jwt_);
  154. // Initialize NodeSync gRPC server for runners
  155. node_sync_server_ = std::make_unique<grpc::NodeSyncServer>(*node_store_, config_.node_sync_port);
  156. // Initialize Credential gRPC server for runners
  157. credential_server_ = std::make_unique<grpc::CredentialServer>(*credential_store_, config_.credential_service_port);
  158. }
  159. WebServerService::~WebServerService() {
  160. stop();
  161. }
  162. WebServerServiceConfig WebServerService::loadConfig(const std::filesystem::path& path) {
  163. WebServerServiceConfig config;
  164. auto result = config::Config::fromFile(path);
  165. if (result.ok()) {
  166. auto& cfg = result.value();
  167. config.http_port = cfg.getOr<int>("http_port", 8080);
  168. config.node_sync_port = cfg.getOr<int>("node_sync_port", 9002);
  169. config.credential_service_port = cfg.getOr<int>("credential_service_port", 9003);
  170. config.static_files_path = cfg.getOr<std::string>("static_files_path", "./webui/dist");
  171. config.database_address = cfg.getOr<std::string>("database_address", "localhost:9004");
  172. config.database_project = cfg.getOr<std::string>("database_project", "smartbotic-automation");
  173. // HttpServer already enforced a body cap before this - a hardcoded
  174. // 16 MB literal at the point it constructs httplib::Server - so this
  175. // is not closing an open door, it is making that cap configurable
  176. // and raising it. The default goes from 16 MB to 32 MB because an
  177. // upload node needs to carry up to a 16 MB image field, and that
  178. // field cannot fit inside a 16 MB total body once multipart framing
  179. // is added on top.
  180. config.max_upload_mb = cfg.getOr<int>("server.max_upload_mb", 32);
  181. // Runner config
  182. config.runner_config.heartbeat_timeout_sec =
  183. cfg.getOr<int>("runners.heartbeat_timeout_sec", 30);
  184. config.runner_config.offline_removal_sec =
  185. cfg.getOr<int>("runners.offline_removal_sec", 60);
  186. // Load balancer config
  187. auto strategy_str = cfg.getOr<std::string>("runners.load_balancing", "least-connections");
  188. config.load_balancer_config.strategy =
  189. runners::loadBalancingStrategyFromString(strategy_str);
  190. // JWT config
  191. config.jwt_config.secret = cfg.getOr<std::string>("auth.jwt_secret",
  192. "dev-secret-change-in-production");
  193. config.jwt_config.access_token_lifetime_sec =
  194. cfg.getOr<int>("auth.access_token_lifetime_sec", 900);
  195. // Credentials config
  196. config.credentials_config.master_key = cfg.getOr<std::string>("credentials.master_key",
  197. "dev-key-change-in-production");
  198. config.credentials_config.pbkdf2_iterations =
  199. cfg.getOr<int>("credentials.pbkdf2_iterations", 100000);
  200. }
  201. return config;
  202. }
  203. void WebServerService::start() {
  204. LOG_INFO("Starting WebServer service...");
  205. // Ensure admin user exists
  206. auth_store_->ensureAdminUser();
  207. // Setup API routes
  208. setupRoutes();
  209. // Start runner registry cleanup
  210. runner_registry_->start();
  211. // Start workflow scheduler
  212. scheduler_->start();
  213. // Keep a history of workflow documents. Without it the version number
  214. // still climbs on every write but nothing older is retained, so there is no
  215. // earlier version to publish, compare against or go back to.
  216. //
  217. // Turned on here rather than at creation because the collection already
  218. // exists on every install that predates this, and createCollection refuses
  219. // a collection that is already there.
  220. {
  221. storage::CollectionConfig cfg;
  222. cfg.versioning_enabled = true;
  223. auto configured = storage_->configureCollection("workflows", cfg);
  224. if (configured.failed()) {
  225. LOG_WARN("Could not enable version history on workflows: {}",
  226. configured.error().message());
  227. } else {
  228. LOG_INFO("Version history enabled on workflows");
  229. }
  230. }
  231. // Everybody gets a personal project, and anything stored before projects
  232. // existed is filed into one - a record with no project is one only an
  233. // instance admin can see.
  234. if (project_ctrl_) {
  235. project_ctrl_->migrateExistingRecords();
  236. }
  237. // Somebody has to be the owner - the account that cannot be demoted or
  238. // deleted, so an installation can never end up with nobody able to
  239. // administer it. If nobody is, the first admin becomes it.
  240. {
  241. auto users = auth_store_->listUsers(1, 1000);
  242. if (users.ok()) {
  243. bool have_owner = false;
  244. for (const auto& user : users.value()) {
  245. if (user.role == "owner") { have_owner = true; break; }
  246. }
  247. if (!have_owner) {
  248. for (const auto& user : users.value()) {
  249. if (user.role != "admin") continue;
  250. auto promoted = storage_->update("users", user.id, {{"role", "owner"}}, 0, true);
  251. if (promoted.ok()) {
  252. LOG_INFO("{} is now the owner of this installation", user.username);
  253. } else {
  254. LOG_WARN("Could not make {} the owner: {}", user.username,
  255. promoted.error().message());
  256. }
  257. break;
  258. }
  259. }
  260. }
  261. }
  262. // Ensure the executions summary view exists before serving requests
  263. ensureExecutionsSummaryView();
  264. // Load active workflows from database and register with scheduler
  265. loadScheduledWorkflows();
  266. // Start NodeSync gRPC server for runners
  267. node_sync_server_->start();
  268. // Start Credential gRPC server for runners
  269. credential_server_->start();
  270. // Start WebSocket server
  271. ws_server_->start();
  272. // Start HTTP server (blocking in background thread)
  273. http_server_->start();
  274. LOG_INFO("WebServer service started on port {}", config_.http_port);
  275. }
  276. void WebServerService::stop() {
  277. LOG_INFO("Stopping WebServer service...");
  278. http_server_->stop();
  279. ws_server_->stop();
  280. credential_server_->stop();
  281. node_sync_server_->stop();
  282. scheduler_->stop();
  283. runner_registry_->stop();
  284. LOG_INFO("WebServer service stopped");
  285. }
  286. void WebServerService::setupRoutes() {
  287. auto& server = http_server_->server();
  288. // Health check
  289. server.Get("/health", [](const httplib::Request& req, httplib::Response& res) {
  290. res.set_content(R"({"status":"ok"})", "application/json");
  291. });
  292. // API version
  293. server.Get("/api/v1", [](const httplib::Request& req, httplib::Response& res) {
  294. res.set_content(R"({"name":"SmartBotic API","version":"1.0.0"})", "application/json");
  295. });
  296. // Register controllers - stored as members to ensure they outlive httplib callbacks
  297. auth_ctrl_ = std::make_unique<api::AuthController>(*auth_store_, *auth_middleware_);
  298. auth_ctrl_->registerRoutes(server);
  299. // Who may do what. Built before any controller that asks it - dereferencing
  300. // this while it was still null bound a reference to nothing, and the first
  301. // request that used it took the server down with a segfault rather than
  302. // failing anywhere near the mistake.
  303. access_ = std::make_unique<auth::AccessControl>(*storage_);
  304. user_ctrl_ = std::make_unique<api::UserController>(
  305. *auth_store_, *auth_middleware_, *access_, *storage_);
  306. user_ctrl_->registerRoutes(server);
  307. // Before every controller that holds a reference to it. Dereferencing an
  308. // empty unique_ptr here does not fail here - it fails later, inside the
  309. // first request that uses it, as a segfault nowhere near the mistake. That
  310. // has happened once on this branch already.
  311. retention_ = std::make_unique<retention::RetentionService>(*storage_);
  312. project_ctrl_ = std::make_unique<api::ProjectController>(
  313. *storage_, *auth_middleware_, *access_, *auth_store_, *retention_);
  314. project_ctrl_->registerRoutes(server);
  315. workflow_ctrl_ = std::make_unique<api::WorkflowController>(
  316. *storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_,
  317. *scheduler_, *db_watcher_, *access_, *node_store_, *retention_);
  318. workflow_ctrl_->registerRoutes(server);
  319. workflow_group_ctrl_ = std::make_unique<api::WorkflowGroupController>(
  320. *storage_, *access_, *auth_middleware_, *ws_server_);
  321. workflow_group_ctrl_->registerRoutes(server);
  322. execution_ctrl_ = std::make_unique<api::ExecutionController>(
  323. *storage_, *auth_middleware_, *access_, *ws_server_, *scheduler_, *load_balancer_,
  324. [this](const std::string& workflow_id, const std::string& execution_id,
  325. bool failed, const std::string& error) {
  326. noteExecutionOutcome(workflow_id, execution_id, failed, error);
  327. });
  328. file_ctrl_ = std::make_unique<api::FileController>(*storage_, *auth_middleware_);
  329. execution_ctrl_->registerRoutes(server);
  330. file_ctrl_->registerRoutes(server);
  331. proxy_ctrl_ = std::make_unique<api::ProxyController>(*auth_middleware_);
  332. proxy_ctrl_->registerRoutes(server);
  333. node_ctrl_ = std::make_unique<api::NodeController>(
  334. *node_store_, *auth_middleware_, *load_balancer_, &node_sync_server_->service());
  335. node_ctrl_->registerRoutes(server);
  336. runner_ctrl_ = std::make_unique<api::RunnerController>(
  337. *runner_registry_, *auth_middleware_,
  338. [this](const std::string& runner_id, const std::string& address) {
  339. reconcileOrphanedExecutions(runner_id, address);
  340. });
  341. runner_ctrl_->registerRoutes(server);
  342. webhook_ctrl_ = std::make_unique<api::WebhookController>(
  343. *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_, *scheduler_);
  344. webhook_ctrl_->registerRoutes(server);
  345. database_ctrl_ = std::make_unique<api::DatabaseController>(*storage_, *auth_middleware_);
  346. database_ctrl_->registerRoutes(server);
  347. credential_ctrl_ = std::make_unique<api::CredentialController>(
  348. *credential_store_, *access_, *storage_, *auth_middleware_);
  349. credential_ctrl_->registerRoutes(server);
  350. LOG_INFO("API routes registered");
  351. }
  352. void WebServerService::ensureExecutionsSummaryView() {
  353. // View names are namespaced per project from smartbotic-database 2.4.2, so a
  354. // view created by this client is also queryable by it. Earlier versions
  355. // registered the name unqualified while queries were project-prefixed, which
  356. // made the view unreachable.
  357. // Recreate rather than reuse. A view is metadata over a collection, so
  358. // rebuilding it costs nothing, and keeping an existing one means the view
  359. // silently outlives the field list below: after a database recovery the old
  360. // view served documents that were no longer in the collection at all.
  361. for (const auto& view : storage_->listViews()) {
  362. if (view.name == kExecutionsSummaryView) {
  363. LOG_INFO("Rebuilding the executions summary view");
  364. storage_->dropView(kExecutionsSummaryView);
  365. break;
  366. }
  367. }
  368. // stopReason belongs here for the same reason error does: it is the message
  369. // explaining how a run ended, and the listing is where someone looks for it.
  370. // A stopped run has status "completed" and an empty error, so without this
  371. // field the list cannot tell a run that finished its work from one that
  372. // deliberately ended early.
  373. auto result = storage_->createView(kExecutionsSummaryView, "executions",
  374. {"workflowId", "workflowName", "status", "triggerType",
  375. "startedAt", "finishedAt", "error", "runnerId",
  376. "stopped", "stopReason", "stoppedNodeId"});
  377. if (result.failed()) {
  378. LOG_ERROR("Could not create the executions summary view: {}. The executions "
  379. "listing will fail until this is resolved.", result.error().message());
  380. return;
  381. }
  382. LOG_INFO("Created server-side view '{}' over executions", kExecutionsSummaryView);
  383. }
  384. void WebServerService::loadScheduledWorkflows() {
  385. LOG_INFO("Loading scheduled workflows from database...");
  386. // Query all active workflows
  387. storage::QueryOptions options;
  388. options.filters.push_back({"active", true});
  389. options.page_size = 1000; // Load up to 1000 workflows
  390. auto result = storage_->query("workflows", options);
  391. if (result.failed()) {
  392. LOG_ERROR("Failed to load workflows: {}", result.error().message());
  393. return;
  394. }
  395. int registered_count = 0;
  396. for (const auto& workflow : result.value().documents) {
  397. std::string workflow_id = workflow.value("_id", "");
  398. std::string workflow_name = workflow.value("name", "");
  399. auto nodes = workflow.value("nodes", nlohmann::json::array());
  400. // Find scheduled trigger nodes
  401. for (const auto& node : nodes) {
  402. std::string node_id = node.value("id", "");
  403. std::string node_type = node.value("type", "");
  404. // Get node definition to check if it's a scheduled trigger
  405. auto node_result = node_store_->get(node_type);
  406. if (node_result.failed()) {
  407. continue;
  408. }
  409. const auto& node_def = node_result.value();
  410. if (!node_def.is_trigger) {
  411. continue;
  412. }
  413. // Told, not asked: the database reports the change, so there is no
  414. // interval. Registered here as well as on activate, or a restart
  415. // would leave every event-driven workflow deaf until somebody
  416. // toggled it.
  417. if (node_type == "database-change") {
  418. auto config = smartbotic::common::applyConfigDefaults(
  419. node.value("config", nlohmann::json::object()), node_def.config_schema);
  420. const std::string collection = config.value("collection", "");
  421. if (collection.empty()) continue;
  422. std::vector<std::string> event_types;
  423. if (config.contains("eventTypes") && config["eventTypes"].is_array()) {
  424. for (const auto& t : config["eventTypes"]) {
  425. if (t.is_string()) event_types.push_back(t.get<std::string>());
  426. }
  427. }
  428. db_watcher_->watch(workflow_id, workflow_name, node_id, collection, event_types);
  429. registered_count++;
  430. continue;
  431. }
  432. if (!node_def.is_scheduled) {
  433. continue;
  434. }
  435. // Get interval from node config, filling in any schema defaults
  436. // the stored config is missing (e.g. an untouched form field).
  437. auto config = smartbotic::common::applyConfigDefaults(
  438. node.value("config", nlohmann::json::object()), node_def.config_schema);
  439. int interval = config.value("pollInterval", 0);
  440. if (interval > 0) {
  441. auto policy = overlapPolicyFromString(
  442. config.value("overlapPolicy", std::string("skip")));
  443. scheduler_->registerWorkflow(
  444. workflow_id,
  445. workflow_name,
  446. node_id,
  447. node_type,
  448. interval,
  449. policy,
  450. config.value("maxConcurrent", 1),
  451. config.value("maxRunMinutes", 0)
  452. );
  453. registered_count++;
  454. LOG_DEBUG("Registered workflow '{}' ({}) for scheduled execution every {} minutes",
  455. workflow_name, workflow_id, interval);
  456. }
  457. }
  458. }
  459. LOG_INFO("Loaded {} scheduled workflows", registered_count);
  460. }
  461. void WebServerService::noteExecutionOutcome(const std::string& workflow_id,
  462. const std::string& execution_id,
  463. bool failed,
  464. const std::string& error) {
  465. if (failed) {
  466. runErrorWorkflow(workflow_id, execution_id, error);
  467. }
  468. auto stored = storage_->get("workflows", workflow_id);
  469. if (stored.failed()) {
  470. return;
  471. }
  472. const auto workflow = stored.value();
  473. const auto settings = workflow.value("settings", nlohmann::json::object());
  474. const int limit = settings.value("deactivateAfterFailures", 0);
  475. const int streak = workflow.value("consecutiveFailures", 0);
  476. if (!failed) {
  477. // Only written when there is something to clear, so an ordinary run does
  478. // not cost a write.
  479. if (streak != 0) {
  480. storage_->update("workflows", workflow_id,
  481. {{"consecutiveFailures", 0}}, 0, true);
  482. }
  483. return;
  484. }
  485. const int next = streak + 1;
  486. nlohmann::json patch = {{"consecutiveFailures", next}};
  487. // Counting is worth doing even when nothing is switched off - it is the
  488. // number someone looks at when asking how long this has been going wrong.
  489. if (limit <= 0 || next < limit || workflow.value("active", false) != true) {
  490. storage_->update("workflows", workflow_id, patch, 0, true);
  491. return;
  492. }
  493. // Switched off, and said out loud. A workflow that simply appeared
  494. // "Inactive" one morning with no reason recorded is worse than one that
  495. // kept failing, because nobody can tell which of the two happened.
  496. patch["active"] = false;
  497. patch["deactivatedReason"] = "Switched off after " + std::to_string(next) +
  498. " failures in a row. The last one: " +
  499. (error.empty() ? "no reason given" : error);
  500. patch["deactivatedAt"] = common::TimeUtils::nowMs();
  501. patch["updatedAt"] = common::TimeUtils::nowMs();
  502. storage_->update("workflows", workflow_id, patch, 0, true);
  503. // Off the schedule too, or it would keep firing while reading as inactive.
  504. scheduler_->unregisterWorkflow(workflow_id);
  505. LOG_WARN("Workflow {} deactivated after {} consecutive failures: {}",
  506. workflow_id, next, error);
  507. ws_server_->broadcast("workflows.deactivated", {
  508. {"id", workflow_id},
  509. {"reason", patch["deactivatedReason"]},
  510. {"consecutiveFailures", next}
  511. });
  512. }
  513. void WebServerService::reconcileOrphanedExecutions(const std::string& runner_id,
  514. const std::string& address) {
  515. // Ask the runner what it is actually running rather than assuming. A runner
  516. // that re-registers while working - a duplicate call, a flapping network -
  517. // must not have its live executions closed underneath it, and only the
  518. // runner knows which those are.
  519. std::unordered_set<std::string> still_running;
  520. {
  521. auto channel = ::grpc::CreateChannel(address, ::grpc::InsecureChannelCredentials());
  522. auto stub = proto::RunnerService::NewStub(channel);
  523. proto::ListActiveExecutionsRequest request;
  524. proto::ListActiveExecutionsResponse response;
  525. ::grpc::ClientContext context;
  526. context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(10));
  527. auto status = stub->ListActiveExecutions(&context, request, &response);
  528. if (!status.ok()) {
  529. // Without an answer there is no way to tell an orphan from a live
  530. // run, and closing a live one is far worse than leaving a stale
  531. // record for someone to press Stop on.
  532. LOG_WARN("Runner {} could not say what it is running ({}), so nothing was reconciled",
  533. runner_id, status.error_message());
  534. return;
  535. }
  536. for (const auto& id : response.execution_ids()) {
  537. still_running.insert(id);
  538. }
  539. }
  540. storage::QueryOptions options;
  541. options.filters.push_back({"runnerId", runner_id});
  542. options.filters.push_back({"status", "running"});
  543. options.page = 1;
  544. options.page_size = 500;
  545. auto found = storage_->query("executions", options);
  546. if (found.failed()) {
  547. LOG_WARN("Could not look for orphaned executions on runner {}: {}",
  548. runner_id, found.error().message());
  549. return;
  550. }
  551. int closed = 0;
  552. for (const auto& record : found.value().documents) {
  553. const std::string id = record.value("_id", "");
  554. if (id.empty() || still_running.contains(id)) {
  555. continue;
  556. }
  557. // Deliberately not touching Waiting. That is a run parked on a person,
  558. // not on a runner, and it is meant to outlive one.
  559. const nlohmann::json patch = {
  560. {"status", "cancelled"},
  561. {"error", "The runner restarted while this was running, so nothing was left to finish it"},
  562. {"finishedAt", common::TimeUtils::nowMs()}
  563. };
  564. if (storage_->update("executions", id, patch, 0, true).ok()) {
  565. ++closed;
  566. ws_server_->broadcast("executions." + id + ".cancelled", {{"executionId", id}});
  567. }
  568. }
  569. if (closed > 0) {
  570. LOG_WARN("Closed {} execution(s) that runner {} was recorded as running but is not",
  571. closed, runner_id);
  572. }
  573. }
  574. void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
  575. const std::string& failed_execution_id,
  576. const std::string& error_message) {
  577. {
  578. std::lock_guard<std::mutex> lock(handled_failures_mutex_);
  579. if (!handled_failures_.insert(failed_execution_id).second) {
  580. return; // already handled this failure
  581. }
  582. if (handled_failures_.size() > 512) {
  583. handled_failures_.erase(handled_failures_.begin());
  584. }
  585. }
  586. auto failed = storage_->get("workflows", failed_workflow_id);
  587. if (failed.failed()) {
  588. return;
  589. }
  590. const auto settings = failed.value().value("settings", nlohmann::json::object());
  591. const std::string handler_id = settings.value("errorWorkflowId", std::string());
  592. if (handler_id.empty()) {
  593. return;
  594. }
  595. // A handler that fails must not summon itself, which would run forever.
  596. if (handler_id == failed_workflow_id) {
  597. LOG_WARN("Workflow {} names itself as its error workflow; not running it",
  598. failed_workflow_id);
  599. return;
  600. }
  601. auto handler = storage_->get("workflows", handler_id);
  602. if (handler.failed()) {
  603. LOG_WARN("Workflow {} names error workflow {}, which no longer exists",
  604. failed_workflow_id, handler_id);
  605. return;
  606. }
  607. auto runner = load_balancer_->selectRunner();
  608. if (!runner) {
  609. LOG_ERROR("No runners available to run error workflow {}", handler_id);
  610. return;
  611. }
  612. // The handler is told what failed rather than having to look it up, so it can
  613. // notify or record without needing read access to the executions collection.
  614. nlohmann::json trigger_data;
  615. trigger_data["errorWorkflow"] = true;
  616. trigger_data["failedWorkflowId"] = failed_workflow_id;
  617. trigger_data["failedWorkflowName"] = failed.value().value("name", std::string());
  618. trigger_data["failedExecutionId"] = failed_execution_id;
  619. trigger_data["error"] = error_message;
  620. trigger_data["failedAt"] = common::TimeUtils::nowMs();
  621. auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
  622. auto stub = proto::RunnerService::NewStub(channel);
  623. proto::ExecuteWorkflowRequest request;
  624. request.set_workflow_id(handler_id);
  625. request.set_trigger_type("error-workflow");
  626. // Started by the system, so it runs what was published.
  627. request.set_use_published(true);
  628. request.set_trigger_data(trigger_data.dump());
  629. request.set_wait_for_completion(false);
  630. proto::ExecuteWorkflowResponse response;
  631. ::grpc::ClientContext context;
  632. context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
  633. auto status = stub->ExecuteWorkflow(&context, request, &response);
  634. if (!status.ok()) {
  635. LOG_ERROR("Error workflow {} could not be started: {}", handler_id, status.error_message());
  636. return;
  637. }
  638. LOG_INFO("Error workflow {} started as {} after {} failed",
  639. handler_id, response.execution_id(), failed_workflow_id);
  640. scheduler_->notifyExecutionStarted(handler_id, response.execution_id());
  641. }
  642. void WebServerService::executeScheduledWorkflow(const std::string& workflow_id,
  643. const std::string& trigger_node_id,
  644. const std::string& trigger_type,
  645. const nlohmann::json& extra_trigger_data) {
  646. // Select a runner
  647. auto runner = load_balancer_->selectRunner();
  648. if (!runner) {
  649. LOG_ERROR("No runners available for scheduled workflow {}", workflow_id);
  650. return;
  651. }
  652. // Create gRPC channel and stub
  653. auto channel = ::grpc::CreateChannel(runner->address, ::grpc::InsecureChannelCredentials());
  654. auto stub = proto::RunnerService::NewStub(channel);
  655. // Prepare request
  656. proto::ExecuteWorkflowRequest request;
  657. request.set_workflow_id(workflow_id);
  658. request.set_trigger_type(trigger_type);
  659. // A schedule firing is not somebody testing an edit: it runs the published
  660. // version, so a half-finished change on somebody's canvas never goes live
  661. // just because it was saved.
  662. request.set_use_published(true);
  663. nlohmann::json trigger_data;
  664. trigger_data["triggerNodeId"] = trigger_node_id;
  665. trigger_data["scheduledExecution"] = true;
  666. // What actually happened, for the triggers that are told rather than the
  667. // ones that ask - a database change carries the document with it.
  668. if (extra_trigger_data.is_object()) {
  669. for (const auto& [key, value] : extra_trigger_data.items()) {
  670. trigger_data[key] = value;
  671. }
  672. }
  673. request.set_trigger_data(trigger_data.dump());
  674. request.set_wait_for_completion(false);
  675. // Execute
  676. proto::ExecuteWorkflowResponse response;
  677. ::grpc::ClientContext context;
  678. context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
  679. auto status = stub->ExecuteWorkflow(&context, request, &response);
  680. if (!status.ok()) {
  681. LOG_ERROR("Failed to execute scheduled workflow {}: {}", workflow_id, status.error_message());
  682. return;
  683. }
  684. LOG_INFO("Scheduled workflow {} execution started: {} on runner {}",
  685. workflow_id, response.execution_id(), runner->id);
  686. // Dispatch is fire-and-forget, so the scheduler only learns about the run here.
  687. scheduler_->notifyExecutionStarted(workflow_id, response.execution_id());
  688. // Broadcast execution started
  689. ws_server_->broadcast("executions." + response.execution_id() + ".started", {
  690. {"executionId", response.execution_id()},
  691. {"workflowId", workflow_id},
  692. {"runnerId", runner->id},
  693. {"triggeredBy", trigger_type},
  694. {"scheduled", true}
  695. });
  696. }
  697. } // namespace smartbotic::webserver