Kaynağa Gözat

feat: an execution can pause, and a paused one outlives the seven day log TTL

fszontagh 1 ay önce
ebeveyn
işleme
24fbd61c5f
2 değiştirilmiş dosya ile 55 ekleme ve 3 silme
  1. 50 2
      src/runner/workflow_engine.cpp
  2. 5 1
      src/runner/workflow_engine.hpp

+ 50 - 2
src/runner/workflow_engine.cpp

@@ -30,6 +30,7 @@ std::string executionStatusToString(ExecutionStatus status) {
         case ExecutionStatus::Completed: return "completed";
         case ExecutionStatus::Failed: return "failed";
         case ExecutionStatus::Cancelled: return "cancelled";
+        case ExecutionStatus::Waiting: return "waiting";
         default: return "unknown";
     }
 }
@@ -163,6 +164,12 @@ nlohmann::json ExecutionResult::toJson() const {
         j["webhookResponse"] = webhook_response;
     }
 
+    if (!paused_node_id.empty()) {
+        j["pausedNodeId"] = paused_node_id;
+        j["pauseToken"] = pause_token;
+        j["pauseExpiresAt"] = pause_expires_at;
+    }
+
     return j;
 }
 
@@ -616,6 +623,35 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
                 result.webhook_response = node_result.output["_webhookResponse"];
             }
 
+            // A node asking to pause ends this pass. The execution is stored as
+            // Waiting with everything computed so far, and a later resume picks
+            // it up from here rather than starting again.
+            if (node_result.status == NodeStatus::Completed &&
+                node_result.output.contains("_pause")) {
+                const auto& pause = node_result.output["_pause"];
+
+                result.status = ExecutionStatus::Waiting;
+                result.paused_node_id = node_id;
+                result.pause_token = pause.value("token", "");
+                result.pause_expires_at = pause.value("expiresAt", static_cast<int64_t>(0));
+
+                // The marker is engine plumbing. What is stored is the request a
+                // person has to answer, not the mechanism that carried it.
+                auto& stored_node = result.node_results[node_id];
+                stored_node.output.erase("_pause");
+
+                if (callback) {
+                    callback("execution.waiting", {
+                        {"executionId", result.execution_id},
+                        {"nodeId", node_id},
+                        {"expiresAt", result.pause_expires_at}
+                    });
+                }
+
+                LOG_INFO("Execution {} paused at node {}", result.execution_id, node_id);
+                break;
+            }
+
             // Check for loop node
             if (node_result.status == NodeStatus::Completed &&
                 node_result.output.contains("_isLoop") &&
@@ -1304,8 +1340,20 @@ nlohmann::json WorkflowEngine::collectNodeInput(
 }
 
 void WorkflowEngine::storeExecution(const ExecutionResult& result) {
-    auto insert_result = storage_.insert("executions", result.toJson(), result.execution_id,
-                   7 * 24 * 60 * 60 * 1000);  // 7-day TTL
+    // A finished execution is a log entry and ages out after a week. One that is
+    // waiting for a person is work still owed an answer, and having it expire
+    // under them loses the run silently, so it lives until its own deadline
+    // plus a day of slack for a late approver.
+    int64_t ttl_ms = 7 * 24 * 60 * 60 * 1000;
+    if (result.status == ExecutionStatus::Waiting && result.pause_expires_at > 0) {
+        const int64_t remaining = result.pause_expires_at - TimeUtils::nowMs();
+        const int64_t grace = 24 * 60 * 60 * 1000;
+        if (remaining + grace > ttl_ms) {
+            ttl_ms = remaining + grace;
+        }
+    }
+
+    auto insert_result = storage_.insert("executions", result.toJson(), result.execution_id, ttl_ms);
 
     if (insert_result.failed()) {
         LOG_ERROR("Failed to store execution {}: {}", result.execution_id, insert_result.error().message());

+ 5 - 1
src/runner/workflow_engine.hpp

@@ -81,7 +81,8 @@ enum class ExecutionStatus {
     Running,
     Completed,
     Failed,
-    Cancelled
+    Cancelled,
+    Waiting
 };
 
 std::string executionStatusToString(ExecutionStatus status);
@@ -102,6 +103,9 @@ struct ExecutionResult {
     nlohmann::json final_output;
     nlohmann::json workflow_snapshot;  // Snapshot of workflow at execution time
     nlohmann::json webhook_response;   // Set by a respond-to-webhook node, if any
+    std::string paused_node_id;        // Node that asked to pause, when Waiting
+    std::string pause_token;           // Must be presented to resume
+    int64_t pause_expires_at = 0;      // Milliseconds since the epoch, 0 for never
 
     nlohmann::json toJson() const;
 };