runner_service.hpp 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156
  1. #pragma once
  2. #include <memory>
  3. #include <thread>
  4. #include <atomic>
  5. #include <mutex>
  6. #include <condition_variable>
  7. #include <grpcpp/grpcpp.h>
  8. #include "proto/runner.grpc.pb.h"
  9. #include "node_registry.hpp"
  10. #include "workflow_engine.hpp"
  11. #include "storage/storage_client.hpp"
  12. #include "credentials/credential_client.hpp"
  13. #include "config/config_loader.hpp"
  14. namespace smartbotic::runner {
  15. // Forward declaration
  16. struct RunnerMetrics;
  17. // Execution event callback type
  18. using ExecutionEventCallback = std::function<void(const std::string& event, const nlohmann::json& data)>;
  19. // gRPC Runner Service implementation
  20. class RunnerServiceImpl final : public proto::RunnerService::Service {
  21. public:
  22. RunnerServiceImpl(WorkflowEngine& engine, NodeRegistry& registry,
  23. storage::StorageClient& storage,
  24. ExecutionEventCallback event_callback = nullptr);
  25. void setEventCallback(ExecutionEventCallback callback) { event_callback_ = callback; }
  26. grpc::Status ExecuteWorkflow(grpc::ServerContext* context,
  27. const proto::ExecuteWorkflowRequest* request,
  28. proto::ExecuteWorkflowResponse* response) override;
  29. grpc::Status CancelExecution(grpc::ServerContext* context,
  30. const proto::CancelExecutionRequest* request,
  31. proto::CancelExecutionResponse* response) override;
  32. grpc::Status ListActiveExecutions(grpc::ServerContext* context,
  33. const proto::ListActiveExecutionsRequest* request,
  34. proto::ListActiveExecutionsResponse* response) override;
  35. grpc::Status ResumeExecution(grpc::ServerContext* context,
  36. const proto::ResumeExecutionRequest* request,
  37. proto::ExecuteWorkflowResponse* response) override;
  38. grpc::Status ListNodes(grpc::ServerContext* context,
  39. const proto::ListNodesRequest* request,
  40. proto::ListNodesResponse* response) override;
  41. grpc::Status ReloadNode(grpc::ServerContext* context,
  42. const proto::ReloadNodeRequest* request,
  43. proto::ReloadNodeResponse* response) override;
  44. grpc::Status ExecuteNode(grpc::ServerContext* context,
  45. const proto::ExecuteNodeRequest* request,
  46. proto::ExecuteNodeResponse* response) override;
  47. grpc::Status GetNodeCode(grpc::ServerContext* context,
  48. const proto::GetNodeCodeRequest* request,
  49. proto::GetNodeCodeResponse* response) override;
  50. grpc::Status SaveNodeCode(grpc::ServerContext* context,
  51. const proto::SaveNodeCodeRequest* request,
  52. proto::SaveNodeCodeResponse* response) override;
  53. grpc::Status CreateNode(grpc::ServerContext* context,
  54. const proto::CreateNodeRequest* request,
  55. proto::CreateNodeResponse* response) override;
  56. grpc::Status DeleteNode(grpc::ServerContext* context,
  57. const proto::DeleteNodeRequest* request,
  58. proto::DeleteNodeResponse* response) override;
  59. private:
  60. WorkflowEngine& engine_;
  61. NodeRegistry& registry_;
  62. storage::StorageClient& storage_;
  63. ExecutionEventCallback event_callback_;
  64. };
  65. // Runner service configuration
  66. struct RunnerServiceConfig {
  67. int grpc_port = 9003; // Runner's own gRPC server port
  68. std::string runner_id = "runner-1";
  69. std::string webserver_address = "localhost:8080"; // HTTP for registration
  70. std::string node_sync_address = "localhost:9002"; // gRPC for node sync
  71. std::string credential_service_address = "localhost:9003"; // gRPC for credentials
  72. std::string database_address = "localhost:9004";
  73. std::string database_project = "smartbotic-automation";
  74. int max_message_size_mb = 64; // Max gRPC message size in MB
  75. NodeRegistryConfig node_registry_config;
  76. WorkflowEngineConfig workflow_engine_config;
  77. int heartbeat_interval_sec = 10;
  78. int max_concurrent_executions = 10;
  79. };
  80. // Main runner service
  81. class RunnerService {
  82. public:
  83. explicit RunnerService(const RunnerServiceConfig& config);
  84. ~RunnerService();
  85. static RunnerServiceConfig loadConfig(const std::filesystem::path& path);
  86. void start();
  87. void stop();
  88. // Get components
  89. WorkflowEngine& engine() { return *engine_; }
  90. NodeRegistry& registry() { return *registry_; }
  91. private:
  92. // Returns false when the webserver could not be reached or rejected the
  93. // registration, so callers can retry.
  94. bool registerWithWebServer();
  95. void heartbeatLoop();
  96. void unregisterFromWebServer();
  97. RunnerMetrics collectMetrics();
  98. RunnerServiceConfig config_;
  99. std::unique_ptr<storage::StorageClient> storage_;
  100. std::unique_ptr<credentials::CredentialClient> credential_client_;
  101. std::unique_ptr<NodeRegistry> registry_;
  102. std::unique_ptr<WorkflowEngine> engine_;
  103. std::unique_ptr<RunnerServiceImpl> service_impl_;
  104. std::unique_ptr<grpc::Server> server_;
  105. std::thread heartbeat_thread_;
  106. std::atomic<bool> running_{false};
  107. // Whether the webserver currently knows about this runner. Drives
  108. // re-registration from the heartbeat loop and keeps the logs to one line per
  109. // state change rather than one per interval.
  110. std::atomic<bool> registered_{false};
  111. // Shutdown signaling
  112. std::mutex shutdown_mutex_;
  113. std::condition_variable shutdown_cv_;
  114. };
  115. // Runner metrics for heartbeat
  116. struct RunnerMetrics {
  117. int active_executions = 0;
  118. int max_executions = 10;
  119. int64_t memory_used_bytes = 0;
  120. int64_t memory_total_bytes = 0;
  121. double cpu_percent = 0.0;
  122. int64_t total_executions = 0;
  123. int64_t failed_executions = 0;
  124. };
  125. } // namespace smartbotic::runner