| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036 |
- #include "runner_service.hpp"
- #include "common/uuid.hpp"
- #include "common/time_utils.hpp"
- #include "logging/logger.hpp"
- #include "proto/runner.grpc.pb.h"
- #include <grpcpp/health_check_service_interface.h>
- #include <sys/resource.h>
- #include <fstream>
- #include <curl/curl.h>
- #ifdef MYSQL_SUPPORT
- #include "runner/mysql/mysql_client.hpp"
- #endif
- #ifdef POSTGRESQL_SUPPORT
- #include "runner/postgresql/postgresql_client.hpp"
- #endif
- namespace smartbotic::runner {
- using namespace common;
- // Maps the engine's ExecutionStatus onto the wire enum explicitly. The two
- // enums are not numerically aligned: proto::ExecutionStatus reserves 0 for
- // EXECUTION_STATUS_UNSPECIFIED, while ExecutionStatus::Pending is 0, so a
- // bare static_cast silently shifts every status by one. Deliberately no
- // default label, so an unhandled case is a compiler warning rather than a
- // silent mismatch.
- static proto::ExecutionStatus toProtoStatus(ExecutionStatus status) {
- switch (status) {
- case ExecutionStatus::Pending: return proto::EXECUTION_STATUS_PENDING;
- case ExecutionStatus::Running: return proto::EXECUTION_STATUS_RUNNING;
- case ExecutionStatus::Completed: return proto::EXECUTION_STATUS_COMPLETED;
- case ExecutionStatus::Failed: return proto::EXECUTION_STATUS_FAILED;
- case ExecutionStatus::Cancelled: return proto::EXECUTION_STATUS_CANCELLED;
- case ExecutionStatus::Waiting: return proto::EXECUTION_STATUS_WAITING;
- }
- return proto::EXECUTION_STATUS_UNSPECIFIED;
- }
- // RunnerServiceImpl implementation
- RunnerServiceImpl::RunnerServiceImpl(WorkflowEngine& engine, NodeRegistry& registry,
- storage::StorageClient& storage,
- ExecutionEventCallback event_callback)
- : engine_(engine), registry_(registry), storage_(storage), event_callback_(event_callback) {}
- grpc::Status RunnerServiceImpl::ExecuteWorkflow(grpc::ServerContext* context,
- const proto::ExecuteWorkflowRequest* request,
- proto::ExecuteWorkflowResponse* response) {
- LOG_INFO("ExecuteWorkflow called for workflow: {}", request->workflow_id());
- // An inline workflow runs as given, without being stored anywhere. This is
- // how the editor asks a node what it can offer - which models a server has,
- // say - while its config is still being edited.
- if (!request->inline_workflow().empty()) {
- nlohmann::json inline_doc;
- try {
- inline_doc = nlohmann::json::parse(request->inline_workflow());
- } catch (const std::exception& e) {
- response->set_status(proto::EXECUTION_STATUS_FAILED);
- auto* error = response->mutable_error();
- error->set_message(std::string("inline_workflow is not valid JSON: ") + e.what());
- return grpc::Status::OK;
- }
- // The id decides which credentials this run may read, so it comes from
- // the request rather than from the caller-supplied document.
- inline_doc["_id"] = request->workflow_id();
- auto inline_workflow = Workflow::fromJson(inline_doc);
- nlohmann::json inline_trigger;
- if (!request->trigger_data().empty()) {
- try {
- inline_trigger = nlohmann::json::parse(request->trigger_data());
- } catch (...) {}
- }
- auto inline_outcome = engine_.execute(inline_workflow, "manual", inline_trigger, nullptr);
- if (inline_outcome.failed()) {
- response->set_status(proto::EXECUTION_STATUS_FAILED);
- auto* error = response->mutable_error();
- error->set_message(inline_outcome.error().message());
- return grpc::Status::OK;
- }
- const auto& inline_result = inline_outcome.value();
- response->set_execution_id(inline_result.execution_id);
- response->set_status(toProtoStatus(inline_result.status));
- response->set_result(inline_result.toJson().dump());
- return grpc::Status::OK;
- }
- // Get workflow from database
- auto workflow_result = storage_.get("workflows", request->workflow_id());
- // A trigger runs what was published. The record itself is the draft.
- if (request->use_published() && workflow_result.ok()) {
- const int64_t published = workflow_result.value().value("publishedVersion", int64_t{0});
- const int64_t current = workflow_result.value().value("_version", int64_t{0});
- if (published > 0 && published != current) {
- auto pinned = storage_.getVersion("workflows", request->workflow_id(), published);
- if (pinned.ok()) {
- LOG_INFO("Workflow {} running published version {} (draft is {})",
- request->workflow_id(), published, current);
- // A stored version is the document as it was, and the database
- // keeps _id outside the document - so a version comes back
- // without one. Everything downstream identifies the run by it,
- // and without it the execution recorded an empty workflowId:
- // the run happened, but the workflow's own execution list never
- // showed it, which reads as "it did not run".
- pinned.value()["_id"] = request->workflow_id();
- workflow_result = pinned;
- } else {
- // Refusing would stop a live workflow because its history has
- // aged out, which is worse than running the draft - but it must
- // be said out loud, because the two can differ.
- LOG_WARN("Workflow {}: published version {} is no longer stored, "
- "running the current draft instead", request->workflow_id(), published);
- }
- }
- }
- if (workflow_result.failed()) {
- LOG_ERROR("Failed to get workflow from database: {}", workflow_result.error().message());
- response->set_status(proto::EXECUTION_STATUS_FAILED);
- auto* error = response->mutable_error();
- error->set_code(static_cast<int32_t>(workflow_result.error().code()));
- error->set_message(workflow_result.error().message());
- return grpc::Status::OK;
- }
- // Parse workflow
- auto workflow = Workflow::fromJson(workflow_result.value());
- // Parse trigger data
- nlohmann::json trigger_data;
- if (!request->trigger_data().empty()) {
- try {
- trigger_data = nlohmann::json::parse(request->trigger_data());
- } catch (...) {
- trigger_data = request->trigger_data();
- }
- }
- // Create callback to forward events
- ExecutionCallback callback;
- if (event_callback_) {
- callback = [this, &workflow](const std::string& event_type, const nlohmann::json& data) {
- nlohmann::json event_data = data;
- event_data["workflowId"] = workflow.id;
- event_callback_(event_type, event_data);
- };
- }
- // Execute workflow with callback
- LOG_INFO("Starting workflow execution...");
- auto result = engine_.execute(workflow, request->trigger_type(), trigger_data, callback);
- if (result.failed()) {
- LOG_ERROR("Workflow execution failed: {}", result.error().message());
- // Emit execution.failed event
- if (event_callback_) {
- event_callback_("execution.failed", {
- {"executionId", ""},
- {"workflowId", workflow.id},
- {"error", result.error().message()}
- });
- }
- response->set_status(proto::EXECUTION_STATUS_FAILED);
- auto* error = response->mutable_error();
- error->set_code(static_cast<int32_t>(result.error().code()));
- error->set_message(result.error().message());
- return grpc::Status::OK;
- }
- LOG_INFO("Workflow execution completed, execution_id: {}, status: {}",
- result.value().execution_id, executionStatusToString(result.value().status));
- // No terminal event is emitted here. The engine already emits one for
- // every finished run, on its own way out, and this was a second copy of
- // the same event: every completed run broadcast execution.completed twice
- // and every failed run broadcast execution.failed twice. Confirmed on the
- // wire before removing - two broadcasts a millisecond apart, from one
- // "Workflow execution completed" in the runner's log.
- //
- // The duplicate was not just noise. The webserver runs an error workflow
- // on execution.failed, so an error handler ran twice per failure, and it
- // releases the scheduler slot on the same event.
- //
- // The engine's is also the better of the two: it covers waiting and
- // cancelled rather than only the two statuses handled here, it carries the
- // status and error fields, it truncates a large final output instead of
- // posting the whole thing over HTTP, and it fires on the resume path too -
- // which is why a resumed run already emitted exactly once and made the
- // asymmetry visible.
- //
- // The failure branch above stays: it covers execute() itself returning an
- // error, where the engine produced no result and emitted nothing.
- response->set_execution_id(result.value().execution_id);
- response->set_status(toProtoStatus(result.value().status));
- if (request->wait_for_completion()) {
- if (!result.value().webhook_response.is_null()) {
- nlohmann::json envelope;
- envelope["_webhookResponse"] = result.value().webhook_response;
- response->set_result(envelope.dump());
- } else {
- response->set_result(result.value().final_output.dump());
- }
- }
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::CancelExecution(grpc::ServerContext* context,
- const proto::CancelExecutionRequest* request,
- proto::CancelExecutionResponse* response) {
- const auto outcome = engine_.cancelExecution(request->execution_id());
- response->set_was_running(outcome.was_running);
- response->set_marked(outcome.marked);
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::ListActiveExecutions(
- grpc::ServerContext* context,
- const proto::ListActiveExecutionsRequest* request,
- proto::ListActiveExecutionsResponse* response) {
- for (const auto& id : engine_.activeExecutionIds()) {
- response->add_execution_ids(id);
- }
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::ResumeExecution(grpc::ServerContext* context,
- const proto::ResumeExecutionRequest* request,
- proto::ExecuteWorkflowResponse* response) {
- nlohmann::json payload = nlohmann::json::object();
- if (!request->payload().empty()) {
- try {
- payload = nlohmann::json::parse(request->payload());
- } catch (const std::exception& e) {
- return grpc::Status(grpc::StatusCode::INVALID_ARGUMENT,
- std::string("payload is not JSON: ") + e.what());
- }
- }
- // The resumed half of the run needs the same event plumbing the initial
- // half gets in ExecuteWorkflow above, or the UI shows nothing for it and
- // execution.failed never reaches the handler that runs error workflows.
- // ExecuteWorkflow's callback stamps workflowId onto every event because
- // most of the engine's per-node events don't carry it themselves; this
- // does the same, reading workflowId from the execution record up front
- // since resume() (unlike execute()) is not handed a parsed Workflow by
- // its caller.
- std::string workflow_id;
- auto stored = storage_.get("executions", request->execution_id());
- if (stored.ok()) {
- workflow_id = stored.value().value("workflowId", "");
- }
- ExecutionCallback callback;
- if (event_callback_) {
- callback = [this, workflow_id](const std::string& event_type, const nlohmann::json& data) {
- nlohmann::json event_data = data;
- event_data["workflowId"] = workflow_id;
- event_callback_(event_type, event_data);
- };
- }
- auto result = engine_.resume(request->execution_id(), request->token(), payload, callback);
- if (result.failed()) {
- return grpc::Status(grpc::StatusCode::FAILED_PRECONDITION, result.error().message());
- }
- response->set_execution_id(result.value().execution_id);
- response->set_status(toProtoStatus(result.value().status));
- response->set_result(result.value().final_output.dump());
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::ListNodes(grpc::ServerContext* context,
- const proto::ListNodesRequest* request,
- proto::ListNodesResponse* response) {
- std::vector<NodeDefinition> nodes;
- if (request->category().empty()) {
- nodes = registry_.getAllNodes();
- } else {
- nodes = registry_.getNodesByCategory(request->category());
- }
- for (const auto& node : nodes) {
- auto* proto_node = response->add_nodes();
- proto_node->set_id(node.id);
- proto_node->set_name(node.name);
- proto_node->set_category(node.category);
- proto_node->set_version(node.version);
- proto_node->set_description(node.description);
- proto_node->set_icon(node.icon);
- proto_node->set_is_trigger(node.is_trigger);
- proto_node->set_config_schema(node.config_schema.dump());
- proto_node->set_input_schema(node.input_schema.dump());
- proto_node->set_output_schema(node.output_schema.dump());
- for (const auto& input : node.inputs) {
- auto* proto_input = proto_node->add_inputs();
- proto_input->set_name(input.name);
- proto_input->set_display_name(input.display_name);
- proto_input->set_type(input.type);
- proto_input->set_required(input.required);
- }
- for (const auto& output : node.outputs) {
- auto* proto_output = proto_node->add_outputs();
- proto_output->set_name(output.name);
- proto_output->set_display_name(output.display_name);
- proto_output->set_type(output.type);
- if (!output.color.empty()) {
- proto_output->set_color(output.color);
- }
- }
- }
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::ReloadNode(grpc::ServerContext* context,
- const proto::ReloadNodeRequest* request,
- proto::ReloadNodeResponse* response) {
- // Nodes are now synced automatically from webserver
- // Manual reload is no longer needed
- auto node = registry_.getNode(request->node_id());
- if (!node) {
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(404);
- error->set_message("Node not found: " + request->node_id());
- return grpc::Status::OK;
- }
- response->set_success(true);
- auto* proto_node = response->mutable_node();
- proto_node->set_id(node->id);
- proto_node->set_name(node->name);
- proto_node->set_version(node->version);
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::ExecuteNode(grpc::ServerContext* context,
- const proto::ExecuteNodeRequest* request,
- proto::ExecuteNodeResponse* response) {
- auto node_def = registry_.getNode(request->node_type());
- if (!node_def) {
- response->set_success(false);
- response->set_error("Node type not found: " + request->node_type());
- return grpc::Status::OK;
- }
- // Parse input and config
- nlohmann::json input, config;
- try {
- if (!request->input().empty()) {
- input = nlohmann::json::parse(request->input());
- }
- if (!request->config().empty()) {
- config = nlohmann::json::parse(request->config());
- }
- } catch (const std::exception& e) {
- response->set_success(false);
- response->set_error("Invalid JSON: " + std::string(e.what()));
- return grpc::Status::OK;
- }
- // Create execution context
- engine::ScriptContext ctx;
- ctx.execution_id = common::UUID::generate();
- ctx.node_id = request->node_type();
- ctx.input = input;
- ctx.config = config;
- // Execute
- engine::ScriptEnginePool pool(1);
- auto* engine = pool.acquire();
- auto result = engine->execute(node_def->code, ctx);
- pool.release(engine);
- response->set_success(result.success);
- response->set_output(result.output.dump());
- response->set_error(result.error);
- response->set_execution_time_ms(result.execution_time_ms);
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::GetNodeCode(grpc::ServerContext* context,
- const proto::GetNodeCodeRequest* request,
- proto::GetNodeCodeResponse* response) {
- auto result = registry_.getNodeCode(request->node_id());
- if (result.failed()) {
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(static_cast<int32_t>(result.error().code()));
- error->set_message(result.error().message());
- return grpc::Status::OK;
- }
- response->set_success(true);
- response->set_code(result.value());
- // file_path no longer applicable - nodes stored in database
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::SaveNodeCode(grpc::ServerContext* context,
- const proto::SaveNodeCodeRequest* request,
- proto::SaveNodeCodeResponse* response) {
- // Node modifications are now handled centrally by the webserver
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(501);
- error->set_message("Node modifications should be done through the webserver API");
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::CreateNode(grpc::ServerContext* context,
- const proto::CreateNodeRequest* request,
- proto::CreateNodeResponse* response) {
- // Node creation is now handled centrally by the webserver
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(501);
- error->set_message("Node creation should be done through the webserver API");
- return grpc::Status::OK;
- }
- grpc::Status RunnerServiceImpl::DeleteNode(grpc::ServerContext* context,
- const proto::DeleteNodeRequest* request,
- proto::DeleteNodeResponse* response) {
- // Node deletion is now handled centrally by the webserver
- response->set_success(false);
- auto* error = response->mutable_error();
- error->set_code(501);
- error->set_message("Node deletion should be done through the webserver API");
- return grpc::Status::OK;
- }
- // RunnerService implementation
- RunnerService::RunnerService(const RunnerServiceConfig& config)
- : config_(config) {
- // Initialize storage client
- storage::StorageClientConfig storage_config;
- storage_config.address = config_.database_address;
- storage_config.project = config_.database_project;
- storage_config.max_message_size_mb = config_.max_message_size_mb;
- storage_ = std::make_unique<storage::StorageClient>(storage_config);
- // Initialize credential client
- credentials::CredentialClientConfig cred_config;
- cred_config.address = config_.credential_service_address;
- credential_client_ = std::make_unique<credentials::CredentialClient>(cred_config);
- // Initialize node registry - now loads from webserver
- config_.node_registry_config.webserver_address = config_.node_sync_address;
- registry_ = std::make_unique<NodeRegistry>(config_.node_registry_config);
- // Initialize workflow engine
- config_.workflow_engine_config.max_concurrent_executions = config_.max_concurrent_executions;
- config_.workflow_engine_config.runner_id = config_.runner_id;
- engine_ = std::make_unique<WorkflowEngine>(*registry_, *storage_, config_.workflow_engine_config);
- // Set up credential auth callback for workflow engine
- engine_->setCredentialAuthCallback(
- [this](const std::string& credential_id, const std::string& workflow_id)
- -> common::Result<engine::CredentialAuth> {
- auto result = credential_client_->getHttpAuth(credential_id, workflow_id);
- if (result.failed()) {
- return result.error();
- }
- engine::CredentialAuth auth;
- auth.header_name = result.value().header_name;
- auth.header_value = result.value().header_value;
- return auth;
- });
- // The mTLS identity a workflow presents, fetched over the same mediated
- // service the other secrets use.
- engine_->setClientCertificateCallback(
- [this](const std::string& credential_id, const std::string& workflow_id)
- -> common::Result<engine::TlsClientIdentity> {
- auto result = credential_client_->getClientCertificate(credential_id, workflow_id);
- if (result.failed()) {
- return result.error();
- }
- engine::TlsClientIdentity identity;
- identity.certificate_pem = result.value().certificate_pem;
- identity.private_key_pem = result.value().private_key_pem;
- identity.passphrase = result.value().passphrase;
- return identity;
- });
- // Set up IMAP credential callback for workflow engine
- engine_->setImapCredentialCallback(
- [this](const std::string& credential_id, const std::string& workflow_id)
- -> common::Result<engine::ImapCredential> {
- auto result = credential_client_->getImapCredentials(credential_id, workflow_id);
- if (result.failed()) {
- return result.error();
- }
- engine::ImapCredential cred;
- cred.host = result.value().host;
- cred.port = result.value().port;
- cred.username = result.value().username;
- cred.password = result.value().password;
- cred.use_ssl = result.value().use_ssl;
- return cred;
- });
- // Set up SMTP credential callback for workflow engine
- engine_->setSmtpCredentialCallback(
- [this](const std::string& credential_id, const std::string& workflow_id)
- -> common::Result<engine::SmtpCredential> {
- auto result = credential_client_->getSmtpCredentials(credential_id, workflow_id);
- if (result.failed()) {
- return result.error();
- }
- engine::SmtpCredential cred;
- cred.host = result.value().host;
- cred.port = result.value().port;
- cred.username = result.value().username;
- cred.password = result.value().password;
- cred.security = result.value().security;
- cred.from_address = result.value().from_address;
- cred.from_name = result.value().from_name;
- return cred;
- });
- #ifdef MYSQL_SUPPORT
- // Set up MySQL credential callback for workflow engine
- engine_->setMysqlCredentialCallback(
- [this](const std::string& credential_id, const std::string& workflow_id)
- -> common::Result<engine::MysqlCredential> {
- auto result = credential_client_->getMysqlCredentials(credential_id, workflow_id);
- if (result.failed()) {
- return result.error();
- }
- engine::MysqlCredential cred;
- cred.host = result.value().host;
- cred.port = result.value().port;
- cred.username = result.value().username;
- cred.password = result.value().password;
- cred.database = result.value().database;
- cred.use_ssl = result.value().use_ssl;
- return cred;
- });
- // Set up MySQL query callback for workflow engine
- engine_->setMysqlQueryCallback(
- [this](const engine::MysqlQueryOptions& options, const std::string& workflow_id)
- -> engine::MysqlQueryResult {
- engine::MysqlQueryResult result;
- // Get MySQL credentials
- auto cred_result = credential_client_->getMysqlCredentials(options.credential_id, workflow_id);
- if (cred_result.failed()) {
- result.success = false;
- result.error = cred_result.error().message();
- return result;
- }
- // Build connection credentials
- mysql::MysqlCredentials creds;
- creds.host = cred_result.value().host;
- creds.port = cred_result.value().port;
- creds.username = cred_result.value().username;
- creds.password = cred_result.value().password;
- creds.database = options.database.empty() ? cred_result.value().database : options.database;
- creds.use_ssl = cred_result.value().use_ssl;
- // Create client and execute query
- mysql::MysqlClient client(creds);
- auto query_result = client.query(options.query, options.params);
- result.success = query_result.success;
- result.rows = query_result.rows;
- result.columns = query_result.columns;
- result.affected_rows = query_result.affected_rows;
- result.insert_id = query_result.insert_id;
- result.error = query_result.error;
- return result;
- });
- #endif
- #ifdef POSTGRESQL_SUPPORT
- // Set up PostgreSQL credential callback for workflow engine
- engine_->setPostgresqlCredentialCallback(
- [this](const std::string& credential_id, const std::string& workflow_id)
- -> common::Result<engine::PostgresqlCredential> {
- auto result = credential_client_->getPostgresqlCredentials(credential_id, workflow_id);
- if (result.failed()) {
- return result.error();
- }
- engine::PostgresqlCredential cred;
- cred.host = result.value().host;
- cred.port = result.value().port;
- cred.username = result.value().username;
- cred.password = result.value().password;
- cred.database = result.value().database;
- cred.use_ssl = result.value().use_ssl;
- return cred;
- });
- // Set up PostgreSQL query callback for workflow engine
- engine_->setPostgresqlQueryCallback(
- [this](const engine::PostgresqlQueryOptions& options, const std::string& workflow_id)
- -> engine::PostgresqlQueryResult {
- engine::PostgresqlQueryResult result;
- // Get PostgreSQL credentials
- auto cred_result = credential_client_->getPostgresqlCredentials(options.credential_id, workflow_id);
- if (cred_result.failed()) {
- result.success = false;
- result.error = cred_result.error().message();
- return result;
- }
- // Build connection credentials
- postgresql::PostgresqlCredentials creds;
- creds.host = cred_result.value().host;
- creds.port = cred_result.value().port;
- creds.username = cred_result.value().username;
- creds.password = cred_result.value().password;
- creds.database = options.database.empty() ? cred_result.value().database : options.database;
- creds.use_ssl = cred_result.value().use_ssl;
- // Create client and execute query
- postgresql::PostgresqlClient client(creds);
- auto query_result = client.query(options.query, options.params);
- result.success = query_result.success;
- result.rows = query_result.rows;
- result.columns = query_result.columns;
- result.affected_rows = query_result.affected_rows;
- result.insert_id = query_result.insert_id;
- result.error = query_result.error;
- return result;
- });
- #endif
- }
- RunnerService::~RunnerService() {
- stop();
- }
- RunnerServiceConfig RunnerService::loadConfig(const std::filesystem::path& path) {
- RunnerServiceConfig config;
- auto result = config::Config::fromFile(path);
- if (result.ok()) {
- auto& cfg = result.value();
- config.grpc_port = cfg.getOr<int>("grpc_port", 9003);
- config.runner_id = cfg.getOr<std::string>("runner_id", "runner-1");
- config.webserver_address = cfg.getOr<std::string>("webserver_address", "localhost:8080");
- config.node_sync_address = cfg.getOr<std::string>("node_sync_address", "localhost:9002");
- config.credential_service_address = cfg.getOr<std::string>("credential_service_address", "localhost:9003");
- config.database_address = cfg.getOr<std::string>("database_address", "localhost:9004");
- config.database_project = cfg.getOr<std::string>("database_project", "smartbotic-automation");
- config.max_message_size_mb = cfg.getOr<int>("max_message_size_mb", 64);
- // Node sync configuration
- config.node_registry_config.webserver_address = config.node_sync_address;
- config.node_registry_config.sync_enabled =
- cfg.getOr<bool>("node_sync.enabled", true);
- config.node_registry_config.reconnect_interval_ms =
- cfg.getOr<int>("node_sync.reconnect_interval_ms", 5000);
- config.heartbeat_interval_sec =
- cfg.getOr<int>("registration.heartbeat_interval_sec", 10);
- config.max_concurrent_executions =
- cfg.getOr<int>("registration.max_concurrent_executions", 10);
- config.workflow_engine_config.default_timeout_ms =
- cfg.getOr<int>("execution.default_timeout_ms", 60000);
- config.workflow_engine_config.script_config.max_memory_mb =
- cfg.getOr<int>("execution.max_memory_per_script_mb", 64);
- }
- return config;
- }
- void RunnerService::start() {
- if (running_) {
- return;
- }
- LOG_INFO("Starting runner service {}...", config_.runner_id);
- // Start node registry (loads nodes from webserver and subscribes to changes)
- registry_->start();
- // Create event callback to send events to webserver
- std::string webserver_url = "http://" + config_.webserver_address + "/api/v1/internal/execution-event";
- ExecutionEventCallback event_callback = [webserver_url](const std::string& event_type, const nlohmann::json& data) {
- LOG_DEBUG("Sending event {} to webserver", event_type);
- // Send event to webserver via HTTP POST (non-blocking in separate thread)
- std::thread([webserver_url, event_type, data]() {
- CURL* curl = curl_easy_init();
- if (!curl) {
- LOG_ERROR("Failed to init curl for event {}", event_type);
- return;
- }
- nlohmann::json payload;
- payload["event"] = event_type;
- payload["data"] = data;
- std::string body = payload.dump();
- curl_easy_setopt(curl, CURLOPT_URL, webserver_url.c_str());
- curl_easy_setopt(curl, CURLOPT_POST, 1L);
- curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body.c_str());
- curl_easy_setopt(curl, CURLOPT_POSTFIELDSIZE, static_cast<long>(body.size()));
- struct curl_slist* headers = nullptr;
- headers = curl_slist_append(headers, "Content-Type: application/json");
- curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
- curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
- // Discard response body (don't write to stdout)
- curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char*, size_t size, size_t nmemb, void*) -> size_t {
- return size * nmemb;
- });
- CURLcode res = curl_easy_perform(curl);
- if (res != CURLE_OK) {
- LOG_ERROR("Failed to send event {} to webserver: {}", event_type, curl_easy_strerror(res));
- }
- curl_slist_free_all(headers);
- curl_easy_cleanup(curl);
- }).detach();
- };
- // Start gRPC server
- service_impl_ = std::make_unique<RunnerServiceImpl>(*engine_, *registry_, *storage_, event_callback);
- grpc::EnableDefaultHealthCheckService(true);
- grpc::ServerBuilder builder;
- builder.AddListeningPort("0.0.0.0:" + std::to_string(config_.grpc_port),
- grpc::InsecureServerCredentials());
- builder.RegisterService(service_impl_.get());
- // grpc++ defaults an unset receive limit to 4 MB, well under a webhook
- // body the webserver now accepts up to server.max_upload_mb (32 MB by
- // default) for. ExecuteWorkflow carries that body from the webserver to
- // this server as part of the request, so without raising this the
- // gRPC hop silently re-imposes a lower cap than the HTTP one already
- // passed. Reuse max_message_size_mb - already read from config above
- // and already applied to the database client - instead of adding a
- // second knob for the same idea.
- const int max_message_bytes = config_.max_message_size_mb * 1024 * 1024;
- builder.SetMaxReceiveMessageSize(max_message_bytes);
- builder.SetMaxSendMessageSize(max_message_bytes);
- server_ = builder.BuildAndStart();
- LOG_INFO("Runner gRPC server listening on port {}", config_.grpc_port);
- running_ = true;
- // Register with webserver
- registered_ = registerWithWebServer();
- // Start heartbeat, which also re-registers whenever the webserver forgets us
- heartbeat_thread_ = std::thread(&RunnerService::heartbeatLoop, this);
- LOG_INFO("Runner service {} started", config_.runner_id);
- }
- void RunnerService::stop() {
- if (!running_) {
- return;
- }
- LOG_INFO("Stopping runner service...");
- // Signal shutdown to background threads
- {
- std::lock_guard<std::mutex> lock(shutdown_mutex_);
- running_ = false;
- }
- shutdown_cv_.notify_all();
- // Unregister
- unregisterFromWebServer();
- // Stop heartbeat (will wake up immediately now)
- if (heartbeat_thread_.joinable()) {
- heartbeat_thread_.join();
- }
- // Stop node registry
- registry_->stop();
- // Stop gRPC server
- if (server_) {
- server_->Shutdown();
- }
- LOG_INFO("Runner service stopped");
- }
- bool RunnerService::registerWithWebServer() {
- // Use HTTP to register with webserver
- nlohmann::json body;
- body["id"] = config_.runner_id;
- body["address"] = "localhost:" + std::to_string(config_.grpc_port);
- nlohmann::json capabilities;
- std::vector<std::string> node_types;
- for (const auto& node : registry_->getAllNodes()) {
- node_types.push_back(node.id);
- }
- capabilities["nodeTypes"] = node_types;
- capabilities["maxMemoryPerScriptMb"] = config_.workflow_engine_config.script_config.max_memory_mb;
- capabilities["maxExecutionTimeoutSec"] = config_.workflow_engine_config.default_timeout_ms / 1000;
- body["capabilities"] = capabilities;
- std::string url = "http://" + config_.webserver_address + "/api/internal/runners/register";
- CURL* curl = curl_easy_init();
- if (!curl) {
- LOG_WARN("Failed to initialize curl for registration");
- return false;
- }
- std::string response_data;
- std::string body_str = body.dump();
- curl_easy_setopt(curl, CURLOPT_URL, url.c_str());
- curl_easy_setopt(curl, CURLOPT_POST, 1L);
- curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body_str.c_str());
- curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
- struct curl_slist* headers = nullptr;
- headers = curl_slist_append(headers, "Content-Type: application/json");
- curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
- curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char* ptr, size_t size, size_t nmemb, void* userdata) -> size_t {
- auto* data = static_cast<std::string*>(userdata);
- data->append(ptr, size * nmemb);
- return size * nmemb;
- });
- curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response_data);
- CURLcode res = curl_easy_perform(curl);
- long http_code = 0;
- if (res == CURLE_OK) {
- curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &http_code);
- }
- // A transport-level success is not enough: the webserver can still reject the
- // registration, and reporting that as success would hide a runner that is not
- // actually reachable for work.
- const bool ok = (res == CURLE_OK && http_code >= 200 && http_code < 300);
- if (ok) {
- LOG_INFO("Runner registered with webserver");
- } else if (res != CURLE_OK) {
- LOG_WARN("Failed to register with webserver: {}", curl_easy_strerror(res));
- } else {
- LOG_WARN("Webserver rejected registration: HTTP {}", http_code);
- }
- curl_slist_free_all(headers);
- curl_easy_cleanup(curl);
- return ok;
- }
- void RunnerService::heartbeatLoop() {
- std::string url = "http://" + config_.webserver_address + "/api/internal/runners/heartbeat";
- while (running_) {
- // Wait for shutdown signal or timeout
- {
- std::unique_lock<std::mutex> lock(shutdown_mutex_);
- if (shutdown_cv_.wait_for(lock,
- std::chrono::seconds(config_.heartbeat_interval_sec),
- [this] { return !running_.load(); })) {
- // Shutdown signaled, exit loop
- break;
- }
- }
- if (!running_) break;
- auto metrics = collectMetrics();
- nlohmann::json body;
- body["id"] = config_.runner_id;
- body["status"] = "online";
- body["metrics"] = {
- {"activeExecutions", metrics.active_executions},
- {"maxExecutions", metrics.max_executions},
- {"memoryUsedBytes", metrics.memory_used_bytes},
- {"memoryTotalBytes", metrics.memory_total_bytes},
- {"cpuPercent", metrics.cpu_percent}
- };
- CURL* curl = curl_easy_init();
- if (!curl) continue;
- std::string body_str = body.dump();
- std::string response_data;
- curl_easy_setopt(curl, CURLOPT_URL, url.c_str());
- curl_easy_setopt(curl, CURLOPT_POST, 1L);
- curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body_str.c_str());
- curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
- struct curl_slist* headers = nullptr;
- headers = curl_slist_append(headers, "Content-Type: application/json");
- curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
- curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char* ptr, size_t size, size_t nmemb, void* userdata) -> size_t {
- auto* data = static_cast<std::string*>(userdata);
- data->append(ptr, size * nmemb);
- return size * nmemb;
- });
- curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response_data);
- CURLcode res = curl_easy_perform(curl);
- long http_code = 0;
- if (res == CURLE_OK) {
- curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &http_code);
- }
- curl_slist_free_all(headers);
- curl_easy_cleanup(curl);
- if (res == CURLE_OK && http_code >= 200 && http_code < 300) {
- if (!registered_) {
- LOG_INFO("Reconnected to webserver");
- registered_ = true;
- }
- continue;
- }
- // The registry lives in webserver memory, so a webserver restart forgets
- // this runner while the runner itself stays healthy. It answers 404 for an
- // unknown runner; without re-registering here the runner would stay
- // invisible and every execution would fail with "no runners available".
- if (registered_) {
- if (res != CURLE_OK) {
- LOG_WARN("Heartbeat failed: {}", curl_easy_strerror(res));
- } else {
- LOG_WARN("Heartbeat rejected: HTTP {}", http_code);
- }
- registered_ = false;
- }
- if (registerWithWebServer()) {
- registered_ = true;
- }
- }
- }
- void RunnerService::unregisterFromWebServer() {
- std::string url = "http://" + config_.webserver_address + "/api/internal/runners/unregister";
- nlohmann::json body;
- body["id"] = config_.runner_id;
- CURL* curl = curl_easy_init();
- if (!curl) return;
- std::string body_str = body.dump();
- std::string response_data;
- curl_easy_setopt(curl, CURLOPT_URL, url.c_str());
- curl_easy_setopt(curl, CURLOPT_POST, 1L);
- curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body_str.c_str());
- curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
- struct curl_slist* headers = nullptr;
- headers = curl_slist_append(headers, "Content-Type: application/json");
- curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
- curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char* ptr, size_t size, size_t nmemb, void* userdata) -> size_t {
- auto* data = static_cast<std::string*>(userdata);
- data->append(ptr, size * nmemb);
- return size * nmemb;
- });
- curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response_data);
- curl_easy_perform(curl);
- curl_slist_free_all(headers);
- curl_easy_cleanup(curl);
- }
- RunnerMetrics RunnerService::collectMetrics() {
- RunnerMetrics metrics;
- metrics.active_executions = engine_->getActiveExecutionCount();
- metrics.max_executions = config_.max_concurrent_executions;
- // Get memory usage
- struct rusage usage;
- if (getrusage(RUSAGE_SELF, &usage) == 0) {
- metrics.memory_used_bytes = usage.ru_maxrss * 1024; // KB to bytes
- }
- // Get total memory from /proc/meminfo
- std::ifstream meminfo("/proc/meminfo");
- std::string line;
- while (std::getline(meminfo, line)) {
- if (line.starts_with("MemTotal:")) {
- std::istringstream iss(line);
- std::string label;
- int64_t kb;
- iss >> label >> kb;
- metrics.memory_total_bytes = kb * 1024;
- break;
- }
- }
- return metrics;
- }
- } // namespace smartbotic::runner
|