Pārlūkot izejas kodu

feat(scheduler): overlap control for scheduled workflows

A schedule trigger can now declare what happens when its time arrives while
a previous run is still in flight: skip the tick, defer exactly one start
until the run finishes, or allow up to maxConcurrent runs.

The existing "skip if already running" guard did not do this. Scheduled runs
dispatch with wait_for_completion(false), so the gRPC call returns once the
runner accepts the job, and checkAndExecute cleared the flag immediately
after. The flag covered the dispatch, a few milliseconds, not the run, so two
runs of the same workflow could overlap unnoticed. That matters here because
run time is dominated by hosted vision latency, measured between 12s and 117s
per image, so a slow patch can push a run past its interval.

The scheduler now tracks runs from dispatch until an execution.completed,
failed or cancelled event arrives. Three hazards are handled explicitly:

- Completion events do not always carry a workflow id, so a finished run is
  also located by its execution id.
- A fast workflow can finish before ExecuteWorkflow returns, delivering the
  finished event before the started one. Such ids are remembered briefly so
  the late registration does not create a slot nobody releases.
- A lost event would otherwise block a skip or queue workflow forever, so
  each run carries a deadline after which its slot is released with a warning.

Also computes next_run from the time after dispatch rather than before, so a
long run no longer leaves next_run in the past and re-fire immediately.
fszontagh 1 mēnesi atpakaļ
vecāks
revīzija
a7b9dd7332

+ 25 - 0
nodes/core/schedule-trigger.js

@@ -39,6 +39,31 @@ const configSchema = {
       description: 'Timezone for cron expression (e.g., UTC, Europe/Budapest)',
       default: 'UTC'
     },
+    overlapPolicy: {
+      type: 'string',
+      title: 'When Already Running',
+      description: 'What to do when the scheduled time arrives while a previous run is still in flight',
+      enum: ['skip', 'queue', 'allow'],
+      enumLabels: ['Skip this run', 'Run once after it finishes', 'Allow concurrent runs'],
+      default: 'skip'
+    },
+    maxConcurrent: {
+      type: 'integer',
+      title: 'Max Concurrent Runs',
+      description: 'Concurrent run limit. Only applies when the policy allows concurrent runs.',
+      default: 1,
+      minimum: 1,
+      maximum: 20,
+      showWhen: { field: 'overlapPolicy', value: 'allow' }
+    },
+    maxRunMinutes: {
+      type: 'integer',
+      title: 'Abandon Run After (minutes)',
+      description: 'Release the slot if a run never reports completion, so a lost event cannot block the schedule forever. 0 means five times the interval.',
+      default: 0,
+      minimum: 0,
+      maximum: 1440
+    },
     triggerName: {
       type: 'string',
       title: 'Trigger Name',

+ 12 - 2
src/webserver/api/execution_controller.cpp

@@ -5,8 +5,10 @@ namespace smartbotic::webserver::api {
 
 ExecutionController::ExecutionController(storage::StorageClient& storage,
                                          auth::AuthMiddleware& middleware,
-                                         WebSocketServer& ws_server)
-    : storage_(storage), middleware_(middleware), ws_server_(ws_server) {}
+                                         WebSocketServer& ws_server,
+                                         WorkflowScheduler& scheduler)
+    : storage_(storage), middleware_(middleware), ws_server_(ws_server),
+      scheduler_(scheduler) {}
 
 void ExecutionController::registerRoutes(httplib::Server& server) {
     server.Get("/api/v1/executions", [this](const httplib::Request& req, httplib::Response& res) {
@@ -171,6 +173,14 @@ void ExecutionController::receiveExecutionEvent(const httplib::Request& req, htt
             channel = "executions." + execution_id + "." + event_type;
         }
 
+        // Release the scheduler slot held by this run. Scheduled dispatch is
+        // fire-and-forget, so these events are the only signal that a run ended.
+        if (!execution_id.empty() &&
+            (event_type == "execution.completed" || event_type == "execution.failed" ||
+             event_type == "execution.cancelled")) {
+            scheduler_.notifyExecutionFinished(workflow_id, execution_id);
+        }
+
         // Broadcast to WebSocket clients
         ws_server_.broadcast(channel, data);
 

+ 3 - 1
src/webserver/api/execution_controller.hpp

@@ -5,13 +5,14 @@
 #include "../auth/auth_middleware.hpp"
 #include "../websocket_server.hpp"
 #include "storage/storage_client.hpp"
+#include "../scheduler/workflow_scheduler.hpp"
 
 namespace smartbotic::webserver::api {
 
 class ExecutionController {
 public:
     ExecutionController(storage::StorageClient& storage, auth::AuthMiddleware& middleware,
-                        WebSocketServer& ws_server);
+                        WebSocketServer& ws_server, WorkflowScheduler& scheduler);
 
     void registerRoutes(httplib::Server& server);
 
@@ -32,6 +33,7 @@ private:
     storage::StorageClient& storage_;
     auth::AuthMiddleware& middleware_;
     WebSocketServer& ws_server_;
+    WorkflowScheduler& scheduler_;
 };
 
 } // namespace smartbotic::webserver::api

+ 173 - 10
src/webserver/scheduler/workflow_scheduler.cpp

@@ -3,6 +3,20 @@
 
 namespace smartbotic::webserver {
 
+OverlapPolicy overlapPolicyFromString(const std::string& value) {
+    if (value == "queue") return OverlapPolicy::Queue;
+    if (value == "allow") return OverlapPolicy::Allow;
+    return OverlapPolicy::Skip;
+}
+
+std::string overlapPolicyToString(OverlapPolicy policy) {
+    switch (policy) {
+        case OverlapPolicy::Queue: return "queue";
+        case OverlapPolicy::Allow: return "allow";
+        default: return "skip";
+    }
+}
+
 WorkflowScheduler::WorkflowScheduler() = default;
 
 WorkflowScheduler::~WorkflowScheduler() {
@@ -41,7 +55,10 @@ void WorkflowScheduler::registerWorkflow(
     const std::string& workflow_name,
     const std::string& trigger_node_id,
     const std::string& trigger_type,
-    int interval_minutes
+    int interval_minutes,
+    OverlapPolicy overlap_policy,
+    int max_concurrent,
+    int max_run_minutes
 ) {
     if (interval_minutes <= 0) {
         spdlog::debug("Workflow {} has interval 0, not scheduling", workflow_id);
@@ -61,11 +78,25 @@ void WorkflowScheduler::registerWorkflow(
     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;
+
+    // Preserve in-flight state across a re-registration, which happens whenever
+    // the workflow is saved while active. Dropping it would leak a slot.
+    auto existing = workflows_.find(workflow_id);
+    if (existing != workflows_.end()) {
+        entry.active_runs = existing->second.active_runs;
+        entry.pending_start = existing->second.pending_start;
+        entry.next_run = existing->second.next_run;
+    }
 
     workflows_[workflow_id] = entry;
 
-    spdlog::info("Scheduled workflow '{}' ({}) with {} trigger, interval: {} minutes",
-                 workflow_name, workflow_id, trigger_type, interval_minutes);
+    spdlog::info("Scheduled workflow '{}' ({}) with {} trigger, interval: {} minutes, "
+                 "overlap: {}, maxConcurrent: {}",
+                 workflow_name, workflow_id, trigger_type, interval_minutes,
+                 overlapPolicyToString(entry.overlap_policy), entry.max_concurrent);
 }
 
 void WorkflowScheduler::unregisterWorkflow(const std::string& workflow_id) {
@@ -140,6 +171,82 @@ void WorkflowScheduler::schedulerLoop() {
     spdlog::info("Scheduler loop ended");
 }
 
+void WorkflowScheduler::notifyExecutionStarted(const std::string& workflow_id,
+                                               const std::string& execution_id) {
+    std::lock_guard<std::mutex> lock(mutex_);
+
+    auto it = workflows_.find(workflow_id);
+    if (it == workflows_.end()) {
+        return;
+    }
+
+    int minutes = it->second.max_run_minutes;
+    if (minutes <= 0) {
+        minutes = it->second.interval_minutes * 5;
+    }
+    if (minutes <= 0) {
+        minutes = 60;
+    }
+
+    // The run may already have reported completion: a fast workflow can finish
+    // before ExecuteWorkflow returns. Consume that and leave the slot free.
+    if (finished_before_start_.erase(execution_id) > 0) {
+        spdlog::debug("Workflow '{}' run {} finished before its dispatch returned",
+                      it->second.workflow_name, execution_id);
+        return;
+    }
+
+    it->second.active_runs[execution_id] =
+        std::chrono::steady_clock::now() + std::chrono::minutes(minutes);
+
+    spdlog::debug("Workflow '{}' run {} started, {} active",
+                  it->second.workflow_name, execution_id, it->second.active_runs.size());
+}
+
+void WorkflowScheduler::notifyExecutionFinished(const std::string& workflow_id,
+                                                const std::string& execution_id) {
+    std::lock_guard<std::mutex> lock(mutex_);
+
+    auto it = workflows_.find(workflow_id);
+
+    // Execution events do not always carry the workflow id, so fall back to
+    // locating the run by its execution id, which the scheduler already tracks.
+    if (it == workflows_.end() || !it->second.active_runs.count(execution_id)) {
+        it = workflows_.end();
+        for (auto candidate = workflows_.begin(); candidate != workflows_.end(); ++candidate) {
+            if (candidate->second.active_runs.count(execution_id)) {
+                it = candidate;
+                break;
+            }
+        }
+    }
+
+    if (it == workflows_.end()) {
+        // Arrived before the dispatch call returned; remember it so the pending
+        // notifyExecutionStarted does not register a slot nobody will release.
+        finished_before_start_.insert(execution_id);
+        if (finished_before_start_.size() > 256) {
+            finished_before_start_.erase(finished_before_start_.begin());
+        }
+        return;
+    }
+
+    if (it->second.active_runs.erase(execution_id) == 0) {
+        finished_before_start_.insert(execution_id);
+        return;
+    }
+
+    spdlog::debug("Workflow '{}' run {} finished, {} active",
+                  it->second.workflow_name, execution_id, it->second.active_runs.size());
+
+    // A deferred start becomes eligible on the next pass, which is at most
+    // CHECK_INTERVAL_SECONDS away.
+    if (it->second.pending_start && it->second.active_runs.empty()) {
+        spdlog::info("Workflow '{}' has a deferred run, starting it now",
+                     it->second.workflow_name);
+    }
+}
+
 void WorkflowScheduler::checkAndExecute() {
     auto now = std::chrono::steady_clock::now();
 
@@ -150,16 +257,68 @@ void WorkflowScheduler::checkAndExecute() {
         std::lock_guard<std::mutex> lock(mutex_);
 
         for (auto& [id, entry] : workflows_) {
-            // Skip if already running
-            if (entry.running) {
+            // Release slots held by runs that never reported completion. Without
+            // this a lost event would block a skip or queue workflow forever.
+            for (auto it = entry.active_runs.begin(); it != entry.active_runs.end(); ) {
+                if (now >= it->second) {
+                    spdlog::warn("Workflow '{}' run {} exceeded its deadline, releasing the slot",
+                                 entry.workflow_name, it->first);
+                    it = entry.active_runs.erase(it);
+                } else {
+                    ++it;
+                }
+            }
+
+            const int active = static_cast<int>(entry.active_runs.size());
+            const bool due = now >= entry.next_run;
+            const bool queued = entry.pending_start && active == 0;
+
+            if (!due && !queued) {
                 continue;
             }
 
-            // Check if it's time to run
-            if (now >= entry.next_run) {
-                entry.running = true;
+            // A deferred start fires as soon as the slot frees, regardless of the
+            // interval boundary.
+            if (queued) {
+                entry.pending_start = false;
                 to_execute.push_back(entry);
+                continue;
             }
+
+            switch (entry.overlap_policy) {
+                case OverlapPolicy::Allow:
+                    if (active >= entry.max_concurrent) {
+                        spdlog::info("Workflow '{}' at concurrency limit {}, skipping tick",
+                                     entry.workflow_name, entry.max_concurrent);
+                        entry.next_run = now + std::chrono::minutes(entry.interval_minutes);
+                        continue;
+                    }
+                    break;
+
+                case OverlapPolicy::Queue:
+                    if (active > 0) {
+                        if (!entry.pending_start) {
+                            spdlog::info("Workflow '{}' still running, deferring this tick",
+                                         entry.workflow_name);
+                            entry.pending_start = true;
+                        }
+                        entry.next_run = now + std::chrono::minutes(entry.interval_minutes);
+                        continue;
+                    }
+                    break;
+
+                case OverlapPolicy::Skip:
+                default:
+                    if (active > 0) {
+                        spdlog::info("Workflow '{}' still running, skipping this tick",
+                                     entry.workflow_name);
+                        entry.next_run = now + std::chrono::minutes(entry.interval_minutes);
+                        continue;
+                    }
+                    break;
+            }
+
+            to_execute.push_back(entry);
         }
     }
 
@@ -182,9 +341,13 @@ void WorkflowScheduler::checkAndExecute() {
             std::lock_guard<std::mutex> lock(mutex_);
             auto it = workflows_.find(entry.workflow_id);
             if (it != workflows_.end()) {
+                // Schedule from the moment dispatch finished, not from the start of
+                // 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 = now;
-                it->second.next_run = now + std::chrono::minutes(it->second.interval_minutes);
+                it->second.last_run = dispatched_at;
+                it->second.next_run = dispatched_at + std::chrono::minutes(it->second.interval_minutes);
 
                 auto next_in_minutes = it->second.interval_minutes;
                 spdlog::debug("Workflow '{}' next run in {} minutes",

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

@@ -2,6 +2,7 @@
 
 #include <string>
 #include <map>
+#include <set>
 #include <mutex>
 #include <thread>
 #include <atomic>
@@ -14,6 +15,16 @@ namespace smartbotic::webserver {
 /**
  * Scheduled workflow entry
  */
+// What a scheduled tick does when a previous run is still in flight.
+enum class OverlapPolicy {
+    Skip,   // drop the tick
+    Queue,  // remember one pending start, fire it when the run finishes
+    Allow   // start while active_runs < max_concurrent
+};
+
+OverlapPolicy overlapPolicyFromString(const std::string& value);
+std::string overlapPolicyToString(OverlapPolicy policy);
+
 struct ScheduledWorkflow {
     std::string workflow_id;
     std::string workflow_name;
@@ -23,6 +34,19 @@ struct ScheduledWorkflow {
     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;
+    int max_run_minutes = 0;  // 0 means derive as 5x the interval
+
+    // Runs dispatched but not yet reported finished, mapped to the point at which
+    // they are treated as abandoned. Dispatch is fire-and-forget, so without a
+    // deadline a lost completion event would block the schedule indefinitely.
+    std::map<std::string, std::chrono::steady_clock::time_point> active_runs;
+
+    // At most one deferred start, so a workflow that consistently overruns its
+    // interval cannot accumulate a backlog it can never drain.
+    bool pending_start = false;
 };
 
 /**
@@ -72,7 +96,24 @@ public:
                           const std::string& workflow_name,
                           const std::string& trigger_node_id,
                           const std::string& trigger_type,
-                          int interval_minutes);
+                          int interval_minutes,
+                          OverlapPolicy overlap_policy = OverlapPolicy::Skip,
+                          int max_concurrent = 1,
+                          int max_run_minutes = 0);
+
+    /**
+     * Record that a dispatched run has started. Dispatch is fire-and-forget, so
+     * the scheduler only learns a run is in flight through this call.
+     */
+    void notifyExecutionStarted(const std::string& workflow_id,
+                                const std::string& execution_id);
+
+    /**
+     * Release the slot held by a run that completed, failed or was cancelled.
+     * A queued start becomes eligible immediately.
+     */
+    void notifyExecutionFinished(const std::string& workflow_id,
+                                 const std::string& execution_id);
 
     /**
      * Unregister a workflow from scheduled execution
@@ -106,6 +147,12 @@ private:
     mutable std::mutex mutex_;
     std::map<std::string, ScheduledWorkflow> workflows_;
 
+    // Runs that reported completion before their dispatch call returned. A fast
+    // workflow can finish before ExecuteWorkflow responds, so the finished event
+    // arrives first; without this the slot would be registered after the fact and
+    // never released.
+    std::set<std::string> finished_before_start_;
+
     std::thread scheduler_thread_;
     std::atomic<bool> running_{false};
 

+ 11 - 2
src/webserver/webserver_service.cpp

@@ -209,7 +209,7 @@ void WebServerService::setupRoutes() {
         *storage_, *auth_middleware_, *ws_server_);
     workflow_group_ctrl_->registerRoutes(server);
 
-    execution_ctrl_ = std::make_unique<api::ExecutionController>(*storage_, *auth_middleware_, *ws_server_);
+    execution_ctrl_ = std::make_unique<api::ExecutionController>(*storage_, *auth_middleware_, *ws_server_, *scheduler_);
     execution_ctrl_->registerRoutes(server);
 
     node_ctrl_ = std::make_unique<api::NodeController>(
@@ -273,12 +273,18 @@ void WebServerService::loadScheduledWorkflows() {
             int interval = config.value("pollInterval", 0);
 
             if (interval > 0) {
+                auto policy = overlapPolicyFromString(
+                    config.value("overlapPolicy", std::string("skip")));
+
                 scheduler_->registerWorkflow(
                     workflow_id,
                     workflow_name,
                     node_id,
                     node_type,
-                    interval
+                    interval,
+                    policy,
+                    config.value("maxConcurrent", 1),
+                    config.value("maxRunMinutes", 0)
                 );
                 registered_count++;
                 LOG_DEBUG("Registered workflow '{}' ({}) for scheduled execution every {} minutes",
@@ -330,6 +336,9 @@ void WebServerService::executeScheduledWorkflow(const std::string& workflow_id,
     LOG_INFO("Scheduled workflow {} execution started: {} on runner {}",
              workflow_id, response.execution_id(), runner->id);
 
+    // Dispatch is fire-and-forget, so the scheduler only learns about the run here.
+    scheduler_->notifyExecutionStarted(workflow_id, response.execution_id());
+
     // Broadcast execution started
     ws_server_->broadcast("executions." + response.execution_id() + ".started", {
         {"executionId", response.execution_id()},