Jelajahi Sumber

feat: resume a paused execution from its stored results and snapshot

fszontagh 1 bulan lalu
induk
melakukan
9bc6219f6b
2 mengubah file dengan 118 tambahan dan 0 penghapusan
  1. 107 0
      src/runner/workflow_engine.cpp
  2. 11 0
      src/runner/workflow_engine.hpp

+ 107 - 0
src/runner/workflow_engine.cpp

@@ -227,6 +227,20 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
         LOG_INFO("Received cached outputs for {} nodes", cached_outputs.size());
     }
 
+    // A resume carries the results of everything that ran before the pause, and
+    // reuses the original execution id so the run reads as one execution rather
+    // than two.
+    nlohmann::json resume_seed;
+    std::string resume_execution_id;
+    if (actual_trigger_data.contains("_resumeSeed")) {
+        resume_seed = actual_trigger_data["_resumeSeed"];
+        actual_trigger_data.erase("_resumeSeed");
+    }
+    if (actual_trigger_data.contains("_resumeExecutionId")) {
+        resume_execution_id = actual_trigger_data["_resumeExecutionId"].get<std::string>();
+        actual_trigger_data.erase("_resumeExecutionId");
+    }
+
     // Create execution record
     ExecutionResult result;
     result.execution_id = UUID::generatePrefixed("exec");
@@ -238,6 +252,20 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
     result.trigger_data = actual_trigger_data;
     result.started_at = TimeUtils::nowMs();
 
+    if (!resume_execution_id.empty()) {
+        result.execution_id = resume_execution_id;
+    }
+    for (auto it = resume_seed.begin(); it != resume_seed.end(); ++it) {
+        NodeExecutionResult seeded;
+        seeded.node_id = it.key();
+        seeded.status = NodeStatus::Completed;
+        seeded.output = it.value();
+        seeded.started_at = TimeUtils::nowMs();
+        seeded.finished_at = seeded.started_at;
+        seeded.from_resume = true;
+        result.node_results[it.key()] = seeded;
+    }
+
     // Create workflow snapshot for pinning feature
     nlohmann::json snapshot;
     snapshot["nodes"] = nlohmann::json::array();
@@ -509,6 +537,16 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
                 input = collectNodeInput(node_id, workflow, result.node_results);
             }
 
+            // On a resume the results of everything that ran before the pause
+            // are already seeded, and they are not run a second time.
+            {
+                auto seeded = result.node_results.find(node_id);
+                if (seeded != result.node_results.end() && seeded->second.from_resume) {
+                    LOG_DEBUG("Node {} already ran before the pause, keeping its result", node_id);
+                    continue;
+                }
+            }
+
             // Check if node should be skipped (inactive branch)
             if (input.contains("_skip") && input["_skip"].get<bool>()) {
                 NodeExecutionResult skip_result;
@@ -765,6 +803,75 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
     return result;
 }
 
+common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execution_id,
+                                                       const std::string& token,
+                                                       const nlohmann::json& payload,
+                                                       ExecutionCallback callback) {
+    auto stored = storage_.get("executions", execution_id);
+    if (stored.failed()) {
+        return common::Error(common::ErrorCode::NotFound, "No execution " + execution_id);
+    }
+
+    const nlohmann::json& record = stored.value();
+
+    if (record.value("status", "") != "waiting") {
+        return common::Error(common::ErrorCode::InvalidArgument,
+            "Execution " + execution_id + " is " + record.value("status", "unknown") +
+            ", not waiting for an answer");
+    }
+
+    const std::string expected_token = record.value("pauseToken", "");
+    if (expected_token.empty() || expected_token != token) {
+        return common::Error(common::ErrorCode::PermissionDenied,
+            "The token does not match the one this execution is waiting on");
+    }
+
+    const int64_t expires_at = record.value("pauseExpiresAt", static_cast<int64_t>(0));
+    if (expires_at > 0 && TimeUtils::nowMs() > expires_at) {
+        return common::Error(common::ErrorCode::InvalidArgument,
+            "This approval expired at " + std::to_string(expires_at) + " and can no longer be answered");
+    }
+
+    const std::string paused_node_id = record.value("pausedNodeId", "");
+    if (paused_node_id.empty()) {
+        return common::Error(common::ErrorCode::Internal,
+            "Execution " + execution_id + " is waiting but records no paused node");
+    }
+
+    if (!record.contains("workflowSnapshot")) {
+        return common::Error(common::ErrorCode::Internal,
+            "Execution " + execution_id + " has no workflow snapshot to resume against");
+    }
+
+    Workflow workflow = Workflow::fromJson(record["workflowSnapshot"]);
+
+    // Seed everything that ran before the pause, including the paused node
+    // itself, whose output becomes the answer that was given.
+    nlohmann::json seed = nlohmann::json::object();
+    if (record.contains("nodeExecutions") && record["nodeExecutions"].is_array()) {
+        for (const auto& entry : record["nodeExecutions"]) {
+            const std::string node_id = entry.value("nodeId", "");
+            if (node_id.empty()) {
+                continue;
+            }
+            seed[node_id] = entry.value("output", nlohmann::json::object());
+        }
+    }
+
+    nlohmann::json answer = payload.is_object() ? payload : nlohmann::json::object();
+    answer["answeredAt"] = TimeUtils::nowMs();
+    seed[paused_node_id] = answer;
+
+    nlohmann::json trigger_data = record.value("triggerData", nlohmann::json::object());
+    trigger_data["_resumeSeed"] = seed;
+    trigger_data["_resumeExecutionId"] = execution_id;
+
+    LOG_INFO("Resuming execution {} from node {} with {} seeded results",
+             execution_id, paused_node_id, seed.size());
+
+    return execute(workflow, record.value("triggerType", "resume"), trigger_data, callback);
+}
+
 void WorkflowEngine::cancelExecution(const std::string& execution_id) {
     std::lock_guard<std::mutex> lock(mutex_);
     cancelled_executions_.insert(execution_id);

+ 11 - 0
src/runner/workflow_engine.hpp

@@ -73,6 +73,7 @@ struct NodeExecutionResult {
     int64_t finished_at = 0;
     int retry_count = 0;
     bool from_cache = false;  // True if output was from cached previous execution
+    bool from_resume = false;   // Seeded from a stored execution rather than run now
 };
 
 // Execution status
@@ -166,6 +167,16 @@ public:
                                             const nlohmann::json& trigger_data,
                                             ExecutionCallback callback = nullptr);
 
+    // Continue an execution that stopped at a pause marker. The workflow is
+    // rebuilt from the snapshot stored with the execution rather than from the
+    // workflow as it stands now, because it may have been edited while the
+    // approval was waiting and finishing a run against a different workflow
+    // than it started under is worse than refusing.
+    common::Result<ExecutionResult> resume(const std::string& execution_id,
+                                           const std::string& token,
+                                           const nlohmann::json& payload,
+                                           ExecutionCallback callback = nullptr);
+
     // Cancel execution
     void cancelExecution(const std::string& execution_id);