فهرست منبع

fix: Stop closes an execution whose runner has restarted

Two executions sat at "running" for fourteen hours with Stop doing nothing.
Neither was running: a runner keeps what it is executing in memory, so a
restart loses it, and the record is left saying "running" with nobody left to
finish it.

Stop looked like it worked. The cancel reached runner-1, which had restarted
and had never heard of that execution id - so it filed the id away in its
cancellation set, answered OK, and the webserver reported "cancelling, the
runner stops at the next node boundary". Every part of that was true except
that nothing could ever act on it.

CancelExecution now answers whether the runner was actually running it,
rather than returning Empty. When it was not, the webserver closes the record
itself and says why. The path for a runner that is not registered at all was
already handled; this is the same situation with the runner still present
under the same name.

Both stuck executions cleared, with the reason recorded on each rather than a
bare "cancelled" that says nothing about why a fourteen hour run ended.

This does not stop it happening - an execution orphaned by a restart still
reads as running until somebody presses Stop. Reconciling those at startup is
worth doing separately.

63 passed.
fszontagh 1 ماه پیش
والد
کامیت
8591e3462b

+ 1 - 1
proto/runner.proto

@@ -214,7 +214,7 @@ service RunnerService {
     rpc ExecuteWorkflow(ExecuteWorkflowRequest) returns (ExecuteWorkflowResponse);
 
     // Cancel an execution
-    rpc CancelExecution(CancelExecutionRequest) returns (Empty);
+    rpc CancelExecution(CancelExecutionRequest) returns (CancelExecutionResponse);
 
     // 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;
 }
 
+// Cancel execution response
+message CancelExecutionResponse {
+    // False when this runner has never heard of the execution - it finished
+    // before the request arrived, or the runner has restarted since it started.
+    // The caller needs to know, because a cancel nobody can act on leaves the
+    // record saying "running" for ever.
+    bool was_running = 1;
+    // How many executions were marked, counting sub-workflows this one started.
+    int32 marked = 2;
+}
+
 // Retry execution request
 message RetryExecutionRequest {
     string execution_id = 1;

+ 4 - 2
src/runner/runner_service.cpp

@@ -194,8 +194,10 @@ grpc::Status RunnerServiceImpl::ExecuteWorkflow(grpc::ServerContext* context,
 
 grpc::Status RunnerServiceImpl::CancelExecution(grpc::ServerContext* context,
                                                 const proto::CancelExecutionRequest* request,
-                                                proto::Empty* response) {
-    engine_.cancelExecution(request->execution_id());
+                                                proto::CancelExecutionResponse* response) {
+    const auto outcome = engine_.cancelExecution(request->execution_id());
+    response->set_was_running(outcome.was_running);
+    response->set_marked(outcome.marked);
     return grpc::Status::OK;
 }
 

+ 1 - 1
src/runner/runner_service.hpp

@@ -36,7 +36,7 @@ public:
 
     grpc::Status CancelExecution(grpc::ServerContext* context,
                                  const proto::CancelExecutionRequest* request,
-                                 proto::Empty* response) override;
+                                 proto::CancelExecutionResponse* response) override;
 
     grpc::Status ResumeExecution(grpc::ServerContext* context,
                                  const proto::ResumeExecutionRequest* request,

+ 10 - 3
src/runner/workflow_engine.cpp

@@ -1096,9 +1096,12 @@ common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execut
     return execute(workflow, record.value("triggerType", "resume"), trigger_data, callback);
 }
 
-void WorkflowEngine::cancelExecution(const std::string& execution_id) {
+WorkflowEngine::CancelOutcome WorkflowEngine::cancelExecution(const std::string& execution_id) {
     std::lock_guard<std::mutex> lock(mutex_);
 
+    CancelOutcome outcome;
+    outcome.was_running = active_executions_.contains(execution_id);
+
     // Cancelling a run cancels everything it started. Without this, stopping a
     // workflow that calls another one stops nothing anyone can see: the called
     // workflow runs to the end, and the caller only notices between steps, so a
@@ -1133,8 +1136,12 @@ void WorkflowEngine::cancelExecution(const std::string& execution_id) {
         }
     }
 
-    LOG_INFO("Cancellation requested for execution {} ({} execution(s) marked, {} interrupted)",
-             execution_id, cancelled, interrupted);
+    LOG_INFO("Cancellation requested for execution {} ({} execution(s) marked, {} interrupted, "
+             "running here: {})",
+             execution_id, cancelled, interrupted, outcome.was_running);
+
+    outcome.marked = static_cast<int>(cancelled);
+    return outcome;
 }
 
 int WorkflowEngine::getActiveExecutionCount() const {

+ 10 - 1
src/runner/workflow_engine.hpp

@@ -202,7 +202,16 @@ public:
                                            ExecutionCallback callback = nullptr);
 
     // Cancel execution
-    void cancelExecution(const std::string& execution_id);
+    // Returns how many executions were marked - this one plus anything it
+    // started - and whether this runner was actually running it. Zero means the
+    // request reached a runner that has never heard of it, which is what
+    // happens after a restart, and the caller has to resolve the record itself
+    // or it says "running" for ever.
+    struct CancelOutcome {
+        bool was_running = false;
+        int marked = 0;
+    };
+    CancelOutcome cancelExecution(const std::string& execution_id);
 
     // Get active execution count
     int getActiveExecutionCount() const;

+ 27 - 1
src/webserver/api/execution_controller.cpp

@@ -1,6 +1,7 @@
 #include "execution_controller.hpp"
 #include "../webserver_service.hpp"
 #include "logging/logger.hpp"
+#include "common/time_utils.hpp"
 #include "proto/runner.grpc.pb.h"
 #include <grpcpp/grpcpp.h>
 #include <algorithm>
@@ -296,7 +297,7 @@ void ExecutionController::cancelExecution(const httplib::Request& req, httplib::
     proto::CancelExecutionRequest grpc_req;
     grpc_req.set_execution_id(id);
 
-    proto::Empty grpc_res;
+    proto::CancelExecutionResponse grpc_res;
     ::grpc::ClientContext grpc_ctx;
     grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(15));
 
@@ -306,6 +307,31 @@ void ExecutionController::cancelExecution(const httplib::Request& req, httplib::
         return;
     }
 
+    // The runner answered, but it has never heard of this execution. A runner
+    // keeps what it is running in memory, so a restart loses it - and the record
+    // is left saying "running" with nobody left to finish it. Reporting
+    // "cancelling" here is how one sat at 14 hours: the button worked, the
+    // request arrived, and nothing could ever act on it.
+    if (!grpc_res.was_running()) {
+        const nlohmann::json patch = {
+            {"status", "cancelled"},
+            {"error", "The runner was no longer running this - it restarted while the "
+                      "execution was in flight"},
+            {"finishedAt", common::TimeUtils::nowMs()}
+        };
+        auto updated = storage_.update("executions", id, patch, 0, true);
+        if (updated.failed()) {
+            sendError(res, "Could not cancel: " + updated.error().message(), 500);
+            return;
+        }
+        LOG_WARN("Execution {} was not running on {}; marked cancelled in the database", id, runner_id);
+        ws_server_.broadcast("executions." + id + ".cancelled", {{"executionId", id}});
+        sendJson(res, {{"success", true}, {"status", "cancelled"},
+                       {"note", "The runner is no longer running this - it restarted while the "
+                                "execution was in flight - so the record was closed"}});
+        return;
+    }
+
     // Deliberately not writing the status here. Cancellation is honoured at the
     // next node boundary, and the runner writes the finished record itself - a
     // status written now would be overwritten moments later, and would claim