| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156 |
- #pragma once
- #include <memory>
- #include <thread>
- #include <atomic>
- #include <mutex>
- #include <condition_variable>
- #include <grpcpp/grpcpp.h>
- #include "proto/runner.grpc.pb.h"
- #include "node_registry.hpp"
- #include "workflow_engine.hpp"
- #include "storage/storage_client.hpp"
- #include "credentials/credential_client.hpp"
- #include "config/config_loader.hpp"
- namespace smartbotic::runner {
- // Forward declaration
- struct RunnerMetrics;
- // Execution event callback type
- using ExecutionEventCallback = std::function<void(const std::string& event, const nlohmann::json& data)>;
- // gRPC Runner Service implementation
- class RunnerServiceImpl final : public proto::RunnerService::Service {
- public:
- RunnerServiceImpl(WorkflowEngine& engine, NodeRegistry& registry,
- storage::StorageClient& storage,
- ExecutionEventCallback event_callback = nullptr);
- void setEventCallback(ExecutionEventCallback callback) { event_callback_ = callback; }
- grpc::Status ExecuteWorkflow(grpc::ServerContext* context,
- const proto::ExecuteWorkflowRequest* request,
- proto::ExecuteWorkflowResponse* response) override;
- grpc::Status CancelExecution(grpc::ServerContext* context,
- const proto::CancelExecutionRequest* request,
- proto::CancelExecutionResponse* response) override;
- grpc::Status ListActiveExecutions(grpc::ServerContext* context,
- const proto::ListActiveExecutionsRequest* request,
- proto::ListActiveExecutionsResponse* response) override;
- grpc::Status ResumeExecution(grpc::ServerContext* context,
- const proto::ResumeExecutionRequest* request,
- proto::ExecuteWorkflowResponse* response) override;
- grpc::Status ListNodes(grpc::ServerContext* context,
- const proto::ListNodesRequest* request,
- proto::ListNodesResponse* response) override;
- grpc::Status ReloadNode(grpc::ServerContext* context,
- const proto::ReloadNodeRequest* request,
- proto::ReloadNodeResponse* response) override;
- grpc::Status ExecuteNode(grpc::ServerContext* context,
- const proto::ExecuteNodeRequest* request,
- proto::ExecuteNodeResponse* response) override;
- grpc::Status GetNodeCode(grpc::ServerContext* context,
- const proto::GetNodeCodeRequest* request,
- proto::GetNodeCodeResponse* response) override;
- grpc::Status SaveNodeCode(grpc::ServerContext* context,
- const proto::SaveNodeCodeRequest* request,
- proto::SaveNodeCodeResponse* response) override;
- grpc::Status CreateNode(grpc::ServerContext* context,
- const proto::CreateNodeRequest* request,
- proto::CreateNodeResponse* response) override;
- grpc::Status DeleteNode(grpc::ServerContext* context,
- const proto::DeleteNodeRequest* request,
- proto::DeleteNodeResponse* response) override;
- private:
- WorkflowEngine& engine_;
- NodeRegistry& registry_;
- storage::StorageClient& storage_;
- ExecutionEventCallback event_callback_;
- };
- // Runner service configuration
- struct RunnerServiceConfig {
- int grpc_port = 9003; // Runner's own gRPC server port
- std::string runner_id = "runner-1";
- std::string webserver_address = "localhost:8080"; // HTTP for registration
- std::string node_sync_address = "localhost:9002"; // gRPC for node sync
- std::string credential_service_address = "localhost:9003"; // gRPC for credentials
- std::string database_address = "localhost:9004";
- std::string database_project = "smartbotic-automation";
- int max_message_size_mb = 64; // Max gRPC message size in MB
- NodeRegistryConfig node_registry_config;
- WorkflowEngineConfig workflow_engine_config;
- int heartbeat_interval_sec = 10;
- int max_concurrent_executions = 10;
- };
- // Main runner service
- class RunnerService {
- public:
- explicit RunnerService(const RunnerServiceConfig& config);
- ~RunnerService();
- static RunnerServiceConfig loadConfig(const std::filesystem::path& path);
- void start();
- void stop();
- // Get components
- WorkflowEngine& engine() { return *engine_; }
- NodeRegistry& registry() { return *registry_; }
- private:
- // Returns false when the webserver could not be reached or rejected the
- // registration, so callers can retry.
- bool registerWithWebServer();
- void heartbeatLoop();
- void unregisterFromWebServer();
- RunnerMetrics collectMetrics();
- RunnerServiceConfig config_;
- std::unique_ptr<storage::StorageClient> storage_;
- std::unique_ptr<credentials::CredentialClient> credential_client_;
- std::unique_ptr<NodeRegistry> registry_;
- std::unique_ptr<WorkflowEngine> engine_;
- std::unique_ptr<RunnerServiceImpl> service_impl_;
- std::unique_ptr<grpc::Server> server_;
- std::thread heartbeat_thread_;
- std::atomic<bool> running_{false};
- // Whether the webserver currently knows about this runner. Drives
- // re-registration from the heartbeat loop and keeps the logs to one line per
- // state change rather than one per interval.
- std::atomic<bool> registered_{false};
- // Shutdown signaling
- std::mutex shutdown_mutex_;
- std::condition_variable shutdown_cv_;
- };
- // Runner metrics for heartbeat
- struct RunnerMetrics {
- int active_executions = 0;
- int max_executions = 10;
- int64_t memory_used_bytes = 0;
- int64_t memory_total_bytes = 0;
- double cpu_percent = 0.0;
- int64_t total_executions = 0;
- int64_t failed_executions = 0;
- };
- } // namespace smartbotic::runner
|