#include "runner_service.hpp" #include "common/uuid.hpp" #include "common/time_utils.hpp" #include "logging/logger.hpp" #include "proto/runner.grpc.pb.h" #include #include #include #include #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()); // Get workflow from database auto workflow_result = storage_.get("workflows", request->workflow_id()); 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(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(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)); // Emit final execution event based on status if (event_callback_) { auto& exec_result = result.value(); if (exec_result.status == ExecutionStatus::Completed) { event_callback_("execution.completed", { {"executionId", exec_result.execution_id}, {"workflowId", workflow.id}, {"output", exec_result.final_output} }); } else if (exec_result.status == ExecutionStatus::Failed) { // The execution already carries the message that explains the // failure, including which node produced it. Fall back to scanning // the node results only when it does not. std::string error_msg = exec_result.error; if (error_msg.empty()) { for (const auto& [node_id, node_result] : exec_result.node_results) { if (!node_result.error.empty()) { error_msg = node_result.error; break; } } } event_callback_("execution.failed", { {"executionId", exec_result.execution_id}, {"workflowId", workflow.id}, {"error", error_msg} }); } } 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::Empty* response) { engine_.cancelExecution(request->execution_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 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(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_config); // Initialize credential client credentials::CredentialClientConfig cred_config; cred_config.address = config_.credential_service_address; credential_client_ = std::make_unique(cred_config); // Initialize node registry - now loads from webserver config_.node_registry_config.webserver_address = config_.node_sync_address; registry_ = std::make_unique(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(*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 { 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; }); // Set up IMAP credential callback for workflow engine engine_->setImapCredentialCallback( [this](const std::string& credential_id, const std::string& workflow_id) -> common::Result { 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 { 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 { 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 { 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("grpc_port", 9003); config.runner_id = cfg.getOr("runner_id", "runner-1"); config.webserver_address = cfg.getOr("webserver_address", "localhost:8080"); config.node_sync_address = cfg.getOr("node_sync_address", "localhost:9002"); config.credential_service_address = cfg.getOr("credential_service_address", "localhost:9003"); config.database_address = cfg.getOr("database_address", "localhost:9004"); config.database_project = cfg.getOr("database_project", "smartbotic-automation"); config.max_message_size_mb = cfg.getOr("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("node_sync.enabled", true); config.node_registry_config.reconnect_interval_ms = cfg.getOr("node_sync.reconnect_interval_ms", 5000); config.heartbeat_interval_sec = cfg.getOr("registration.heartbeat_interval_sec", 10); config.max_concurrent_executions = cfg.getOr("registration.max_concurrent_executions", 10); config.workflow_engine_config.default_timeout_ms = cfg.getOr("execution.default_timeout_ms", 60000); config.workflow_engine_config.script_config.max_memory_mb = cfg.getOr("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(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(*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()); 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 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 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(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 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(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(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