Browse Source

fix: a webhook run counts towards the schedule's overlap policy

A workflow on a schedule that is also reachable by webhook could get a
scheduled run on top of a webhook run, however firmly its overlap policy
said skip. Webhook dispatch blocks until the run has finished, so by the
time it came back with an execution id there was nothing left to claim a
slot for, and the tick in between saw the workflow as idle.

The webhook now names the execution id up front and claims the slot before
it dispatches, using the same handover the sub-workflow call already uses -
so the marker is now _assignedExecutionId rather than _childExecutionId.
Measured: firing a webhook 12s before a tick on a 25s run, the tick logged
"still running, skipping this tick" instead of starting a second run.

The slot is given back if the runner never took the work. A gRPC deadline
is the exception - the call gave up, the workflow did not, and it will
report its own end. And if the runner ever stops honouring the offered id,
the slot moves to the id that came back rather than being stranded until
the watchdog.

Scheduler status reported "running" from a bool that was only ever assigned
false, so it said not-running throughout every run. It now reports
activeRuns, the count an overlap policy actually acts on, which is also how
the behaviour above was verified.

No fixture: the shortest honest test spans a scheduler tick, which is a
minute of wall time. Verified by hand against the running services. 60/60.
fszontagh 1 tháng trước cách đây
mục cha
commit
87fc0f3390

+ 10 - 8
src/runner/workflow_engine.cpp

@@ -283,12 +283,14 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
         actual_trigger_data.erase("_callDepth");
     }
 
-    // An id chosen by the caller, so it could record the child before the child
-    // existed. Read here and taken as this run's id.
+    // An id chosen by whoever started this run, so they could record it before
+    // the run existed - a caller tracking a sub-workflow it may need to cancel,
+    // or a webhook claiming its schedule slot before it blocks. Read here and
+    // taken as this run's id.
     std::string assigned_execution_id;
-    if (actual_trigger_data.contains("_childExecutionId")) {
-        assigned_execution_id = actual_trigger_data.value("_childExecutionId", "");
-        actual_trigger_data.erase("_childExecutionId");
+    if (actual_trigger_data.contains("_assignedExecutionId")) {
+        assigned_execution_id = actual_trigger_data.value("_assignedExecutionId", "");
+        actual_trigger_data.erase("_assignedExecutionId");
     }
 
     if (actual_trigger_data.contains("_resumeExecutionId")) {
@@ -1762,12 +1764,12 @@ bool WorkflowEngine::runSubWorkflow(NodeExecutionResult& node_result, int call_d
     // No callback: the sub-workflow reports its own progress against its own
     // execution, and forwarding its node events to the parent's subscribers
     // would make the parent's canvas light up nodes it does not have.
-    // The child's id is generated inside execute(), so the caller cannot know it
-    // in advance. It is handed in on the trigger data instead, and recorded
+    // The id is generated inside execute(), so a caller cannot know it in
+    // advance. It is handed in on the trigger data instead, and recorded
     // against the caller as soon as it is known - which is what lets a cancel
     // arriving mid-call reach the workflow actually doing the work.
     const std::string child_execution_id = UUID::generatePrefixed("exec");
-    sub_trigger["_childExecutionId"] = child_execution_id;
+    sub_trigger["_assignedExecutionId"] = child_execution_id;
     {
         std::lock_guard<std::mutex> lock(mutex_);
         child_executions_[parent_execution_id].insert(child_execution_id);

+ 34 - 2
src/webserver/api/webhook_controller.cpp

@@ -1,6 +1,7 @@
 #include "webhook_controller.hpp"
 #include "logging/logger.hpp"
 #include "common/time_utils.hpp"
+#include "common/uuid.hpp"
 #include "proto/runner.grpc.pb.h"
 #include <grpcpp/grpcpp.h>
 #include <cctype>
@@ -14,12 +15,14 @@ WebhookController::WebhookController(storage::StorageClient& storage,
                                      runners::RunnerRegistry& registry,
                                      runners::LoadBalancer& load_balancer,
                                      WebSocketServer& ws_server,
-                                     nodes::NodeStore& node_store)
+                                     nodes::NodeStore& node_store,
+                                     WorkflowScheduler& scheduler)
     : storage_(storage)
     , registry_(registry)
     , load_balancer_(load_balancer)
     , ws_server_(ws_server)
-    , node_store_(node_store) {}
+    , node_store_(node_store)
+    , scheduler_(scheduler) {}
 
 void WebhookController::registerRoutes(httplib::Server& server) {
     // Match any HTTP method for webhooks
@@ -144,6 +147,16 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
     auto channel = grpc::CreateChannel(runner->address, grpc::InsecureChannelCredentials());
     auto stub = proto::RunnerService::NewStub(channel);
 
+    // A schedule's overlap policy counts what is running, and this call blocks
+    // until the run has finished - so by the time it returns with an execution
+    // id there is nothing left to claim a slot for, and a scheduled tick in the
+    // meantime saw the workflow as idle. Naming the id here rather than letting
+    // the runner pick one is what makes the slot claimable up front; the runner
+    // takes this as the run's id.
+    const std::string execution_id = UUID::generatePrefixed("exec");
+    trigger_data["_assignedExecutionId"] = execution_id;
+    scheduler_.notifyExecutionStarted(workflow_id, execution_id);
+
     proto::ExecuteWorkflowRequest grpc_req;
     grpc_req.set_workflow_id(workflow_id);
     grpc_req.set_trigger_type("webhook");
@@ -158,11 +171,30 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
     auto status = stub->ExecuteWorkflow(&grpc_ctx, grpc_req, &grpc_res);
 
     if (!status.ok()) {
+        // No run means no slot. Left held, this would block every scheduled
+        // tick until the watchdog deadline - up to an hour - because a runner
+        // was briefly unreachable. A deadline is the exception: the call gave
+        // up, the workflow did not, and it will report its own end - releasing
+        // the slot here would say idle while it is still working.
+        if (status.error_code() != grpc::StatusCode::DEADLINE_EXCEEDED) {
+            scheduler_.notifyExecutionFinished(workflow_id, execution_id);
+        }
         LOG_ERROR("Webhook execution failed: {}", status.error_message());
         sendError(res, "Execution failed: " + status.error_message(), 500);
         return;
     }
 
+    // The runner is expected to take the id offered above. If it ever stops
+    // doing so, the slot claimed under our id has nobody to release it and the
+    // schedule would stall until the watchdog deadline, so move it rather than
+    // leave a stuck workflow behind a silent assumption.
+    if (grpc_res.execution_id() != execution_id) {
+        LOG_WARN("Webhook run came back as {} but was dispatched as {} - moving the schedule slot",
+                 grpc_res.execution_id(), execution_id);
+        scheduler_.notifyExecutionFinished(workflow_id, execution_id);
+        scheduler_.notifyExecutionStarted(workflow_id, grpc_res.execution_id());
+    }
+
     // A workflow that paused mid-run is not done: final_output is still null
     // (the body would just be the literal string "null"), and labelling this
     // ".completed" is false. Report it honestly - 202 Accepted, a status of

+ 6 - 1
src/webserver/api/webhook_controller.hpp

@@ -6,6 +6,7 @@
 #include "../runners/load_balancer.hpp"
 #include "../websocket_server.hpp"
 #include "../nodes/node_store.hpp"
+#include "../scheduler/workflow_scheduler.hpp"
 
 namespace smartbotic::webserver::api {
 
@@ -15,7 +16,8 @@ public:
                      runners::RunnerRegistry& registry,
                      runners::LoadBalancer& load_balancer,
                      WebSocketServer& ws_server,
-                     nodes::NodeStore& node_store);
+                     nodes::NodeStore& node_store,
+                     WorkflowScheduler& scheduler);
 
     void registerRoutes(httplib::Server& server);
 
@@ -35,6 +37,9 @@ private:
     runners::LoadBalancer& load_balancer_;
     WebSocketServer& ws_server_;
     nodes::NodeStore& node_store_;
+    // Held so a webhook run counts towards a schedule's overlap policy, the
+    // same as any other way of starting the workflow.
+    WorkflowScheduler& scheduler_;
 };
 
 } // namespace smartbotic::webserver::api

+ 5 - 3
src/webserver/scheduler/workflow_scheduler.cpp

@@ -77,7 +77,6 @@ void WorkflowScheduler::registerWorkflow(
     entry.interval_minutes = interval_minutes;
     entry.last_run = now;  // Consider it just ran to avoid immediate execution
     entry.next_run = now + std::chrono::minutes(interval_minutes);
-    entry.running = false;
     entry.overlap_policy = overlap_policy;
     entry.max_concurrent = max_concurrent > 0 ? max_concurrent : 1;
     entry.max_run_minutes = max_run_minutes;
@@ -132,7 +131,11 @@ nlohmann::json WorkflowScheduler::getScheduledWorkflows() const {
             {"triggerNodeId", entry.trigger_node_id},
             {"triggerType", entry.trigger_type},
             {"intervalMinutes", entry.interval_minutes},
-            {"running", entry.running},
+            // What is running right now, which is what an overlap policy acts
+            // on. This used to report a bool that was only ever assigned false,
+            // so the status said "not running" throughout every run.
+            {"activeRuns", entry.active_runs.size()},
+            {"running", !entry.active_runs.empty()},
             {"secondsUntilNextRun", seconds_until_next > 0 ? seconds_until_next : 0}
         });
     }
@@ -345,7 +348,6 @@ void WorkflowScheduler::checkAndExecute() {
                 // this pass: otherwise a slow dispatch leaves next_run in the past
                 // and the following tick fires immediately.
                 auto dispatched_at = std::chrono::steady_clock::now();
-                it->second.running = false;
                 it->second.last_run = dispatched_at;
                 it->second.next_run = dispatched_at + std::chrono::minutes(it->second.interval_minutes);
 

+ 0 - 1
src/webserver/scheduler/workflow_scheduler.hpp

@@ -33,7 +33,6 @@ struct ScheduledWorkflow {
     int interval_minutes;
     std::chrono::steady_clock::time_point last_run;
     std::chrono::steady_clock::time_point next_run;
-    bool running = false;
 
     OverlapPolicy overlap_policy = OverlapPolicy::Skip;
     int max_concurrent = 1;

+ 1 - 1
src/webserver/webserver_service.cpp

@@ -237,7 +237,7 @@ void WebServerService::setupRoutes() {
     runner_ctrl_->registerRoutes(server);
 
     webhook_ctrl_ = std::make_unique<api::WebhookController>(
-        *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_);
+        *storage_, *runner_registry_, *load_balancer_, *ws_server_, *node_store_, *scheduler_);
     webhook_ctrl_->registerRoutes(server);
 
     database_ctrl_ = std::make_unique<api::DatabaseController>(*storage_, *auth_middleware_);