| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908 |
- #include "webserver_service.hpp"
- #include <algorithm>
- #include "api/auth_controller.hpp"
- #include "api/user_controller.hpp"
- #include "api/workflow_controller.hpp"
- #include "api/workflow_group_controller.hpp"
- #include "api/execution_controller.hpp"
- #include "api/node_controller.hpp"
- #include "api/runner_controller.hpp"
- #include "api/webhook_controller.hpp"
- #include "api/database_controller.hpp"
- #include "api/file_controller.hpp"
- #include "api/proxy_controller.hpp"
- #include "api/credential_controller.hpp"
- #include "api/settings_controller.hpp"
- #include "nodes/node_store.hpp"
- #include "grpc/node_sync_service.hpp"
- #include "grpc/credential_service.hpp"
- #include "credentials/credential_store.hpp"
- #include "scheduler/workflow_scheduler.hpp"
- #include "common/time_utils.hpp"
- #include <unordered_set>
- #include "common/config_defaults.hpp"
- #include "logging/logger.hpp"
- #include <grpcpp/grpcpp.h>
- #include "proto/runner.grpc.pb.h"
- namespace smartbotic::webserver {
- namespace {
- // server.form_dispatch_threads / server.form_dispatch_queue_capacity are read
- // as plain int from JSON config and handed to DispatchConfig fields that are
- // std::size_t: 0 threads means every immediate-mode form submission gets a
- // "thanks, your answer was received" that nothing will ever run, and a
- // negative value silently wraps to an enormous unsigned size on the
- // conversion. Clamp to a sane range and say so, rather than let either
- // mistake pass through as if it were what the operator meant.
- std::size_t clampDispatchSetting(const char* name, int configured, int min_value, int max_value) {
- int clamped = std::clamp(configured, min_value, max_value);
- if (clamped != configured) {
- LOG_WARN("Configured {} of {} is out of the allowed range [{}, {}]; using {} instead",
- name, configured, min_value, max_value, clamped);
- }
- return static_cast<std::size_t>(clamped);
- }
- } // namespace
- WebServerService::WebServerService(const WebServerServiceConfig& config)
- : config_(config) {
- // Initialize storage client
- storage::StorageClientConfig storage_config;
- storage_config.address = config_.database_address;
- storage_config.project = config_.database_project;
- storage_ = std::make_unique<storage::StorageClient>(storage_config);
- // Runtime settings: overrides live in the database, read at the point of
- // use. Built right after storage_ so every component below that takes a
- // live setting can be handed a working accessor rather than a null one.
- settings_store_ = std::make_unique<settings::SettingsStore>(*storage_);
- settings_ = std::make_unique<settings::SettingsAccessor>(*settings_store_);
- // Initialize JWT
- jwt_ = std::make_unique<auth::JwtUtils>(config_.jwt_config, settings_.get());
- // Initialize auth store
- auth_store_ = std::make_unique<auth::AuthStore>(*storage_, *jwt_);
- // The indexes this installation's queries actually need. Declared here, and
- // idempotently, so a fresh install or a restored backup gets them without
- // anyone remembering to - they live in the database, not in the config.
- //
- // Only what was measured to help. An index costs write throughput, so a
- // declaration that buys nothing is a real cost paid for ever:
- //
- // sessions.refreshToken every token refresh scanned the whole
- // collection - 24 ms over 9,661 sessions, 1 ms now
- // executions.workflowId a lookup that missed took 1,124 ms and now takes
- // 1. What is left when it hits is the size of the
- // execution documents, about 30 ms each, which no
- // index can help with
- // workflows.projectId small today; the listing filters by it on every
- // page load and workflows are written by hand
- // users.username/email login looks up by both
- //
- // executions.startedAt was declared, dropped, and declared again. On 2.9 it
- // changed nothing: ranges and sorts were not served from an index, and a
- // sorted single row still took 422 ms over 10,000 rows. Re-measured on
- // 2.10, which fixed that, it makes no measurable difference either - but the
- // collection is 1,500 rows now rather than 10,000, having had its orphans
- // removed, so the two measurements are not comparable and neither is a
- // verdict. It is kept because the field is perfectly selective (every value
- // distinct), the planner declines an index it cannot use, and the collection
- // will grow again.
- for (const auto& [collection, field] : std::initializer_list<std::pair<const char*, const char*>>{
- {"executions", "workflowId"},
- {"executions", "startedAt"},
- {"workflows", "projectId"},
- {"users", "username"},
- {"users", "email"},
- {"sessions", "refreshToken"}}) {
- auto created = storage_->createIndex(collection, field);
- if (created.failed()) {
- // Not fatal: an older database has no index support and every query
- // still works, just by scanning.
- LOG_DEBUG("Index {}.{} not declared: {}", collection, field,
- created.error().message());
- } else if (created.value() > 0) {
- LOG_INFO("Index {}.{} declared, {} rows backfilled", collection, field,
- created.value());
- }
- }
- // Sessions issued under a longer lifetime than is configured now would
- // otherwise keep it until they ran out, so shortening the setting would not
- // take effect for as long as the old one lasted.
- auth_store_->enforceSessionLifetime();
- // Email became the login credential and is matched case-insensitively;
- // this brings any row written under the old case-sensitive scheme into
- // line so it is still findable by lookup.
- auth_store_->normalizeUserEmails();
- // Initialize auth middleware
- auth_middleware_ = std::make_unique<auth::AuthMiddleware>(*jwt_, *auth_store_);
- // Initialize runner registry
- runner_registry_ = std::make_unique<runners::RunnerRegistry>(
- *storage_, config_.runner_config, settings_.get());
- // Initialize load balancer
- load_balancer_ = std::make_unique<runners::LoadBalancer>(
- *runner_registry_, config_.load_balancer_config, settings_.get());
- // Initialize node store
- node_store_ = std::make_unique<nodes::NodeStore>(*storage_);
- // Migrate/sync nodes from disk - updates existing nodes if code has changed
- auto migrate_result = node_store_->migrateFromFiles("./nodes");
- if (migrate_result.ok()) {
- LOG_INFO("Node migration: {} nodes imported/updated, {} rejected",
- migrate_result.value().migrated.size(), migrate_result.value().rejected.size());
- for (const auto& rejection : migrate_result.value().rejected) {
- for (const auto& reason : rejection.reasons) {
- LOG_WARN("Node migration rejected {}: {}", rejection.node_id, reason);
- }
- }
- } else {
- LOG_WARN("Node migration failed: {}", migrate_result.error().message());
- }
- // Initialize credential store
- credentials::CredentialStoreConfig cred_config;
- cred_config.master_key = config_.credentials_config.master_key;
- cred_config.pbkdf2_iterations = config_.credentials_config.pbkdf2_iterations;
- credential_store_ = std::make_unique<credentials::CredentialStore>(*storage_, cred_config);
- credential_store_->initialize();
- // Initialize workflow scheduler
- scheduler_ = std::make_unique<WorkflowScheduler>();
- // The scheduler holds a slot per run and frees it when the completion event
- // arrives. Events get lost - a database outage mid-run is enough - so it can
- // also ask whether a run has finished. Answering from the execution record
- // rather than from memory turns a 75-minute stall into one tick.
- scheduler_->setExecutionFinishedCheck([this](const std::string& execution_id) {
- auto record = storage_->get("executions", execution_id);
- if (record.failed()) {
- // Unreadable is not finished. Saying otherwise here would start a
- // second run of a workflow that is still going, which is worse than
- // waiting for the deadline the scheduler already has.
- return false;
- }
- const std::string status = record.value().value("status", "");
- return status == "completed" || status == "failed" || status == "cancelled";
- });
- scheduler_->setExecuteCallback([this](const std::string& workflow_id,
- const std::string& trigger_node_id,
- const std::string& trigger_type) {
- executeScheduledWorkflow(workflow_id, trigger_node_id, trigger_type);
- });
- // Triggers that are told rather than asking. The database says when a
- // collection changed, so these workflows do not poll for it.
- db_watcher_ = std::make_unique<DatabaseWatcher>(*storage_);
- db_watcher_->setExecuteCallback([this](const std::string& workflow_id,
- const std::string& trigger_node_id,
- const std::string& trigger_type,
- const nlohmann::json& event) {
- executeScheduledWorkflow(workflow_id, trigger_node_id, trigger_type, event);
- });
- // Initialize HTTP server
- HttpServerConfig http_config;
- http_config.port = config_.http_port;
- http_config.static_files_path = config_.static_files_path;
- http_config.max_upload_mb = config_.max_upload_mb;
- http_server_ = std::make_unique<HttpServer>(http_config);
- // Initialize WebSocket server on separate port
- WebSocketServerConfig ws_config;
- ws_config.port = config_.http_port + 1; // WebSocket on next port (8091)
- ws_server_ = std::make_unique<WebSocketServer>(ws_config, *jwt_);
- // Initialize NodeSync gRPC server for runners
- node_sync_server_ = std::make_unique<grpc::NodeSyncServer>(*node_store_, config_.node_sync_port);
- // Initialize Credential gRPC server for runners
- credential_server_ = std::make_unique<grpc::CredentialServer>(*credential_store_, config_.credential_service_port);
- }
- WebServerService::~WebServerService() {
- stop();
- }
- WebServerServiceConfig WebServerService::loadConfig(const std::filesystem::path& path) {
- WebServerServiceConfig config;
- auto result = config::Config::fromFile(path);
- if (result.ok()) {
- auto& cfg = result.value();
- config.http_port = cfg.getOr<int>("http_port", 8080);
- config.node_sync_port = cfg.getOr<int>("node_sync_port", 9002);
- config.credential_service_port = cfg.getOr<int>("credential_service_port", 9003);
- config.static_files_path = cfg.getOr<std::string>("static_files_path", "./webui/dist");
- config.database_address = cfg.getOr<std::string>("database_address", "localhost:9004");
- config.database_project = cfg.getOr<std::string>("database_project", "smartbotic-automation");
- // HttpServer already enforced a body cap before this - a hardcoded
- // 16 MB literal at the point it constructs httplib::Server - so this
- // is not closing an open door, it is making that cap configurable
- // and raising it. The default goes from 16 MB to 32 MB because an
- // upload node needs to carry up to a 16 MB image field, and that
- // field cannot fit inside a 16 MB total body once multipart framing
- // is added on top.
- config.max_upload_mb = cfg.getOr<int>("server.max_upload_mb", 32);
- config.form_dispatch_threads = cfg.getOr<int>("server.form_dispatch_threads", 4);
- config.form_dispatch_queue_capacity =
- cfg.getOr<int>("server.form_dispatch_queue_capacity", 32);
- config.workflow_dispatch_threads =
- cfg.getOr<int>("server.workflow_dispatch_threads", 4);
- config.workflow_dispatch_queue_capacity =
- cfg.getOr<int>("server.workflow_dispatch_queue_capacity", 32);
- // Runner config
- config.runner_config.heartbeat_timeout_sec =
- cfg.getOr<int>("runners.heartbeat_timeout_sec", 30);
- config.runner_config.offline_removal_sec =
- cfg.getOr<int>("runners.offline_removal_sec", 60);
- // Load balancer config
- auto strategy_str = cfg.getOr<std::string>("runners.load_balancing", "least-connections");
- config.load_balancer_config.strategy =
- runners::loadBalancingStrategyFromString(strategy_str);
- // Every webserver -> runner gRPC channel is sized off this, not a
- // separate knob - see the comment on LoadBalancerConfig for why
- // server.max_upload_mb (already loaded above) is the right source of
- // truth rather than a second config key that could drift from it.
- config.load_balancer_config.max_message_size_mb = config.max_upload_mb;
- // JWT config
- config.jwt_config.secret = cfg.getOr<std::string>("auth.jwt_secret",
- "dev-secret-change-in-production");
- config.jwt_config.access_token_lifetime_sec =
- cfg.getOr<int>("auth.access_token_lifetime_sec", 900);
- // Was declared in JwtUtils::Config and present in webserver.json, but
- // never actually read here - the file value was silently ignored and
- // the struct default (86400s) always won. Fixed as part of wiring
- // this key into the live settings accessor.
- config.jwt_config.refresh_token_lifetime_sec =
- cfg.getOr<int64_t>("auth.refresh_token_lifetime_sec", 86400);
- // Credentials config
- config.credentials_config.master_key = cfg.getOr<std::string>("credentials.master_key",
- "dev-key-change-in-production");
- config.credentials_config.pbkdf2_iterations =
- cfg.getOr<int>("credentials.pbkdf2_iterations", 100000);
- }
- return config;
- }
- void WebServerService::start() {
- LOG_INFO("Starting WebServer service...");
- // Ensure admin user exists
- auth_store_->ensureAdminUser();
- // Setup API routes
- setupRoutes();
- // Start runner registry cleanup
- runner_registry_->start();
- // Start workflow scheduler
- scheduler_->start();
- // Keep a history of workflow documents. Without it the version number
- // still climbs on every write but nothing older is retained, so there is no
- // earlier version to publish, compare against or go back to.
- //
- // Turned on here rather than at creation because the collection already
- // exists on every install that predates this, and createCollection refuses
- // a collection that is already there.
- {
- storage::CollectionConfig cfg;
- cfg.versioning_enabled = true;
- auto configured = storage_->configureCollection("workflows", cfg);
- if (configured.failed()) {
- LOG_WARN("Could not enable version history on workflows: {}",
- configured.error().message());
- } else {
- LOG_INFO("Version history enabled on workflows");
- }
- }
- // Everybody gets a personal project, and anything stored before projects
- // existed is filed into one - a record with no project is one only an
- // instance admin can see.
- if (project_ctrl_) {
- project_ctrl_->migrateExistingRecords();
- }
- // Somebody has to be the owner - the account that cannot be demoted or
- // deleted, so an installation can never end up with nobody able to
- // administer it. If nobody is, the first admin becomes it.
- {
- auto users = auth_store_->listUsers(1, 1000);
- if (users.ok()) {
- bool have_owner = false;
- for (const auto& user : users.value()) {
- if (user.role == "owner") { have_owner = true; break; }
- }
- if (!have_owner) {
- for (const auto& user : users.value()) {
- if (user.role != "admin") continue;
- auto promoted = storage_->update("users", user.id, {{"role", "owner"}}, 0, true);
- if (promoted.ok()) {
- LOG_INFO("{} is now the owner of this installation", user.username);
- } else {
- LOG_WARN("Could not make {} the owner: {}", user.username,
- promoted.error().message());
- }
- break;
- }
- }
- }
- }
- // Ensure the executions summary view exists before serving requests
- ensureExecutionsSummaryView();
- // Load active workflows from database and register with scheduler
- loadScheduledWorkflows();
- // Start NodeSync gRPC server for runners
- node_sync_server_->start();
- // Start Credential gRPC server for runners
- credential_server_->start();
- // Start WebSocket server
- ws_server_->start();
- // Start HTTP server (blocking in background thread)
- http_server_->start();
- LOG_INFO("WebServer service started on port {}", config_.http_port);
- }
- void WebServerService::stop() {
- LOG_INFO("Stopping WebServer service...");
- http_server_->stop();
- // No new requests can arrive once the HTTP server is stopped, but each
- // controller's own dispatch pool workers may still be mid-run - stop()
- // joins them, so nothing is left running against a controller about to
- // be destroyed.
- if (webhook_ctrl_) {
- webhook_ctrl_->stop();
- }
- if (workflow_ctrl_) {
- workflow_ctrl_->stop();
- }
- ws_server_->stop();
- credential_server_->stop();
- node_sync_server_->stop();
- scheduler_->stop();
- runner_registry_->stop();
- LOG_INFO("WebServer service stopped");
- }
- void WebServerService::setupRoutes() {
- auto& server = http_server_->server();
- // Health check
- server.Get("/health", [](const httplib::Request& req, httplib::Response& res) {
- res.set_content(R"({"status":"ok"})", "application/json");
- });
- // API version
- server.Get("/api/v1", [](const httplib::Request& req, httplib::Response& res) {
- res.set_content(R"({"name":"SmartBotic API","version":"1.0.0"})", "application/json");
- });
- // Register controllers - stored as members to ensure they outlive httplib callbacks
- auth_ctrl_ = std::make_unique<api::AuthController>(*auth_store_, *auth_middleware_);
- auth_ctrl_->registerRoutes(server);
- // Who may do what. Built before any controller that asks it - dereferencing
- // this while it was still null bound a reference to nothing, and the first
- // request that used it took the server down with a segfault rather than
- // failing anywhere near the mistake.
- access_ = std::make_unique<auth::AccessControl>(*storage_);
- user_ctrl_ = std::make_unique<api::UserController>(
- *auth_store_, *auth_middleware_, *access_, *storage_);
- user_ctrl_->registerRoutes(server);
- // Before every controller that holds a reference to it. Dereferencing an
- // empty unique_ptr here does not fail here - it fails later, inside the
- // first request that uses it, as a segfault nowhere near the mistake. That
- // has happened once on this branch already.
- retention_ = std::make_unique<retention::RetentionService>(*storage_);
- project_ctrl_ = std::make_unique<api::ProjectController>(
- *storage_, *auth_middleware_, *access_, *auth_store_, *retention_);
- project_ctrl_->registerRoutes(server);
- api::WorkflowController::DispatchConfig workflow_dispatch_config;
- workflow_dispatch_config.threads = clampDispatchSetting(
- "server.workflow_dispatch_threads", config_.workflow_dispatch_threads, 1, 64);
- workflow_dispatch_config.queue_capacity = clampDispatchSetting(
- "server.workflow_dispatch_queue_capacity", config_.workflow_dispatch_queue_capacity, 1, 10000);
- workflow_dispatch_config.label_kind = "editor execute dispatch";
- workflow_ctrl_ = std::make_unique<api::WorkflowController>(
- *storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_,
- *scheduler_, *db_watcher_, *access_, *node_store_, *retention_, workflow_dispatch_config);
- workflow_ctrl_->registerRoutes(server);
- workflow_group_ctrl_ = std::make_unique<api::WorkflowGroupController>(
- *storage_, *access_, *auth_middleware_, *ws_server_);
- workflow_group_ctrl_->registerRoutes(server);
- execution_ctrl_ = std::make_unique<api::ExecutionController>(
- *storage_, *auth_middleware_, *access_, *ws_server_, *scheduler_, *load_balancer_,
- [this](const std::string& workflow_id, const std::string& execution_id,
- bool failed, const std::string& error) {
- noteExecutionOutcome(workflow_id, execution_id, failed, error);
- });
- file_ctrl_ = std::make_unique<api::FileController>(*storage_, *auth_middleware_);
- execution_ctrl_->registerRoutes(server);
- file_ctrl_->registerRoutes(server);
- proxy_ctrl_ = std::make_unique<api::ProxyController>(*auth_middleware_);
- proxy_ctrl_->registerRoutes(server);
- node_ctrl_ = std::make_unique<api::NodeController>(
- *node_store_, *auth_middleware_, *load_balancer_, &node_sync_server_->service());
- node_ctrl_->registerRoutes(server);
- runner_ctrl_ = std::make_unique<api::RunnerController>(
- *runner_registry_, *auth_middleware_,
- [this](const std::string& runner_id, const std::string& address) {
- reconcileOrphanedExecutions(runner_id, address);
- });
- runner_ctrl_->registerRoutes(server);
- api::WebhookController::DispatchConfig form_dispatch_config;
- form_dispatch_config.threads =
- clampDispatchSetting("server.form_dispatch_threads", config_.form_dispatch_threads, 1, 64);
- form_dispatch_config.queue_capacity = clampDispatchSetting(
- "server.form_dispatch_queue_capacity", config_.form_dispatch_queue_capacity, 1, 10000);
- webhook_ctrl_ = std::make_unique<api::WebhookController>(
- *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_, *scheduler_,
- *jwt_, form_dispatch_config);
- webhook_ctrl_->registerRoutes(server);
- database_ctrl_ = std::make_unique<api::DatabaseController>(*storage_, *auth_middleware_);
- database_ctrl_->registerRoutes(server);
- credential_ctrl_ = std::make_unique<api::CredentialController>(
- *credential_store_, *access_, *storage_, *auth_middleware_);
- credential_ctrl_->registerRoutes(server);
- settings_ctrl_ = std::make_unique<api::SettingsController>(
- *settings_store_, *auth_middleware_, config_);
- settings_ctrl_->registerRoutes(server);
- LOG_INFO("API routes registered");
- }
- void WebServerService::ensureExecutionsSummaryView() {
- // View names are namespaced per project from smartbotic-database 2.4.2, so a
- // view created by this client is also queryable by it. Earlier versions
- // registered the name unqualified while queries were project-prefixed, which
- // made the view unreachable.
- // Recreate rather than reuse. A view is metadata over a collection, so
- // rebuilding it costs nothing, and keeping an existing one means the view
- // silently outlives the field list below: after a database recovery the old
- // view served documents that were no longer in the collection at all.
- for (const auto& view : storage_->listViews()) {
- if (view.name == kExecutionsSummaryView) {
- LOG_INFO("Rebuilding the executions summary view");
- storage_->dropView(kExecutionsSummaryView);
- break;
- }
- }
- // stopReason belongs here for the same reason error does: it is the message
- // explaining how a run ended, and the listing is where someone looks for it.
- // A stopped run has status "completed" and an empty error, so without this
- // field the list cannot tell a run that finished its work from one that
- // deliberately ended early.
- //
- // toleratedErrorCount is the same problem one level down: a node can
- // swallow its own failure (e.g. skipOnError) so the run keeps going, and
- // the whole execution still ends up status "completed" with an empty
- // top-level error. Without this field a run that quietly tolerated a
- // failed feed is indistinguishable, in the list, from a run where nothing
- // went wrong.
- auto result = storage_->createView(kExecutionsSummaryView, "executions",
- {"workflowId", "workflowName", "status", "triggerType",
- "startedAt", "finishedAt", "error", "runnerId",
- "stopped", "stopReason", "stoppedNodeId",
- "toleratedErrorCount"});
- if (result.failed()) {
- LOG_ERROR("Could not create the executions summary view: {}. The executions "
- "listing will fail until this is resolved.", result.error().message());
- return;
- }
- LOG_INFO("Created server-side view '{}' over executions", kExecutionsSummaryView);
- }
- void WebServerService::loadScheduledWorkflows() {
- LOG_INFO("Loading scheduled workflows from database...");
- // Query all active workflows
- storage::QueryOptions options;
- options.filters.push_back({"active", true});
- options.page_size = 1000; // Load up to 1000 workflows
- auto result = storage_->query("workflows", options);
- if (result.failed()) {
- LOG_ERROR("Failed to load workflows: {}", result.error().message());
- return;
- }
- int registered_count = 0;
- for (const auto& workflow : result.value().documents) {
- std::string workflow_id = workflow.value("_id", "");
- std::string workflow_name = workflow.value("name", "");
- auto nodes = workflow.value("nodes", nlohmann::json::array());
- // Find scheduled trigger nodes
- for (const auto& node : nodes) {
- std::string node_id = node.value("id", "");
- std::string node_type = node.value("type", "");
- // Get node definition to check if it's a scheduled trigger
- auto node_result = node_store_->get(node_type);
- if (node_result.failed()) {
- continue;
- }
- const auto& node_def = node_result.value();
- if (!node_def.is_trigger) {
- continue;
- }
- // Told, not asked: the database reports the change, so there is no
- // interval. Registered here as well as on activate, or a restart
- // would leave every event-driven workflow deaf until somebody
- // toggled it.
- if (node_type == "database-change") {
- auto config = smartbotic::common::applyConfigDefaults(
- node.value("config", nlohmann::json::object()), node_def.config_schema);
- const std::string collection = config.value("collection", "");
- if (collection.empty()) continue;
- std::vector<std::string> event_types;
- if (config.contains("eventTypes") && config["eventTypes"].is_array()) {
- for (const auto& t : config["eventTypes"]) {
- if (t.is_string()) event_types.push_back(t.get<std::string>());
- }
- }
- db_watcher_->watch(workflow_id, workflow_name, node_id, collection, event_types);
- registered_count++;
- continue;
- }
- if (!node_def.is_scheduled) {
- continue;
- }
- // Get interval from node config, filling in any schema defaults
- // the stored config is missing (e.g. an untouched form field).
- auto config = smartbotic::common::applyConfigDefaults(
- node.value("config", nlohmann::json::object()), node_def.config_schema);
- int interval = config.value("pollInterval", 0);
- if (interval > 0) {
- auto policy = overlapPolicyFromString(
- config.value("overlapPolicy", std::string("skip")));
- scheduler_->registerWorkflow(
- workflow_id,
- workflow_name,
- node_id,
- node_type,
- interval,
- policy,
- config.value("maxConcurrent", 1),
- config.value("maxRunMinutes", 0)
- );
- registered_count++;
- LOG_DEBUG("Registered workflow '{}' ({}) for scheduled execution every {} minutes",
- workflow_name, workflow_id, interval);
- }
- }
- }
- LOG_INFO("Loaded {} scheduled workflows", registered_count);
- }
- void WebServerService::noteExecutionOutcome(const std::string& workflow_id,
- const std::string& execution_id,
- bool failed,
- const std::string& error) {
- if (failed) {
- runErrorWorkflow(workflow_id, execution_id, error);
- }
- auto stored = storage_->get("workflows", workflow_id);
- if (stored.failed()) {
- return;
- }
- const auto workflow = stored.value();
- const auto settings = workflow.value("settings", nlohmann::json::object());
- const int limit = settings.value("deactivateAfterFailures", 0);
- const int streak = workflow.value("consecutiveFailures", 0);
- if (!failed) {
- // Only written when there is something to clear, so an ordinary run does
- // not cost a write.
- if (streak != 0) {
- storage_->update("workflows", workflow_id,
- {{"consecutiveFailures", 0}}, 0, true);
- }
- return;
- }
- const int next = streak + 1;
- nlohmann::json patch = {{"consecutiveFailures", next}};
- // Counting is worth doing even when nothing is switched off - it is the
- // number someone looks at when asking how long this has been going wrong.
- if (limit <= 0 || next < limit || workflow.value("active", false) != true) {
- storage_->update("workflows", workflow_id, patch, 0, true);
- return;
- }
- // Switched off, and said out loud. A workflow that simply appeared
- // "Inactive" one morning with no reason recorded is worse than one that
- // kept failing, because nobody can tell which of the two happened.
- patch["active"] = false;
- patch["deactivatedReason"] = "Switched off after " + std::to_string(next) +
- " failures in a row. The last one: " +
- (error.empty() ? "no reason given" : error);
- patch["deactivatedAt"] = common::TimeUtils::nowMs();
- patch["updatedAt"] = common::TimeUtils::nowMs();
- storage_->update("workflows", workflow_id, patch, 0, true);
- // Off the schedule too, or it would keep firing while reading as inactive.
- scheduler_->unregisterWorkflow(workflow_id);
- LOG_WARN("Workflow {} deactivated after {} consecutive failures: {}",
- workflow_id, next, error);
- ws_server_->broadcast("workflows.deactivated", {
- {"id", workflow_id},
- {"reason", patch["deactivatedReason"]},
- {"consecutiveFailures", next}
- });
- }
- void WebServerService::reconcileOrphanedExecutions(const std::string& runner_id,
- const std::string& address) {
- // Ask the runner what it is actually running rather than assuming. A runner
- // that re-registers while working - a duplicate call, a flapping network -
- // must not have its live executions closed underneath it, and only the
- // runner knows which those are.
- std::unordered_set<std::string> still_running;
- {
- auto channel = load_balancer_->createChannel(address);
- auto stub = proto::RunnerService::NewStub(channel);
- proto::ListActiveExecutionsRequest request;
- proto::ListActiveExecutionsResponse response;
- ::grpc::ClientContext context;
- context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(10));
- auto status = stub->ListActiveExecutions(&context, request, &response);
- if (!status.ok()) {
- // Without an answer there is no way to tell an orphan from a live
- // run, and closing a live one is far worse than leaving a stale
- // record for someone to press Stop on.
- LOG_WARN("Runner {} could not say what it is running ({}), so nothing was reconciled",
- runner_id, status.error_message());
- return;
- }
- for (const auto& id : response.execution_ids()) {
- still_running.insert(id);
- }
- }
- storage::QueryOptions options;
- options.filters.push_back({"runnerId", runner_id});
- options.filters.push_back({"status", "running"});
- options.page = 1;
- options.page_size = 500;
- auto found = storage_->query("executions", options);
- if (found.failed()) {
- LOG_WARN("Could not look for orphaned executions on runner {}: {}",
- runner_id, found.error().message());
- return;
- }
- int closed = 0;
- for (const auto& record : found.value().documents) {
- const std::string id = record.value("_id", "");
- if (id.empty() || still_running.contains(id)) {
- continue;
- }
- // Deliberately not touching Waiting. That is a run parked on a person,
- // not on a runner, and it is meant to outlive one.
- const nlohmann::json patch = {
- {"status", "cancelled"},
- {"error", "The runner restarted while this was running, so nothing was left to finish it"},
- {"finishedAt", common::TimeUtils::nowMs()}
- };
- if (storage_->update("executions", id, patch, 0, true).ok()) {
- ++closed;
- ws_server_->broadcast("executions." + id + ".cancelled", {{"executionId", id}});
- }
- }
- if (closed > 0) {
- LOG_WARN("Closed {} execution(s) that runner {} was recorded as running but is not",
- closed, runner_id);
- }
- }
- void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
- const std::string& failed_execution_id,
- const std::string& error_message) {
- {
- std::lock_guard<std::mutex> lock(handled_failures_mutex_);
- if (!handled_failures_.insert(failed_execution_id).second) {
- return; // already handled this failure
- }
- if (handled_failures_.size() > 512) {
- handled_failures_.erase(handled_failures_.begin());
- }
- }
- auto failed = storage_->get("workflows", failed_workflow_id);
- if (failed.failed()) {
- return;
- }
- const auto settings = failed.value().value("settings", nlohmann::json::object());
- const std::string handler_id = settings.value("errorWorkflowId", std::string());
- if (handler_id.empty()) {
- return;
- }
- // A handler that fails must not summon itself, which would run forever.
- if (handler_id == failed_workflow_id) {
- LOG_WARN("Workflow {} names itself as its error workflow; not running it",
- failed_workflow_id);
- return;
- }
- auto handler = storage_->get("workflows", handler_id);
- if (handler.failed()) {
- LOG_WARN("Workflow {} names error workflow {}, which no longer exists",
- failed_workflow_id, handler_id);
- return;
- }
- auto runner = load_balancer_->selectRunner();
- if (!runner) {
- LOG_ERROR("No runners available to run error workflow {}", handler_id);
- return;
- }
- // The handler is told what failed rather than having to look it up, so it can
- // notify or record without needing read access to the executions collection.
- nlohmann::json trigger_data;
- trigger_data["errorWorkflow"] = true;
- trigger_data["failedWorkflowId"] = failed_workflow_id;
- trigger_data["failedWorkflowName"] = failed.value().value("name", std::string());
- trigger_data["failedExecutionId"] = failed_execution_id;
- trigger_data["error"] = error_message;
- trigger_data["failedAt"] = common::TimeUtils::nowMs();
- auto channel = load_balancer_->createChannel(runner->address);
- auto stub = proto::RunnerService::NewStub(channel);
- proto::ExecuteWorkflowRequest request;
- request.set_workflow_id(handler_id);
- request.set_trigger_type("error-workflow");
- // Started by the system, so it runs what was published.
- request.set_use_published(true);
- request.set_trigger_data(trigger_data.dump());
- request.set_wait_for_completion(false);
- proto::ExecuteWorkflowResponse response;
- ::grpc::ClientContext context;
- context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
- auto status = stub->ExecuteWorkflow(&context, request, &response);
- if (!status.ok()) {
- LOG_ERROR("Error workflow {} could not be started: {}", handler_id, status.error_message());
- return;
- }
- LOG_INFO("Error workflow {} started as {} after {} failed",
- handler_id, response.execution_id(), failed_workflow_id);
- scheduler_->notifyExecutionStarted(handler_id, response.execution_id());
- }
- void WebServerService::executeScheduledWorkflow(const std::string& workflow_id,
- const std::string& trigger_node_id,
- const std::string& trigger_type,
- const nlohmann::json& extra_trigger_data) {
- // Select a runner
- auto runner = load_balancer_->selectRunner();
- if (!runner) {
- LOG_ERROR("No runners available for scheduled workflow {}", workflow_id);
- return;
- }
- // Create gRPC channel and stub
- auto channel = load_balancer_->createChannel(runner->address);
- auto stub = proto::RunnerService::NewStub(channel);
- // Prepare request
- proto::ExecuteWorkflowRequest request;
- request.set_workflow_id(workflow_id);
- request.set_trigger_type(trigger_type);
- // A schedule firing is not somebody testing an edit: it runs the published
- // version, so a half-finished change on somebody's canvas never goes live
- // just because it was saved.
- request.set_use_published(true);
- nlohmann::json trigger_data;
- trigger_data["triggerNodeId"] = trigger_node_id;
- trigger_data["scheduledExecution"] = true;
- // What actually happened, for the triggers that are told rather than the
- // ones that ask - a database change carries the document with it.
- if (extra_trigger_data.is_object()) {
- for (const auto& [key, value] : extra_trigger_data.items()) {
- trigger_data[key] = value;
- }
- }
- request.set_trigger_data(trigger_data.dump());
- request.set_wait_for_completion(false);
- // Execute
- proto::ExecuteWorkflowResponse response;
- ::grpc::ClientContext context;
- context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(30));
- auto status = stub->ExecuteWorkflow(&context, request, &response);
- if (!status.ok()) {
- LOG_ERROR("Failed to execute scheduled workflow {}: {}", workflow_id, status.error_message());
- return;
- }
- LOG_INFO("Scheduled workflow {} execution started: {} on runner {}",
- workflow_id, response.execution_id(), runner->id);
- // Dispatch is fire-and-forget, so the scheduler only learns about the run here.
- scheduler_->notifyExecutionStarted(workflow_id, response.execution_id());
- // Broadcast execution started
- ws_server_->broadcast("executions." + response.execution_id() + ".started", {
- {"executionId", response.execution_id()},
- {"workflowId", workflow_id},
- {"runnerId", runner->id},
- {"triggeredBy", trigger_type},
- {"scheduled", true}
- });
- }
- } // namespace smartbotic::webserver
|