runner_service.cpp 41 KB

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