runner_service.cpp 36 KB

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