runner_service.cpp 42 KB

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