webserver_service.cpp 39 KB

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