|
|
@@ -706,13 +706,11 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
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}
|
|
|
- });
|
|
|
- }
|
|
|
+ // No callback here - the generic terminal callback below (after
|
|
|
+ // storeExecution) already emits "execution.waiting" once this
|
|
|
+ // pass ends, and it carries nodeId/expiresAt too. Emitting here
|
|
|
+ // as well produced the same event twice with two different
|
|
|
+ // payload shapes.
|
|
|
|
|
|
LOG_INFO("Execution {} paused at node {}", result.execution_id, node_id);
|
|
|
break;
|
|
|
@@ -816,13 +814,23 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
// error workflow, a notification - needs to know what went wrong, and
|
|
|
// this event is the first one out, so omitting it left listeners with a
|
|
|
// failure and no reason for it.
|
|
|
- callback("execution." + executionStatusToString(result.status), {
|
|
|
+ nlohmann::json event_data = {
|
|
|
{"executionId", result.execution_id},
|
|
|
{"workflowId", result.workflow_id},
|
|
|
{"status", executionStatusToString(result.status)},
|
|
|
{"error", result.error},
|
|
|
{"output", truncateLargeValues(result.final_output)}
|
|
|
- });
|
|
|
+ };
|
|
|
+ // A Waiting result is the one terminal status that carries reader-
|
|
|
+ // relevant fields the generic shape above doesn't have - which node
|
|
|
+ // is asking, and when the wait itself expires. This is the only
|
|
|
+ // "execution.waiting" emission (the pause block above deliberately
|
|
|
+ // does not emit its own), so those fields land here.
|
|
|
+ if (result.status == ExecutionStatus::Waiting) {
|
|
|
+ event_data["nodeId"] = result.paused_node_id;
|
|
|
+ event_data["expiresAt"] = result.pause_expires_at;
|
|
|
+ }
|
|
|
+ callback("execution." + executionStatusToString(result.status), event_data);
|
|
|
}
|
|
|
|
|
|
LOG_INFO("Workflow execution {} completed with status: {}",
|
|
|
@@ -871,6 +879,42 @@ 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.
|
|
|
+ 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
|
|
|
+ // stamped by the database on every write. Its absence means either a
|
|
|
+ // very old record predating that guarantee or something reading the
|
|
|
+ // record incorrectly - either way, resuming without a real
|
|
|
+ // compare-and-set would silently reopen the replay window this claim
|
|
|
+ // exists to close, so refuse rather than proceed unprotected.
|
|
|
+ return common::Error(common::ErrorCode::Internal,
|
|
|
+ "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
|
|
|
@@ -1535,6 +1579,36 @@ void WorkflowEngine::storeExecution(const ExecutionResult& result) {
|
|
|
// update when it is already there.
|
|
|
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()) {
|
|
|
+ LOG_ERROR("Failed to re-store waiting execution {} with a fresh TTL: {}",
|
|
|
+ result.execution_id, reinserted.error().message());
|
|
|
+ } else {
|
|
|
+ LOG_DEBUG("Re-stored waiting execution {} with a fresh TTL", result.execution_id);
|
|
|
+ }
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
auto updated = storage_.update("executions", result.execution_id, result.toJson());
|
|
|
if (updated.failed()) {
|
|
|
LOG_ERROR("Failed to update execution {}: {}", result.execution_id, updated.error().message());
|