runner_service.cpp 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847
  1. #include "runner_service.hpp"
  2. #include "common/uuid.hpp"
  3. #include "common/time_utils.hpp"
  4. #include "logging/logger.hpp"
  5. #include "proto/runner.grpc.pb.h"
  6. #include <grpcpp/health_check_service_interface.h>
  7. #include <sys/resource.h>
  8. #include <fstream>
  9. #include <curl/curl.h>
  10. #ifdef MYSQL_SUPPORT
  11. #include "runner/mysql/mysql_client.hpp"
  12. #endif
  13. #ifdef POSTGRESQL_SUPPORT
  14. #include "runner/postgresql/postgresql_client.hpp"
  15. #endif
  16. namespace smartbotic::runner {
  17. using namespace common;
  18. // RunnerServiceImpl implementation
  19. RunnerServiceImpl::RunnerServiceImpl(WorkflowEngine& engine, NodeRegistry& registry,
  20. storage::StorageClient& storage,
  21. ExecutionEventCallback event_callback)
  22. : engine_(engine), registry_(registry), storage_(storage), event_callback_(event_callback) {}
  23. grpc::Status RunnerServiceImpl::ExecuteWorkflow(grpc::ServerContext* context,
  24. const proto::ExecuteWorkflowRequest* request,
  25. proto::ExecuteWorkflowResponse* response) {
  26. LOG_INFO("ExecuteWorkflow called for workflow: {}", request->workflow_id());
  27. // Get workflow from database
  28. auto workflow_result = storage_.get("workflows", request->workflow_id());
  29. if (workflow_result.failed()) {
  30. LOG_ERROR("Failed to get workflow from database: {}", workflow_result.error().message());
  31. response->set_status(proto::EXECUTION_STATUS_FAILED);
  32. auto* error = response->mutable_error();
  33. error->set_code(static_cast<int32_t>(workflow_result.error().code()));
  34. error->set_message(workflow_result.error().message());
  35. return grpc::Status::OK;
  36. }
  37. // Parse workflow
  38. auto workflow = Workflow::fromJson(workflow_result.value());
  39. // Parse trigger data
  40. nlohmann::json trigger_data;
  41. if (!request->trigger_data().empty()) {
  42. try {
  43. trigger_data = nlohmann::json::parse(request->trigger_data());
  44. } catch (...) {
  45. trigger_data = request->trigger_data();
  46. }
  47. }
  48. // Create callback to forward events
  49. ExecutionCallback callback;
  50. if (event_callback_) {
  51. callback = [this, &workflow](const std::string& event_type, const nlohmann::json& data) {
  52. nlohmann::json event_data = data;
  53. event_data["workflowId"] = workflow.id;
  54. event_callback_(event_type, event_data);
  55. };
  56. }
  57. // Execute workflow with callback
  58. LOG_INFO("Starting workflow execution...");
  59. auto result = engine_.execute(workflow, request->trigger_type(), trigger_data, callback);
  60. if (result.failed()) {
  61. LOG_ERROR("Workflow execution failed: {}", result.error().message());
  62. // Emit execution.failed event
  63. if (event_callback_) {
  64. event_callback_("execution.failed", {
  65. {"executionId", ""},
  66. {"workflowId", workflow.id},
  67. {"error", result.error().message()}
  68. });
  69. }
  70. response->set_status(proto::EXECUTION_STATUS_FAILED);
  71. auto* error = response->mutable_error();
  72. error->set_code(static_cast<int32_t>(result.error().code()));
  73. error->set_message(result.error().message());
  74. return grpc::Status::OK;
  75. }
  76. LOG_INFO("Workflow execution completed, execution_id: {}, status: {}",
  77. result.value().execution_id, executionStatusToString(result.value().status));
  78. // Emit final execution event based on status
  79. if (event_callback_) {
  80. auto& exec_result = result.value();
  81. if (exec_result.status == ExecutionStatus::Completed) {
  82. event_callback_("execution.completed", {
  83. {"executionId", exec_result.execution_id},
  84. {"workflowId", workflow.id},
  85. {"output", exec_result.final_output}
  86. });
  87. } else if (exec_result.status == ExecutionStatus::Failed) {
  88. // The execution already carries the message that explains the
  89. // failure, including which node produced it. Fall back to scanning
  90. // the node results only when it does not.
  91. std::string error_msg = exec_result.error;
  92. if (error_msg.empty()) {
  93. for (const auto& [node_id, node_result] : exec_result.node_results) {
  94. if (!node_result.error.empty()) {
  95. error_msg = node_result.error;
  96. break;
  97. }
  98. }
  99. }
  100. event_callback_("execution.failed", {
  101. {"executionId", exec_result.execution_id},
  102. {"workflowId", workflow.id},
  103. {"error", error_msg}
  104. });
  105. }
  106. }
  107. response->set_execution_id(result.value().execution_id);
  108. response->set_status(static_cast<proto::ExecutionStatus>(result.value().status));
  109. if (request->wait_for_completion()) {
  110. response->set_result(result.value().final_output.dump());
  111. }
  112. return grpc::Status::OK;
  113. }
  114. grpc::Status RunnerServiceImpl::CancelExecution(grpc::ServerContext* context,
  115. const proto::CancelExecutionRequest* request,
  116. proto::Empty* response) {
  117. engine_.cancelExecution(request->execution_id());
  118. return grpc::Status::OK;
  119. }
  120. grpc::Status RunnerServiceImpl::ListNodes(grpc::ServerContext* context,
  121. const proto::ListNodesRequest* request,
  122. proto::ListNodesResponse* response) {
  123. std::vector<NodeDefinition> nodes;
  124. if (request->category().empty()) {
  125. nodes = registry_.getAllNodes();
  126. } else {
  127. nodes = registry_.getNodesByCategory(request->category());
  128. }
  129. for (const auto& node : nodes) {
  130. auto* proto_node = response->add_nodes();
  131. proto_node->set_id(node.id);
  132. proto_node->set_name(node.name);
  133. proto_node->set_category(node.category);
  134. proto_node->set_version(node.version);
  135. proto_node->set_description(node.description);
  136. proto_node->set_icon(node.icon);
  137. proto_node->set_is_trigger(node.is_trigger);
  138. proto_node->set_config_schema(node.config_schema.dump());
  139. proto_node->set_input_schema(node.input_schema.dump());
  140. proto_node->set_output_schema(node.output_schema.dump());
  141. for (const auto& input : node.inputs) {
  142. auto* proto_input = proto_node->add_inputs();
  143. proto_input->set_name(input.name);
  144. proto_input->set_display_name(input.display_name);
  145. proto_input->set_type(input.type);
  146. proto_input->set_required(input.required);
  147. }
  148. for (const auto& output : node.outputs) {
  149. auto* proto_output = proto_node->add_outputs();
  150. proto_output->set_name(output.name);
  151. proto_output->set_display_name(output.display_name);
  152. proto_output->set_type(output.type);
  153. if (!output.color.empty()) {
  154. proto_output->set_color(output.color);
  155. }
  156. }
  157. }
  158. return grpc::Status::OK;
  159. }
  160. grpc::Status RunnerServiceImpl::ReloadNode(grpc::ServerContext* context,
  161. const proto::ReloadNodeRequest* request,
  162. proto::ReloadNodeResponse* response) {
  163. // Nodes are now synced automatically from webserver
  164. // Manual reload is no longer needed
  165. auto node = registry_.getNode(request->node_id());
  166. if (!node) {
  167. response->set_success(false);
  168. auto* error = response->mutable_error();
  169. error->set_code(404);
  170. error->set_message("Node not found: " + request->node_id());
  171. return grpc::Status::OK;
  172. }
  173. response->set_success(true);
  174. auto* proto_node = response->mutable_node();
  175. proto_node->set_id(node->id);
  176. proto_node->set_name(node->name);
  177. proto_node->set_version(node->version);
  178. return grpc::Status::OK;
  179. }
  180. grpc::Status RunnerServiceImpl::ExecuteNode(grpc::ServerContext* context,
  181. const proto::ExecuteNodeRequest* request,
  182. proto::ExecuteNodeResponse* response) {
  183. auto node_def = registry_.getNode(request->node_type());
  184. if (!node_def) {
  185. response->set_success(false);
  186. response->set_error("Node type not found: " + request->node_type());
  187. return grpc::Status::OK;
  188. }
  189. // Parse input and config
  190. nlohmann::json input, config;
  191. try {
  192. if (!request->input().empty()) {
  193. input = nlohmann::json::parse(request->input());
  194. }
  195. if (!request->config().empty()) {
  196. config = nlohmann::json::parse(request->config());
  197. }
  198. } catch (const std::exception& e) {
  199. response->set_success(false);
  200. response->set_error("Invalid JSON: " + std::string(e.what()));
  201. return grpc::Status::OK;
  202. }
  203. // Create execution context
  204. engine::ScriptContext ctx;
  205. ctx.execution_id = common::UUID::generate();
  206. ctx.node_id = request->node_type();
  207. ctx.input = input;
  208. ctx.config = config;
  209. // Execute
  210. engine::ScriptEnginePool pool(1);
  211. auto* engine = pool.acquire();
  212. auto result = engine->execute(node_def->code, ctx);
  213. pool.release(engine);
  214. response->set_success(result.success);
  215. response->set_output(result.output.dump());
  216. response->set_error(result.error);
  217. response->set_execution_time_ms(result.execution_time_ms);
  218. return grpc::Status::OK;
  219. }
  220. grpc::Status RunnerServiceImpl::GetNodeCode(grpc::ServerContext* context,
  221. const proto::GetNodeCodeRequest* request,
  222. proto::GetNodeCodeResponse* response) {
  223. auto result = registry_.getNodeCode(request->node_id());
  224. if (result.failed()) {
  225. response->set_success(false);
  226. auto* error = response->mutable_error();
  227. error->set_code(static_cast<int32_t>(result.error().code()));
  228. error->set_message(result.error().message());
  229. return grpc::Status::OK;
  230. }
  231. response->set_success(true);
  232. response->set_code(result.value());
  233. // file_path no longer applicable - nodes stored in database
  234. return grpc::Status::OK;
  235. }
  236. grpc::Status RunnerServiceImpl::SaveNodeCode(grpc::ServerContext* context,
  237. const proto::SaveNodeCodeRequest* request,
  238. proto::SaveNodeCodeResponse* response) {
  239. // Node modifications are now handled centrally by the webserver
  240. response->set_success(false);
  241. auto* error = response->mutable_error();
  242. error->set_code(501);
  243. error->set_message("Node modifications should be done through the webserver API");
  244. return grpc::Status::OK;
  245. }
  246. grpc::Status RunnerServiceImpl::CreateNode(grpc::ServerContext* context,
  247. const proto::CreateNodeRequest* request,
  248. proto::CreateNodeResponse* response) {
  249. // Node creation is now handled centrally by the webserver
  250. response->set_success(false);
  251. auto* error = response->mutable_error();
  252. error->set_code(501);
  253. error->set_message("Node creation should be done through the webserver API");
  254. return grpc::Status::OK;
  255. }
  256. grpc::Status RunnerServiceImpl::DeleteNode(grpc::ServerContext* context,
  257. const proto::DeleteNodeRequest* request,
  258. proto::DeleteNodeResponse* response) {
  259. // Node deletion is now handled centrally by the webserver
  260. response->set_success(false);
  261. auto* error = response->mutable_error();
  262. error->set_code(501);
  263. error->set_message("Node deletion should be done through the webserver API");
  264. return grpc::Status::OK;
  265. }
  266. // RunnerService implementation
  267. RunnerService::RunnerService(const RunnerServiceConfig& config)
  268. : config_(config) {
  269. // Initialize storage client
  270. storage::StorageClientConfig storage_config;
  271. storage_config.address = config_.database_address;
  272. storage_config.project = config_.database_project;
  273. storage_config.max_message_size_mb = config_.max_message_size_mb;
  274. storage_ = std::make_unique<storage::StorageClient>(storage_config);
  275. // Initialize credential client
  276. credentials::CredentialClientConfig cred_config;
  277. cred_config.address = config_.credential_service_address;
  278. credential_client_ = std::make_unique<credentials::CredentialClient>(cred_config);
  279. // Initialize node registry - now loads from webserver
  280. config_.node_registry_config.webserver_address = config_.node_sync_address;
  281. registry_ = std::make_unique<NodeRegistry>(config_.node_registry_config);
  282. // Initialize workflow engine
  283. config_.workflow_engine_config.max_concurrent_executions = config_.max_concurrent_executions;
  284. config_.workflow_engine_config.runner_id = config_.runner_id;
  285. engine_ = std::make_unique<WorkflowEngine>(*registry_, *storage_, config_.workflow_engine_config);
  286. // Set up credential auth callback for workflow engine
  287. engine_->setCredentialAuthCallback(
  288. [this](const std::string& credential_id, const std::string& workflow_id)
  289. -> common::Result<engine::CredentialAuth> {
  290. auto result = credential_client_->getHttpAuth(credential_id, workflow_id);
  291. if (result.failed()) {
  292. return result.error();
  293. }
  294. engine::CredentialAuth auth;
  295. auth.header_name = result.value().header_name;
  296. auth.header_value = result.value().header_value;
  297. return auth;
  298. });
  299. // Set up IMAP credential callback for workflow engine
  300. engine_->setImapCredentialCallback(
  301. [this](const std::string& credential_id, const std::string& workflow_id)
  302. -> common::Result<engine::ImapCredential> {
  303. auto result = credential_client_->getImapCredentials(credential_id, workflow_id);
  304. if (result.failed()) {
  305. return result.error();
  306. }
  307. engine::ImapCredential cred;
  308. cred.host = result.value().host;
  309. cred.port = result.value().port;
  310. cred.username = result.value().username;
  311. cred.password = result.value().password;
  312. cred.use_ssl = result.value().use_ssl;
  313. return cred;
  314. });
  315. #ifdef MYSQL_SUPPORT
  316. // Set up MySQL credential callback for workflow engine
  317. engine_->setMysqlCredentialCallback(
  318. [this](const std::string& credential_id, const std::string& workflow_id)
  319. -> common::Result<engine::MysqlCredential> {
  320. auto result = credential_client_->getMysqlCredentials(credential_id, workflow_id);
  321. if (result.failed()) {
  322. return result.error();
  323. }
  324. engine::MysqlCredential cred;
  325. cred.host = result.value().host;
  326. cred.port = result.value().port;
  327. cred.username = result.value().username;
  328. cred.password = result.value().password;
  329. cred.database = result.value().database;
  330. cred.use_ssl = result.value().use_ssl;
  331. return cred;
  332. });
  333. // Set up MySQL query callback for workflow engine
  334. engine_->setMysqlQueryCallback(
  335. [this](const engine::MysqlQueryOptions& options, const std::string& workflow_id)
  336. -> engine::MysqlQueryResult {
  337. engine::MysqlQueryResult result;
  338. // Get MySQL credentials
  339. auto cred_result = credential_client_->getMysqlCredentials(options.credential_id, workflow_id);
  340. if (cred_result.failed()) {
  341. result.success = false;
  342. result.error = cred_result.error().message();
  343. return result;
  344. }
  345. // Build connection credentials
  346. mysql::MysqlCredentials creds;
  347. creds.host = cred_result.value().host;
  348. creds.port = cred_result.value().port;
  349. creds.username = cred_result.value().username;
  350. creds.password = cred_result.value().password;
  351. creds.database = options.database.empty() ? cred_result.value().database : options.database;
  352. creds.use_ssl = cred_result.value().use_ssl;
  353. // Create client and execute query
  354. mysql::MysqlClient client(creds);
  355. auto query_result = client.query(options.query, options.params);
  356. result.success = query_result.success;
  357. result.rows = query_result.rows;
  358. result.columns = query_result.columns;
  359. result.affected_rows = query_result.affected_rows;
  360. result.insert_id = query_result.insert_id;
  361. result.error = query_result.error;
  362. return result;
  363. });
  364. #endif
  365. #ifdef POSTGRESQL_SUPPORT
  366. // Set up PostgreSQL credential callback for workflow engine
  367. engine_->setPostgresqlCredentialCallback(
  368. [this](const std::string& credential_id, const std::string& workflow_id)
  369. -> common::Result<engine::PostgresqlCredential> {
  370. auto result = credential_client_->getPostgresqlCredentials(credential_id, workflow_id);
  371. if (result.failed()) {
  372. return result.error();
  373. }
  374. engine::PostgresqlCredential cred;
  375. cred.host = result.value().host;
  376. cred.port = result.value().port;
  377. cred.username = result.value().username;
  378. cred.password = result.value().password;
  379. cred.database = result.value().database;
  380. cred.use_ssl = result.value().use_ssl;
  381. return cred;
  382. });
  383. // Set up PostgreSQL query callback for workflow engine
  384. engine_->setPostgresqlQueryCallback(
  385. [this](const engine::PostgresqlQueryOptions& options, const std::string& workflow_id)
  386. -> engine::PostgresqlQueryResult {
  387. engine::PostgresqlQueryResult result;
  388. // Get PostgreSQL credentials
  389. auto cred_result = credential_client_->getPostgresqlCredentials(options.credential_id, workflow_id);
  390. if (cred_result.failed()) {
  391. result.success = false;
  392. result.error = cred_result.error().message();
  393. return result;
  394. }
  395. // Build connection credentials
  396. postgresql::PostgresqlCredentials creds;
  397. creds.host = cred_result.value().host;
  398. creds.port = cred_result.value().port;
  399. creds.username = cred_result.value().username;
  400. creds.password = cred_result.value().password;
  401. creds.database = options.database.empty() ? cred_result.value().database : options.database;
  402. creds.use_ssl = cred_result.value().use_ssl;
  403. // Create client and execute query
  404. postgresql::PostgresqlClient client(creds);
  405. auto query_result = client.query(options.query, options.params);
  406. result.success = query_result.success;
  407. result.rows = query_result.rows;
  408. result.columns = query_result.columns;
  409. result.affected_rows = query_result.affected_rows;
  410. result.insert_id = query_result.insert_id;
  411. result.error = query_result.error;
  412. return result;
  413. });
  414. #endif
  415. }
  416. RunnerService::~RunnerService() {
  417. stop();
  418. }
  419. RunnerServiceConfig RunnerService::loadConfig(const std::filesystem::path& path) {
  420. RunnerServiceConfig config;
  421. auto result = config::Config::fromFile(path);
  422. if (result.ok()) {
  423. auto& cfg = result.value();
  424. config.grpc_port = cfg.getOr<int>("grpc_port", 9003);
  425. config.runner_id = cfg.getOr<std::string>("runner_id", "runner-1");
  426. config.webserver_address = cfg.getOr<std::string>("webserver_address", "localhost:8080");
  427. config.node_sync_address = cfg.getOr<std::string>("node_sync_address", "localhost:9002");
  428. config.credential_service_address = cfg.getOr<std::string>("credential_service_address", "localhost:9003");
  429. config.database_address = cfg.getOr<std::string>("database_address", "localhost:9004");
  430. config.database_project = cfg.getOr<std::string>("database_project", "smartbotic-automation");
  431. config.max_message_size_mb = cfg.getOr<int>("max_message_size_mb", 64);
  432. // Node sync configuration
  433. config.node_registry_config.webserver_address = config.node_sync_address;
  434. config.node_registry_config.sync_enabled =
  435. cfg.getOr<bool>("node_sync.enabled", true);
  436. config.node_registry_config.reconnect_interval_ms =
  437. cfg.getOr<int>("node_sync.reconnect_interval_ms", 5000);
  438. config.heartbeat_interval_sec =
  439. cfg.getOr<int>("registration.heartbeat_interval_sec", 10);
  440. config.max_concurrent_executions =
  441. cfg.getOr<int>("registration.max_concurrent_executions", 10);
  442. config.workflow_engine_config.default_timeout_ms =
  443. cfg.getOr<int>("execution.default_timeout_ms", 60000);
  444. config.workflow_engine_config.script_config.max_memory_mb =
  445. cfg.getOr<int>("execution.max_memory_per_script_mb", 64);
  446. }
  447. return config;
  448. }
  449. void RunnerService::start() {
  450. if (running_) {
  451. return;
  452. }
  453. LOG_INFO("Starting runner service {}...", config_.runner_id);
  454. // Start node registry (loads nodes from webserver and subscribes to changes)
  455. registry_->start();
  456. // Create event callback to send events to webserver
  457. std::string webserver_url = "http://" + config_.webserver_address + "/api/v1/internal/execution-event";
  458. ExecutionEventCallback event_callback = [webserver_url](const std::string& event_type, const nlohmann::json& data) {
  459. LOG_DEBUG("Sending event {} to webserver", event_type);
  460. // Send event to webserver via HTTP POST (non-blocking in separate thread)
  461. std::thread([webserver_url, event_type, data]() {
  462. CURL* curl = curl_easy_init();
  463. if (!curl) {
  464. LOG_ERROR("Failed to init curl for event {}", event_type);
  465. return;
  466. }
  467. nlohmann::json payload;
  468. payload["event"] = event_type;
  469. payload["data"] = data;
  470. std::string body = payload.dump();
  471. curl_easy_setopt(curl, CURLOPT_URL, webserver_url.c_str());
  472. curl_easy_setopt(curl, CURLOPT_POST, 1L);
  473. curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body.c_str());
  474. curl_easy_setopt(curl, CURLOPT_POSTFIELDSIZE, static_cast<long>(body.size()));
  475. struct curl_slist* headers = nullptr;
  476. headers = curl_slist_append(headers, "Content-Type: application/json");
  477. curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
  478. curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
  479. // Discard response body (don't write to stdout)
  480. curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char*, size_t size, size_t nmemb, void*) -> size_t {
  481. return size * nmemb;
  482. });
  483. CURLcode res = curl_easy_perform(curl);
  484. if (res != CURLE_OK) {
  485. LOG_ERROR("Failed to send event {} to webserver: {}", event_type, curl_easy_strerror(res));
  486. }
  487. curl_slist_free_all(headers);
  488. curl_easy_cleanup(curl);
  489. }).detach();
  490. };
  491. // Start gRPC server
  492. service_impl_ = std::make_unique<RunnerServiceImpl>(*engine_, *registry_, *storage_, event_callback);
  493. grpc::EnableDefaultHealthCheckService(true);
  494. grpc::ServerBuilder builder;
  495. builder.AddListeningPort("0.0.0.0:" + std::to_string(config_.grpc_port),
  496. grpc::InsecureServerCredentials());
  497. builder.RegisterService(service_impl_.get());
  498. server_ = builder.BuildAndStart();
  499. LOG_INFO("Runner gRPC server listening on port {}", config_.grpc_port);
  500. running_ = true;
  501. // Register with webserver
  502. registered_ = registerWithWebServer();
  503. // Start heartbeat, which also re-registers whenever the webserver forgets us
  504. heartbeat_thread_ = std::thread(&RunnerService::heartbeatLoop, this);
  505. LOG_INFO("Runner service {} started", config_.runner_id);
  506. }
  507. void RunnerService::stop() {
  508. if (!running_) {
  509. return;
  510. }
  511. LOG_INFO("Stopping runner service...");
  512. // Signal shutdown to background threads
  513. {
  514. std::lock_guard<std::mutex> lock(shutdown_mutex_);
  515. running_ = false;
  516. }
  517. shutdown_cv_.notify_all();
  518. // Unregister
  519. unregisterFromWebServer();
  520. // Stop heartbeat (will wake up immediately now)
  521. if (heartbeat_thread_.joinable()) {
  522. heartbeat_thread_.join();
  523. }
  524. // Stop node registry
  525. registry_->stop();
  526. // Stop gRPC server
  527. if (server_) {
  528. server_->Shutdown();
  529. }
  530. LOG_INFO("Runner service stopped");
  531. }
  532. bool RunnerService::registerWithWebServer() {
  533. // Use HTTP to register with webserver
  534. nlohmann::json body;
  535. body["id"] = config_.runner_id;
  536. body["address"] = "localhost:" + std::to_string(config_.grpc_port);
  537. nlohmann::json capabilities;
  538. std::vector<std::string> node_types;
  539. for (const auto& node : registry_->getAllNodes()) {
  540. node_types.push_back(node.id);
  541. }
  542. capabilities["nodeTypes"] = node_types;
  543. capabilities["maxMemoryPerScriptMb"] = config_.workflow_engine_config.script_config.max_memory_mb;
  544. capabilities["maxExecutionTimeoutSec"] = config_.workflow_engine_config.default_timeout_ms / 1000;
  545. body["capabilities"] = capabilities;
  546. std::string url = "http://" + config_.webserver_address + "/api/internal/runners/register";
  547. CURL* curl = curl_easy_init();
  548. if (!curl) {
  549. LOG_WARN("Failed to initialize curl for registration");
  550. return false;
  551. }
  552. std::string response_data;
  553. std::string body_str = body.dump();
  554. curl_easy_setopt(curl, CURLOPT_URL, url.c_str());
  555. curl_easy_setopt(curl, CURLOPT_POST, 1L);
  556. curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body_str.c_str());
  557. curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
  558. struct curl_slist* headers = nullptr;
  559. headers = curl_slist_append(headers, "Content-Type: application/json");
  560. curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
  561. curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char* ptr, size_t size, size_t nmemb, void* userdata) -> size_t {
  562. auto* data = static_cast<std::string*>(userdata);
  563. data->append(ptr, size * nmemb);
  564. return size * nmemb;
  565. });
  566. curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response_data);
  567. CURLcode res = curl_easy_perform(curl);
  568. long http_code = 0;
  569. if (res == CURLE_OK) {
  570. curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &http_code);
  571. }
  572. // A transport-level success is not enough: the webserver can still reject the
  573. // registration, and reporting that as success would hide a runner that is not
  574. // actually reachable for work.
  575. const bool ok = (res == CURLE_OK && http_code >= 200 && http_code < 300);
  576. if (ok) {
  577. LOG_INFO("Runner registered with webserver");
  578. } else if (res != CURLE_OK) {
  579. LOG_WARN("Failed to register with webserver: {}", curl_easy_strerror(res));
  580. } else {
  581. LOG_WARN("Webserver rejected registration: HTTP {}", http_code);
  582. }
  583. curl_slist_free_all(headers);
  584. curl_easy_cleanup(curl);
  585. return ok;
  586. }
  587. void RunnerService::heartbeatLoop() {
  588. std::string url = "http://" + config_.webserver_address + "/api/internal/runners/heartbeat";
  589. while (running_) {
  590. // Wait for shutdown signal or timeout
  591. {
  592. std::unique_lock<std::mutex> lock(shutdown_mutex_);
  593. if (shutdown_cv_.wait_for(lock,
  594. std::chrono::seconds(config_.heartbeat_interval_sec),
  595. [this] { return !running_.load(); })) {
  596. // Shutdown signaled, exit loop
  597. break;
  598. }
  599. }
  600. if (!running_) break;
  601. auto metrics = collectMetrics();
  602. nlohmann::json body;
  603. body["id"] = config_.runner_id;
  604. body["status"] = "online";
  605. body["metrics"] = {
  606. {"activeExecutions", metrics.active_executions},
  607. {"maxExecutions", metrics.max_executions},
  608. {"memoryUsedBytes", metrics.memory_used_bytes},
  609. {"memoryTotalBytes", metrics.memory_total_bytes},
  610. {"cpuPercent", metrics.cpu_percent}
  611. };
  612. CURL* curl = curl_easy_init();
  613. if (!curl) continue;
  614. std::string body_str = body.dump();
  615. std::string response_data;
  616. curl_easy_setopt(curl, CURLOPT_URL, url.c_str());
  617. curl_easy_setopt(curl, CURLOPT_POST, 1L);
  618. curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body_str.c_str());
  619. curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
  620. struct curl_slist* headers = nullptr;
  621. headers = curl_slist_append(headers, "Content-Type: application/json");
  622. curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
  623. curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char* ptr, size_t size, size_t nmemb, void* userdata) -> size_t {
  624. auto* data = static_cast<std::string*>(userdata);
  625. data->append(ptr, size * nmemb);
  626. return size * nmemb;
  627. });
  628. curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response_data);
  629. CURLcode res = curl_easy_perform(curl);
  630. long http_code = 0;
  631. if (res == CURLE_OK) {
  632. curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &http_code);
  633. }
  634. curl_slist_free_all(headers);
  635. curl_easy_cleanup(curl);
  636. if (res == CURLE_OK && http_code >= 200 && http_code < 300) {
  637. if (!registered_) {
  638. LOG_INFO("Reconnected to webserver");
  639. registered_ = true;
  640. }
  641. continue;
  642. }
  643. // The registry lives in webserver memory, so a webserver restart forgets
  644. // this runner while the runner itself stays healthy. It answers 404 for an
  645. // unknown runner; without re-registering here the runner would stay
  646. // invisible and every execution would fail with "no runners available".
  647. if (registered_) {
  648. if (res != CURLE_OK) {
  649. LOG_WARN("Heartbeat failed: {}", curl_easy_strerror(res));
  650. } else {
  651. LOG_WARN("Heartbeat rejected: HTTP {}", http_code);
  652. }
  653. registered_ = false;
  654. }
  655. if (registerWithWebServer()) {
  656. registered_ = true;
  657. }
  658. }
  659. }
  660. void RunnerService::unregisterFromWebServer() {
  661. std::string url = "http://" + config_.webserver_address + "/api/internal/runners/unregister";
  662. nlohmann::json body;
  663. body["id"] = config_.runner_id;
  664. CURL* curl = curl_easy_init();
  665. if (!curl) return;
  666. std::string body_str = body.dump();
  667. std::string response_data;
  668. curl_easy_setopt(curl, CURLOPT_URL, url.c_str());
  669. curl_easy_setopt(curl, CURLOPT_POST, 1L);
  670. curl_easy_setopt(curl, CURLOPT_POSTFIELDS, body_str.c_str());
  671. curl_easy_setopt(curl, CURLOPT_TIMEOUT, 5L);
  672. struct curl_slist* headers = nullptr;
  673. headers = curl_slist_append(headers, "Content-Type: application/json");
  674. curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
  675. curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, +[](char* ptr, size_t size, size_t nmemb, void* userdata) -> size_t {
  676. auto* data = static_cast<std::string*>(userdata);
  677. data->append(ptr, size * nmemb);
  678. return size * nmemb;
  679. });
  680. curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response_data);
  681. curl_easy_perform(curl);
  682. curl_slist_free_all(headers);
  683. curl_easy_cleanup(curl);
  684. }
  685. RunnerMetrics RunnerService::collectMetrics() {
  686. RunnerMetrics metrics;
  687. metrics.active_executions = engine_->getActiveExecutionCount();
  688. metrics.max_executions = config_.max_concurrent_executions;
  689. // Get memory usage
  690. struct rusage usage;
  691. if (getrusage(RUSAGE_SELF, &usage) == 0) {
  692. metrics.memory_used_bytes = usage.ru_maxrss * 1024; // KB to bytes
  693. }
  694. // Get total memory from /proc/meminfo
  695. std::ifstream meminfo("/proc/meminfo");
  696. std::string line;
  697. while (std::getline(meminfo, line)) {
  698. if (line.starts_with("MemTotal:")) {
  699. std::istringstream iss(line);
  700. std::string label;
  701. int64_t kb;
  702. iss >> label >> kb;
  703. metrics.memory_total_bytes = kb * 1024;
  704. break;
  705. }
  706. }
  707. return metrics;
  708. }
  709. } // namespace smartbotic::runner