runner_service.cpp 44 KB

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