Procházet zdrojové kódy

fix: atomic TTL refresh for waiting executions, claim resume only after validation

storeExecution used remove-then-insert to refresh the TTL on a second pause,
which could lose the record outright on a crash between the two calls, or
leave it on a stale pause state with a token the caller will never match if
remove failed. Added StorageClient::upsert, mirroring insert exactly, which
wraps the upstream atomic upsert and swaps the expiration-index entry and TTL
in one locked operation. storeExecution now uses it instead.

WorkflowEngine::resume claimed the execution (clearing its token and paused
node id) before running the validation checks below it, so a resume that got
refused by one of those checks had already bricked the record: invisible on
GET /executions/pending and unrecoverable. Moved the claim to immediately
before execute(), after every check that can refuse the resume. The version
captured at the original read is still valid to compare against, since
nothing between the read and the (now later) claim write touches the record.
fszontagh před 1 měsícem
rodič
revize
9adfccc4eb

+ 12 - 0
lib/storage/storage_client.cpp

@@ -106,6 +106,18 @@ Result<std::string> StorageClient::insert(const std::string& collection,
     return new_id;
 }
 
+Result<std::string> StorageClient::upsert(const std::string& collection,
+                                          const nlohmann::json& data,
+                                          const std::string& id,
+                                          int64_t ttl_ms) {
+    auto [new_id, is_new] = impl_->client_->upsert(collection, data, id, msToSec(ttl_ms));
+    (void)is_new;
+    if (new_id.empty()) {
+        return Error(ErrorCode::DatabaseError, "Upsert failed: " + collection);
+    }
+    return new_id;
+}
+
 Result<int64_t> StorageClient::update(const std::string& collection,
                                       const std::string& id,
                                       const nlohmann::json& data,

+ 9 - 0
lib/storage/storage_client.hpp

@@ -125,6 +125,15 @@ public:
                                        const std::string& id = "",
                                        int64_t ttl_ms = 0);
 
+    // Atomic insert-or-replace: swaps the expiration-index entry and installs
+    // the new TTL in one locked server-side operation, unlike a remove()
+    // followed by an insert() which leaves a window where the record can be
+    // lost or left stale between the two calls.
+    common::Result<std::string> upsert(const std::string& collection,
+                                       const nlohmann::json& data,
+                                       const std::string& id = "",
+                                       int64_t ttl_ms = 0);
+
     // Returns: the new version on optimistic-lock success (expected_version+1),
     //          the new version on partial/patch success,
     //          -1 when the update succeeded but the new version is unknown

+ 46 - 46
src/runner/workflow_engine.cpp

@@ -879,18 +879,10 @@ common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execut
             "Execution " + execution_id + " has no usable workflow snapshot to resume against");
     }
 
-    // Claim the execution before running it. Everything above this point only
-    // reads the record - two concurrent resumes with the same token both pass
-    // every check above and would otherwise both call execute() on the second
-    // half of the run, producing duplicate side effects (a resumed run is not
-    // idempotent - it sends emails, charges cards, posts messages). The write
-    // below moves status out of "waiting" and clears the token/pausedNodeId so
-    // a second resume cannot match either the status check or the token check.
-    // storage_.update() with a non-zero expected_version goes through upstream
-    // updateIfVersion(), a genuine server-side compare-and-set (verified against
-    // smartbotic-database's client.cpp: it sends expected_version on the wire
-    // and the server rejects a stale write), so a racing second resume loses
-    // this write and its execute() call never happens.
+    // The version to claim with is captured here, at the same point the record
+    // itself was read - nothing between here and the claim write (below, right
+    // before execute()) writes to this record, so the version stays valid to
+    // compare against no matter how much read-only validation runs in between.
     const int64_t expected_version = record.value("_version", static_cast<int64_t>(0));
     if (expected_version <= 0) {
         // Every document returned by get() carries a _version metadata field
@@ -903,18 +895,6 @@ common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execut
             "Execution " + execution_id + " has no version metadata to claim it safely with");
     }
 
-    nlohmann::json claim = record;
-    claim.erase("_version");
-    claim["status"] = "resuming";
-    claim["pauseToken"] = "";
-    claim["pausedNodeId"] = "";
-
-    auto claim_result = storage_.update("executions", execution_id, claim, expected_version);
-    if (claim_result.failed()) {
-        return common::Error(common::ErrorCode::FailedPrecondition,
-            "Execution " + execution_id + " is already being resumed");
-    }
-
     Workflow workflow = Workflow::fromJson(record["workflowSnapshot"]);
 
     // Defend against snapshots stored before id/name/settings were added to
@@ -971,6 +951,34 @@ common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execut
     LOG_INFO("Resuming execution {} from node {} with {} seeded results",
              execution_id, paused_node_id, seed.size());
 
+    // Claim the execution immediately before running it, only now that every
+    // check that can refuse the resume has passed. Two concurrent resumes with
+    // the same token both pass every check above and would otherwise both call
+    // execute() on the second half of the run, producing duplicate side effects
+    // (a resumed run is not idempotent - it sends emails, charges cards, posts
+    // messages). The write below moves status out of "waiting" and clears the
+    // token/pausedNodeId so a second resume cannot match either the status
+    // check or the token check. storage_.update() with a non-zero
+    // expected_version goes through upstream updateIfVersion(), a genuine
+    // server-side compare-and-set (verified against smartbotic-database's
+    // client.cpp: it sends expected_version on the wire and the server rejects
+    // a stale write), so a racing second resume loses this write and its
+    // execute() call never happens. Claiming this late, instead of before
+    // validation, means a resume that gets refused below never touches the
+    // record: it stays "waiting", still visible on GET /executions/pending,
+    // and can be retried.
+    nlohmann::json claim = record;
+    claim.erase("_version");
+    claim["status"] = "resuming";
+    claim["pauseToken"] = "";
+    claim["pausedNodeId"] = "";
+
+    auto claim_result = storage_.update("executions", execution_id, claim, expected_version);
+    if (claim_result.failed()) {
+        return common::Error(common::ErrorCode::FailedPrecondition,
+            "Execution " + execution_id + " is already being resumed");
+    }
+
     return execute(workflow, record.value("triggerType", "resume"), trigger_data, callback);
 }
 
@@ -1580,29 +1588,21 @@ void WorkflowEngine::storeExecution(const ExecutionResult& result) {
     auto existing = storage_.get("executions", result.execution_id);
     if (existing.ok()) {
         if (result.status == ExecutionStatus::Waiting) {
-            // storage_.update() has no ttl_ms parameter - only insert() can set
-            // one. A second pause on an already-existing record recomputes
-            // ttl_ms above from its own (possibly much later) pause_expires_at,
-            // but a plain update() would leave the record on whatever TTL its
-            // very first insert got, undoing the floor above for every pause
-            // after the first - exactly the silent-data-loss shape called out
-            // elsewhere on this branch. Remove and reinsert so the fresh
-            // deadline actually lands. This opens a narrow window between the
-            // two calls where a process crash mid-write would lose the record
-            // outright (the remove lands, the insert never runs); accepted
-            // here because it only applies to a Waiting execution, the window
-            // is a single pair of adjacent calls rather than a standing gap,
-            // and the alternative is a certainty rather than a small chance:
-            // every long-lived second approval would otherwise expire early.
-            auto removed = storage_.remove("executions", result.execution_id);
-            if (removed.failed()) {
-                LOG_ERROR("Failed to remove execution {} before re-inserting with a fresh TTL: {}",
-                          result.execution_id, removed.error().message());
-            }
-            auto reinserted = storage_.insert("executions", result.toJson(), result.execution_id, ttl_ms);
-            if (reinserted.failed()) {
+            // storage_.update() has no ttl_ms parameter - only insert() and
+            // upsert() can set one. A second pause on an already-existing
+            // record recomputes ttl_ms above from its own (possibly much
+            // later) pause_expires_at, but a plain update() would leave the
+            // record on whatever TTL its very first insert got, undoing the
+            // floor above for every pause after the first - exactly the
+            // silent-data-loss shape called out elsewhere on this branch.
+            // upsert() swaps the expiration-index entry and installs the new
+            // TTL in one locked server-side operation, so the fresh deadline
+            // lands without the remove-then-insert window where a crash
+            // between the two calls could lose the record outright.
+            auto upserted = storage_.upsert("executions", result.toJson(), result.execution_id, ttl_ms);
+            if (upserted.failed()) {
                 LOG_ERROR("Failed to re-store waiting execution {} with a fresh TTL: {}",
-                          result.execution_id, reinserted.error().message());
+                          result.execution_id, upserted.error().message());
             } else {
                 LOG_DEBUG("Re-stored waiting execution {} with a fresh TTL", result.execution_id);
             }