Ver código fonte

feat: close executions a restarted runner cannot still be running

A runner keeps what it is executing in memory, so a restart leaves every
in-flight execution recorded as running with nobody left to finish it. Two of
them sat that way for fourteen hours. Pressing Stop closes one now, but only
because somebody noticed it.

When a runner registers, the executions recorded against it are checked and
the ones it is not running are closed, with the reason on the record rather
than a bare "cancelled".

It asks the runner rather than assuming, which needed a new RPC. The obvious
shortcut - a runner that just registered cannot be running anything, so close
everything - is wrong the moment a runner re-registers while working, and
would close live executions underneath it. ListActiveExecutions returns what
the engine actually has in flight, and only what is missing from that list is
closed. If the runner cannot be reached, nothing is reconciled: without an
answer there is no way to tell an orphan from a live run, and closing a live
one is worse than leaving a stale record.

Waiting is deliberately left alone. That is a run parked on a person, not on
a runner, and it is meant to outlive one.

Both halves measured. A workflow left mid-wait, then only the runner
restarted: the record closed itself and the log said so. A live run, then a
genuine re-registration of the same runner: it kept running. The first attempt
at that second test posted to the wrong path and got a 404, so it proved
nothing until it was redone against the route that exists.

66 passed.
fszontagh 1 mês atrás
pai
commit
5c59a40790

+ 3 - 0
proto/runner.proto

@@ -216,6 +216,9 @@ service RunnerService {
     // Cancel an execution
     rpc CancelExecution(CancelExecutionRequest) returns (CancelExecutionResponse);
 
+    // What this runner is running right now
+    rpc ListActiveExecutions(ListActiveExecutionsRequest) returns (ListActiveExecutionsResponse);
+
     // Continue an execution that paused for an answer
     rpc ResumeExecution(ResumeExecutionRequest) returns (ExecuteWorkflowResponse);
 

+ 11 - 0
proto/workflow.proto

@@ -140,6 +140,17 @@ message CancelExecutionRequest {
     string reason = 2;
 }
 
+// Which executions a runner is running right now
+message ListActiveExecutionsRequest {
+}
+
+message ListActiveExecutionsResponse {
+    // Execution ids this runner currently has in flight. A runner keeps them in
+    // memory, so a restart empties this - which is exactly how an execution
+    // recorded as running is found to have nobody left to finish it.
+    repeated string execution_ids = 1;
+}
+
 // Cancel execution response
 message CancelExecutionResponse {
     // False when this runner has never heard of the execution - it finished

+ 10 - 0
src/runner/runner_service.cpp

@@ -201,6 +201,16 @@ grpc::Status RunnerServiceImpl::CancelExecution(grpc::ServerContext* context,
     return grpc::Status::OK;
 }
 
+grpc::Status RunnerServiceImpl::ListActiveExecutions(
+        grpc::ServerContext* context,
+        const proto::ListActiveExecutionsRequest* request,
+        proto::ListActiveExecutionsResponse* response) {
+    for (const auto& id : engine_.activeExecutionIds()) {
+        response->add_execution_ids(id);
+    }
+    return grpc::Status::OK;
+}
+
 grpc::Status RunnerServiceImpl::ResumeExecution(grpc::ServerContext* context,
                                                 const proto::ResumeExecutionRequest* request,
                                                 proto::ExecuteWorkflowResponse* response) {

+ 4 - 0
src/runner/runner_service.hpp

@@ -38,6 +38,10 @@ public:
                                  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;

+ 10 - 0
src/runner/workflow_engine.cpp

@@ -1144,6 +1144,16 @@ WorkflowEngine::CancelOutcome WorkflowEngine::cancelExecution(const std::string&
     return outcome;
 }
 
+std::vector<std::string> WorkflowEngine::activeExecutionIds() const {
+    std::lock_guard<std::mutex> lock(mutex_);
+    std::vector<std::string> ids;
+    ids.reserve(active_executions_.size());
+    for (const auto& [id, _] : active_executions_) {
+        ids.push_back(id);
+    }
+    return ids;
+}
+
 int WorkflowEngine::getActiveExecutionCount() const {
     return active_count_.load();
 }

+ 5 - 0
src/runner/workflow_engine.hpp

@@ -213,6 +213,11 @@ public:
     };
     CancelOutcome cancelExecution(const std::string& execution_id);
 
+    // The executions this engine is running. Held in memory, so it is empty
+    // after a restart - which is what lets a caller tell an execution that is
+    // still going from one nobody is going to finish.
+    std::vector<std::string> activeExecutionIds() const;
+
     // Get active execution count
     int getActiveExecutionCount() const;
 

+ 8 - 2
src/webserver/api/runner_controller.cpp

@@ -4,8 +4,9 @@
 namespace smartbotic::webserver::api {
 
 RunnerController::RunnerController(runners::RunnerRegistry& registry,
-                                   auth::AuthMiddleware& middleware)
-    : registry_(registry), middleware_(middleware) {}
+                                   auth::AuthMiddleware& middleware,
+                                   RegisteredHandler on_registered)
+    : registry_(registry), middleware_(middleware), on_registered_(std::move(on_registered)) {}
 
 void RunnerController::registerRoutes(httplib::Server& server) {
     server.Get("/api/v1/runners", [this](const httplib::Request& req, httplib::Response& res) {
@@ -99,6 +100,11 @@ void RunnerController::registerRunner(const httplib::Request& req, httplib::Resp
         }
 
         LOG_INFO("Runner registered via HTTP: {} at {}", runner.id, runner.address);
+
+        if (on_registered_) {
+            on_registered_(runner.id, runner.address);
+        }
+
         sendJson(res, {{"success", true}, {"id", runner.id}});
     } catch (const std::exception& e) {
         sendError(res, "Invalid request body", 400);

+ 12 - 1
src/webserver/api/runner_controller.hpp

@@ -2,6 +2,7 @@
 
 #include <httplib.h>
 #include <nlohmann/json.hpp>
+#include <functional>
 #include "../auth/auth_middleware.hpp"
 #include "../runners/runner_registry.hpp"
 
@@ -9,7 +10,16 @@ namespace smartbotic::webserver::api {
 
 class RunnerController {
 public:
-    RunnerController(runners::RunnerRegistry& registry, auth::AuthMiddleware& middleware);
+    // Called after a runner registers. A runner keeps what it is executing in
+    // memory, so one that has just come up cannot be running anything recorded
+    // against it earlier - and those records would otherwise say "running" for
+    // ever. Kept as a callback to avoid the controller depending on the whole
+    // service.
+    using RegisteredHandler = std::function<void(const std::string& runner_id,
+                                                 const std::string& address)>;
+
+    RunnerController(runners::RunnerRegistry& registry, auth::AuthMiddleware& middleware,
+                     RegisteredHandler on_registered = nullptr);
 
     void registerRoutes(httplib::Server& server);
 
@@ -29,6 +39,7 @@ private:
 
     runners::RunnerRegistry& registry_;
     auth::AuthMiddleware& middleware_;
+    RegisteredHandler on_registered_;
 };
 
 } // namespace smartbotic::webserver::api

+ 75 - 1
src/webserver/webserver_service.cpp

@@ -17,6 +17,7 @@
 #include "credentials/credential_store.hpp"
 #include "scheduler/workflow_scheduler.hpp"
 #include "common/time_utils.hpp"
+#include <unordered_set>
 #include "common/config_defaults.hpp"
 #include "logging/logger.hpp"
 #include <grpcpp/grpcpp.h>
@@ -233,7 +234,11 @@ void WebServerService::setupRoutes() {
         *node_store_, *auth_middleware_, *load_balancer_, &node_sync_server_->service());
     node_ctrl_->registerRoutes(server);
 
-    runner_ctrl_ = std::make_unique<api::RunnerController>(*runner_registry_, *auth_middleware_);
+    runner_ctrl_ = std::make_unique<api::RunnerController>(
+        *runner_registry_, *auth_middleware_,
+        [this](const std::string& runner_id, const std::string& address) {
+            reconcileOrphanedExecutions(runner_id, address);
+        });
     runner_ctrl_->registerRoutes(server);
 
     webhook_ctrl_ = std::make_unique<api::WebhookController>(
@@ -412,6 +417,75 @@ void WebServerService::noteExecutionOutcome(const std::string& workflow_id,
     });
 }
 
+void WebServerService::reconcileOrphanedExecutions(const std::string& runner_id,
+                                                   const std::string& address) {
+    // Ask the runner what it is actually running rather than assuming. A runner
+    // that re-registers while working - a duplicate call, a flapping network -
+    // must not have its live executions closed underneath it, and only the
+    // runner knows which those are.
+    std::unordered_set<std::string> still_running;
+    {
+        auto channel = ::grpc::CreateChannel(address, ::grpc::InsecureChannelCredentials());
+        auto stub = proto::RunnerService::NewStub(channel);
+
+        proto::ListActiveExecutionsRequest request;
+        proto::ListActiveExecutionsResponse response;
+        ::grpc::ClientContext context;
+        context.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(10));
+
+        auto status = stub->ListActiveExecutions(&context, request, &response);
+        if (!status.ok()) {
+            // Without an answer there is no way to tell an orphan from a live
+            // run, and closing a live one is far worse than leaving a stale
+            // record for someone to press Stop on.
+            LOG_WARN("Runner {} could not say what it is running ({}), so nothing was reconciled",
+                     runner_id, status.error_message());
+            return;
+        }
+        for (const auto& id : response.execution_ids()) {
+            still_running.insert(id);
+        }
+    }
+
+    storage::QueryOptions options;
+    options.filters.push_back({"runnerId", runner_id});
+    options.filters.push_back({"status", "running"});
+    options.page = 1;
+    options.page_size = 500;
+
+    auto found = storage_->query("executions", options);
+    if (found.failed()) {
+        LOG_WARN("Could not look for orphaned executions on runner {}: {}",
+                 runner_id, found.error().message());
+        return;
+    }
+
+    int closed = 0;
+    for (const auto& record : found.value().documents) {
+        const std::string id = record.value("_id", "");
+        if (id.empty() || still_running.contains(id)) {
+            continue;
+        }
+
+        // Deliberately not touching Waiting. That is a run parked on a person,
+        // not on a runner, and it is meant to outlive one.
+        const nlohmann::json patch = {
+            {"status", "cancelled"},
+            {"error", "The runner restarted while this was running, so nothing was left to finish it"},
+            {"finishedAt", common::TimeUtils::nowMs()}
+        };
+        if (storage_->update("executions", id, patch, 0, true).ok()) {
+            ++closed;
+            ws_server_->broadcast("executions." + id + ".cancelled", {{"executionId", id}});
+        }
+    }
+
+    if (closed > 0) {
+        LOG_WARN("Closed {} execution(s) that runner {} was recorded as running but is not",
+                 closed, runner_id);
+    }
+}
+
 void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
                                         const std::string& failed_execution_id,
                                         const std::string& error_message) {

+ 7 - 0
src/webserver/webserver_service.hpp

@@ -102,6 +102,13 @@ private:
                               bool failed,
                               const std::string& error);
 
+    // Close executions a runner cannot possibly still be running. A runner
+    // holds what it is executing in memory, so anything recorded as running
+    // against one that has just come up has nobody left to finish it - and
+    // would sit at "running" until somebody noticed and pressed Stop.
+    void reconcileOrphanedExecutions(const std::string& runner_id,
+                                     const std::string& address);
+
     void runErrorWorkflow(const std::string& failed_workflow_id,
                           const std::string& failed_execution_id,
                           const std::string& error_message);