浏览代码

fix: a runner with no nodes says so, and a lost completion no longer stalls a schedule

35photo2anime stopped running. Two separate faults, both of the same
shape: a failure was logged as a warning and normal operation carried on
top of it.

**The runner had no node definitions at all.** It loads them from the
webserver over gRPC at startup and then subscribes to changes. The load
was attempted exactly once - when the two services start together the
webserver's gRPC port is not always up yet, so it failed, the
subscription succeeded, and the runner ran on with an empty registry. A
subscription carries changes, not the current set, so nothing ever filled
it. Every workflow failed with "Node type not found: click-trigger" and
would have until somebody happened to edit a node. It had been that way
for hours.

The sync loop repairs an empty registry now, before subscribing, and
retries until it succeeds - saying each time that every workflow will
fail until it does. An empty registry is a fault to fix, not a state to
run in.

**The scheduler skipped every tick for over an hour.** It holds a slot per
run and frees it when the completion event arrives; a lost event leaves
the slot held until its deadline, which for a workflow that sets no
maximum is five times its interval - so a quarter-hourly schedule stops
for 75 minutes. The database went down mid-run here, which is enough to
lose one. All it said was "still running, skipping this tick".

The execution record knows whether a run finished, so the scheduler asks
it, and a stall becomes one tick rather than 75 minutes. The question is
put through a callback rather than giving the scheduler a storage client,
so the call happens outside its lock: the case being fixed is a database
that is slow or down, and blocking the scheduler's mutex on it would
stall every other workflow too. Unreadable counts as not finished -
guessing the other way would start a second run of a workflow that is
still going. The skip message now says how many runs are in flight.

Verified: a runner pointed at a dead sync address retries every five
seconds and names the consequence, then loads all 85 definitions the
moment the address answers - where before it fell silent and broke every
run. The schedule fired again on time. I could not manufacture a lost
completion event on demand, so the scheduler's recovery path is reasoned
and reviewed rather than reproduced; the stall it addresses is in the
logs. 68 passed, 0 failed, 2 skipped for a model this machine is not
running.
fszontagh 1 月之前
父节点
当前提交
986d98c212

+ 31 - 2
src/runner/node_registry.cpp

@@ -1,4 +1,5 @@
 #include "node_registry.hpp"
+#include <shared_mutex>
 #include "logging/logger.hpp"
 #include "common/time_utils.hpp"
 
@@ -123,10 +124,19 @@ void NodeRegistry::start() {
 
     running_ = true;
 
-    // Load initial nodes
+    // Load initial nodes. A failure here is not fatal and not final: the
+    // subscribe loop below retries it until it succeeds.
+    //
+    // It used to be attempted exactly once. When the runner and the webserver
+    // start together the webserver's gRPC port is not always up yet, so the
+    // load failed, the subscription - which carries changes, not the current
+    // set - succeeded, and the runner ran on with an empty registry. Every
+    // workflow then failed with "Node type not found: click-trigger" until
+    // somebody happened to edit a node. It had been that way for hours here.
     auto result = loadNodesFromWebServer();
     if (result.failed()) {
-        LOG_WARN("Failed to load initial nodes: {}", result.error().message());
+        LOG_WARN("Could not load nodes at startup ({}) - retrying in the background",
+                 result.error().message());
     }
 
     // Start subscription thread for live updates
@@ -186,8 +196,27 @@ Result<void> NodeRegistry::loadNodesFromWebServer() {
     return {};
 }
 
+bool NodeRegistry::isEmpty() const {
+    std::shared_lock<std::shared_mutex> lock(mutex_);
+    return nodes_.empty();
+}
+
 void NodeRegistry::subscribeLoop() {
     while (running_) {
+        // The subscription delivers changes, so it cannot fill an empty
+        // registry. Anything that left it empty - a startup race, a webserver
+        // that was down - is repaired here before subscribing again.
+        if (isEmpty()) {
+            auto reload = loadNodesFromWebServer();
+            if (reload.failed()) {
+                LOG_WARN("Node registry is empty and could not be filled ({}); every workflow "
+                         "will fail until it is - retrying", reload.error().message());
+                std::this_thread::sleep_for(
+                    std::chrono::milliseconds(config_.reconnect_interval_ms));
+                continue;
+            }
+        }
+
         proto::SubscribeToNodeChangesRequest request;
         request.set_runner_id("runner_" + std::to_string(reinterpret_cast<uintptr_t>(this)));
 

+ 4 - 0
src/runner/node_registry.hpp

@@ -65,6 +65,10 @@ public:
 
     // Start/stop node sync with webserver
     void start();
+
+    // Whether anything is loaded. An empty registry fails every workflow, so
+    // the sync loop treats it as a fault to repair rather than a state to keep.
+    [[nodiscard]] bool isEmpty() const;
     void stop();
 
     // Initial load of all nodes from webserver

+ 42 - 2
src/webserver/scheduler/workflow_scheduler.cpp

@@ -253,6 +253,45 @@ void WorkflowScheduler::notifyExecutionFinished(const std::string& workflow_id,
 void WorkflowScheduler::checkAndExecute() {
     auto now = std::chrono::steady_clock::now();
 
+    // A slot is held on an in-memory event, and an event can be lost - the
+    // database going down mid-run is enough. Until this, a lost completion held
+    // the slot until its deadline, which for a workflow that sets no maximum is
+    // five times its interval: a quarter-hourly schedule stopped for 75 minutes
+    // and only said "still running, skipping this tick" while it did.
+    //
+    // The execution record knows the truth, so ask it. Outside the lock,
+    // deliberately: the case being fixed is a database that is slow or down, and
+    // blocking the scheduler's mutex on it would stall every other workflow too.
+    if (execution_finished_check_) {
+        std::vector<std::pair<std::string, std::string>> held;  // workflow id, execution id
+        {
+            std::lock_guard<std::mutex> lock(mutex_);
+            for (const auto& [id, entry] : workflows_) {
+                for (const auto& [execution_id, deadline] : entry.active_runs) {
+                    held.emplace_back(id, execution_id);
+                }
+            }
+        }
+        std::vector<std::pair<std::string, std::string>> finished;
+        for (const auto& [workflow_id, execution_id] : held) {
+            if (execution_finished_check_(execution_id)) {
+                finished.emplace_back(workflow_id, execution_id);
+            }
+        }
+        if (!finished.empty()) {
+            std::lock_guard<std::mutex> lock(mutex_);
+            for (const auto& [workflow_id, execution_id] : finished) {
+                auto it = workflows_.find(workflow_id);
+                if (it == workflows_.end()) continue;
+                if (it->second.active_runs.erase(execution_id) > 0) {
+                    spdlog::info("Workflow '{}' run {} had already finished - releasing the slot "
+                                 "its completion event never freed",
+                                 it->second.workflow_name, execution_id);
+                }
+            }
+        }
+    }
+
     // Collect workflows that need to run
     std::vector<ScheduledWorkflow> to_execute;
 
@@ -313,8 +352,9 @@ void WorkflowScheduler::checkAndExecute() {
                 case OverlapPolicy::Skip:
                 default:
                     if (active > 0) {
-                        spdlog::info("Workflow '{}' still running, skipping this tick",
-                                     entry.workflow_name);
+                        spdlog::info("Workflow '{}' still running ({} in flight), skipping this "
+                                     "tick",
+                                     entry.workflow_name, active);
                         entry.next_run = now + std::chrono::minutes(entry.interval_minutes);
                         continue;
                     }

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

@@ -76,6 +76,10 @@ public:
     /**
      * Start the scheduler thread
      */
+    void setExecutionFinishedCheck(std::function<bool(const std::string&)> check) {
+        execution_finished_check_ = std::move(check);
+    }
+
     void start();
 
     /**
@@ -143,6 +147,12 @@ private:
     void schedulerLoop();
     void checkAndExecute();
 
+    // Asks whether an execution has finished, for slots whose completion event
+    // never arrived. Set by the owner, which has the storage client; the
+    // scheduler deliberately does not, so that a database call cannot be made
+    // while its lock is held.
+    std::function<bool(const std::string& execution_id)> execution_finished_check_;
+
     mutable std::mutex mutex_;
     std::map<std::string, ScheduledWorkflow> workflows_;
 

+ 16 - 0
src/webserver/webserver_service.cpp

@@ -75,6 +75,22 @@ WebServerService::WebServerService(const WebServerServiceConfig& config)
 
     // Initialize workflow scheduler
     scheduler_ = std::make_unique<WorkflowScheduler>();
+    // The scheduler holds a slot per run and frees it when the completion event
+    // arrives. Events get lost - a database outage mid-run is enough - so it can
+    // also ask whether a run has finished. Answering from the execution record
+    // rather than from memory turns a 75-minute stall into one tick.
+    scheduler_->setExecutionFinishedCheck([this](const std::string& execution_id) {
+        auto record = storage_->get("executions", execution_id);
+        if (record.failed()) {
+            // Unreadable is not finished. Saying otherwise here would start a
+            // second run of a workflow that is still going, which is worse than
+            // waiting for the deadline the scheduler already has.
+            return false;
+        }
+        const std::string status = record.value().value("status", "");
+        return status == "completed" || status == "failed" || status == "cancelled";
+    });
+
     scheduler_->setExecuteCallback([this](const std::string& workflow_id,
                                           const std::string& trigger_node_id,
                                           const std::string& trigger_type) {