runner_service.cpp 30 KB

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