|
@@ -23,6 +23,20 @@ std::string nodeStatusToString(NodeStatus status) {
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+// Inverse of nodeStatusToString, used to seed a resumed execution's node
|
|
|
|
|
+// results with the status each node actually finished with rather than
|
|
|
|
|
+// forcing everything to Completed. Anything unrecognised defaults to
|
|
|
|
|
+// Completed, matching prior behavior for records that predate this field.
|
|
|
|
|
+static NodeStatus nodeStatusFromString(const std::string& s) {
|
|
|
|
|
+ if (s == "pending") return NodeStatus::Pending;
|
|
|
|
|
+ if (s == "running") return NodeStatus::Running;
|
|
|
|
|
+ if (s == "completed") return NodeStatus::Completed;
|
|
|
|
|
+ if (s == "failed") return NodeStatus::Failed;
|
|
|
|
|
+ if (s == "skipped") return NodeStatus::Skipped;
|
|
|
|
|
+ if (s == "disabled") return NodeStatus::Disabled;
|
|
|
|
|
+ return NodeStatus::Completed;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
std::string executionStatusToString(ExecutionStatus status) {
|
|
std::string executionStatusToString(ExecutionStatus status) {
|
|
|
switch (status) {
|
|
switch (status) {
|
|
|
case ExecutionStatus::Pending: return "pending";
|
|
case ExecutionStatus::Pending: return "pending";
|
|
@@ -30,6 +44,7 @@ std::string executionStatusToString(ExecutionStatus status) {
|
|
|
case ExecutionStatus::Completed: return "completed";
|
|
case ExecutionStatus::Completed: return "completed";
|
|
|
case ExecutionStatus::Failed: return "failed";
|
|
case ExecutionStatus::Failed: return "failed";
|
|
|
case ExecutionStatus::Cancelled: return "cancelled";
|
|
case ExecutionStatus::Cancelled: return "cancelled";
|
|
|
|
|
+ case ExecutionStatus::Waiting: return "waiting";
|
|
|
default: return "unknown";
|
|
default: return "unknown";
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
@@ -139,6 +154,14 @@ nlohmann::json ExecutionResult::toJson() const {
|
|
|
j["error"] = error;
|
|
j["error"] = error;
|
|
|
j["output"] = truncateLargeValues(final_output);
|
|
j["output"] = truncateLargeValues(final_output);
|
|
|
|
|
|
|
|
|
|
+ // A finished execution's node outputs are a log, so large strings are
|
|
|
|
|
+ // truncated. A Waiting execution's node outputs are the resume state a
|
|
|
|
|
+ // later resume seeds itself from - truncating them would hand the resume
|
|
|
|
|
+ // a literal "[omitted N bytes]" placeholder in place of real data, so
|
|
|
|
|
+ // they are kept whole. Do not "restore consistency" here later; that
|
|
|
|
|
+ // would silently reintroduce truncated resume data.
|
|
|
|
|
+ const bool keep_full_outputs = (status == ExecutionStatus::Waiting);
|
|
|
|
|
+
|
|
|
j["nodeExecutions"] = nlohmann::json::array();
|
|
j["nodeExecutions"] = nlohmann::json::array();
|
|
|
for (const auto& [id, result] : node_results) {
|
|
for (const auto& [id, result] : node_results) {
|
|
|
nlohmann::json nr;
|
|
nlohmann::json nr;
|
|
@@ -148,7 +171,7 @@ nlohmann::json ExecutionResult::toJson() const {
|
|
|
nr["finishedAt"] = result.finished_at;
|
|
nr["finishedAt"] = result.finished_at;
|
|
|
// Note: input is intentionally not stored to avoid data duplication
|
|
// Note: input is intentionally not stored to avoid data duplication
|
|
|
// Each node's input can be reconstructed from upstream node outputs + connections
|
|
// Each node's input can be reconstructed from upstream node outputs + connections
|
|
|
- nr["output"] = truncateLargeValues(result.output);
|
|
|
|
|
|
|
+ nr["output"] = keep_full_outputs ? result.output : truncateLargeValues(result.output);
|
|
|
nr["error"] = result.error;
|
|
nr["error"] = result.error;
|
|
|
nr["retryCount"] = result.retry_count;
|
|
nr["retryCount"] = result.retry_count;
|
|
|
j["nodeExecutions"].push_back(nr);
|
|
j["nodeExecutions"].push_back(nr);
|
|
@@ -159,6 +182,16 @@ nlohmann::json ExecutionResult::toJson() const {
|
|
|
j["workflowSnapshot"] = workflow_snapshot;
|
|
j["workflowSnapshot"] = workflow_snapshot;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ if (!webhook_response.is_null()) {
|
|
|
|
|
+ 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;
|
|
return j;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -216,6 +249,20 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
LOG_INFO("Received cached outputs for {} nodes", cached_outputs.size());
|
|
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
|
|
// Create execution record
|
|
|
ExecutionResult result;
|
|
ExecutionResult result;
|
|
|
result.execution_id = UUID::generatePrefixed("exec");
|
|
result.execution_id = UUID::generatePrefixed("exec");
|
|
@@ -227,8 +274,36 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
result.trigger_data = actual_trigger_data;
|
|
result.trigger_data = actual_trigger_data;
|
|
|
result.started_at = TimeUtils::nowMs();
|
|
result.started_at = TimeUtils::nowMs();
|
|
|
|
|
|
|
|
- // Create workflow snapshot for pinning feature
|
|
|
|
|
|
|
+ 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();
|
|
|
|
|
+ // The seed carries the real status the node finished with before the
|
|
|
|
|
+ // pause (completed, failed, skipped, disabled). Seeding everything as
|
|
|
|
|
+ // Completed would make collectNodeInput stop propagating a skip, and a
|
|
|
|
|
+ // node downstream of a branch that was not taken could run on the
|
|
|
|
|
+ // resumed pass having never run on the original one.
|
|
|
|
|
+ seeded.status = nodeStatusFromString(it.value().value("status", "completed"));
|
|
|
|
|
+ seeded.output = it.value().value("output", nlohmann::json::object());
|
|
|
|
|
+ seeded.started_at = TimeUtils::nowMs();
|
|
|
|
|
+ seeded.finished_at = seeded.started_at;
|
|
|
|
|
+ result.node_results[it.key()] = seeded;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Create workflow snapshot for pinning feature. id/name/settings are
|
|
|
|
|
+ // included, not only nodes and connections: Workflow::fromJson rebuilds a
|
|
|
|
|
+ // Workflow from this on resume, and an empty workflow.id is treated by
|
|
|
|
|
+ // credential access checks as an admin operation with unrestricted
|
|
|
|
|
+ // credential access (see CredentialStore::hasWorkflowAccess), so an empty
|
|
|
|
|
+ // id here is a privilege escalation, not a cosmetic gap. settings carries
|
|
|
|
|
+ // storagePermissions and continueOnError, which the resumed half of the
|
|
|
|
|
+ // run must honor the same way the first half did.
|
|
|
nlohmann::json snapshot;
|
|
nlohmann::json snapshot;
|
|
|
|
|
+ snapshot["id"] = workflow.id;
|
|
|
|
|
+ snapshot["name"] = workflow.name;
|
|
|
|
|
+ snapshot["settings"] = workflow.settings;
|
|
|
snapshot["nodes"] = nlohmann::json::array();
|
|
snapshot["nodes"] = nlohmann::json::array();
|
|
|
for (const auto& node : workflow.nodes) {
|
|
for (const auto& node : workflow.nodes) {
|
|
|
nlohmann::json n;
|
|
nlohmann::json n;
|
|
@@ -457,8 +532,10 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- // Skip nodes that were already executed (e.g., loop body nodes)
|
|
|
|
|
- // They have results stored during loop execution
|
|
|
|
|
|
|
+ // Skip nodes that were already executed (e.g., loop body nodes, or -
|
|
|
|
|
+ // on a resume - nodes seeded from the stored execution that ran
|
|
|
|
|
+ // before the pause). This is the mechanism that keeps a resume from
|
|
|
|
|
+ // re-running work: every seeded node already has a result here.
|
|
|
if (result.node_results.contains(node_id)) {
|
|
if (result.node_results.contains(node_id)) {
|
|
|
LOG_DEBUG("Node {} already has results, skipping in main loop", node_id);
|
|
LOG_DEBUG("Node {} already has results, skipping in main loop", node_id);
|
|
|
continue;
|
|
continue;
|
|
@@ -600,6 +677,45 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // Any node may declare the HTTP response, not just the last one to
|
|
|
|
|
+ // run, so appending a node to a workflow cannot silently change
|
|
|
|
|
+ // what its webhook returns.
|
|
|
|
|
+ if (node_result.status == NodeStatus::Completed &&
|
|
|
|
|
+ node_result.output.contains("_webhookResponse")) {
|
|
|
|
|
+ if (!result.webhook_response.is_null()) {
|
|
|
|
|
+ LOG_WARN("Node {} overrides a webhook response already set by an earlier node",
|
|
|
|
|
+ node_id);
|
|
|
|
|
+ }
|
|
|
|
|
+ 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");
|
|
|
|
|
+
|
|
|
|
|
+ // 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;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// Check for loop node
|
|
// Check for loop node
|
|
|
if (node_result.status == NodeStatus::Completed &&
|
|
if (node_result.status == NodeStatus::Completed &&
|
|
|
node_result.output.contains("_isLoop") &&
|
|
node_result.output.contains("_isLoop") &&
|
|
@@ -698,13 +814,23 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
// error workflow, a notification - needs to know what went wrong, and
|
|
// 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
|
|
// this event is the first one out, so omitting it left listeners with a
|
|
|
// failure and no reason for it.
|
|
// failure and no reason for it.
|
|
|
- callback("execution." + executionStatusToString(result.status), {
|
|
|
|
|
|
|
+ nlohmann::json event_data = {
|
|
|
{"executionId", result.execution_id},
|
|
{"executionId", result.execution_id},
|
|
|
{"workflowId", result.workflow_id},
|
|
{"workflowId", result.workflow_id},
|
|
|
{"status", executionStatusToString(result.status)},
|
|
{"status", executionStatusToString(result.status)},
|
|
|
{"error", result.error},
|
|
{"error", result.error},
|
|
|
{"output", truncateLargeValues(result.final_output)}
|
|
{"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: {}",
|
|
LOG_INFO("Workflow execution {} completed with status: {}",
|
|
@@ -713,6 +839,149 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
return result;
|
|
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") || !record["workflowSnapshot"].is_object()) {
|
|
|
|
|
+ return common::Error(common::ErrorCode::Internal,
|
|
|
|
|
+ "Execution " + execution_id + " has no usable workflow snapshot to resume against");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // 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
|
|
|
|
|
+ // 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");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ Workflow workflow = Workflow::fromJson(record["workflowSnapshot"]);
|
|
|
|
|
+
|
|
|
|
|
+ // Defend against snapshots stored before id/name/settings were added to
|
|
|
|
|
+ // them: an execution paused under the old shape still has workflowId and
|
|
|
|
|
+ // workflowName on the record itself, so those two are recoverable. An
|
|
|
|
|
+ // empty workflow.id would otherwise be read by CredentialStore::
|
|
|
|
|
+ // hasWorkflowAccess as an admin operation with unrestricted credential
|
|
|
|
|
+ // access, so this is a security backstop, not tidiness.
|
|
|
|
|
+ if (workflow.id.empty()) {
|
|
|
|
|
+ workflow.id = record.value("workflowId", "");
|
|
|
|
|
+ }
|
|
|
|
|
+ if (workflow.name.empty()) {
|
|
|
|
|
+ workflow.name = record.value("workflowName", "");
|
|
|
|
|
+ }
|
|
|
|
|
+ // settings is not recoverable from the record - it was never stored
|
|
|
|
|
+ // anywhere else. If the snapshot carried no settings, refuse rather than
|
|
|
|
|
+ // run the resumed half of the execution under different storage
|
|
|
|
|
+ // permissions and continueOnError than the half that ran before the pause.
|
|
|
|
|
+ if (workflow.settings.empty() && !record["workflowSnapshot"].contains("settings")) {
|
|
|
|
|
+ return common::Error(common::ErrorCode::FailedPrecondition,
|
|
|
|
|
+ "Execution " + execution_id + " was paused before workflow settings were captured "
|
|
|
|
|
+ "in its snapshot and cannot be safely resumed");
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Seed everything that ran before the pause, including the paused node
|
|
|
|
|
+ // itself, whose output becomes the answer that was given. Status is
|
|
|
|
|
+ // carried along with output so a node that was skipped or failed before
|
|
|
|
|
+ // the pause is not reported as Completed on the resumed pass.
|
|
|
|
|
+ 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] = {
|
|
|
|
|
+ {"status", entry.value("status", "completed")},
|
|
|
|
|
+ {"output", entry.value("output", nlohmann::json::object())}
|
|
|
|
|
+ };
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ nlohmann::json answer = payload.is_object() ? payload : nlohmann::json::object();
|
|
|
|
|
+ answer["answeredAt"] = TimeUtils::nowMs();
|
|
|
|
|
+ seed[paused_node_id] = {
|
|
|
|
|
+ {"status", "completed"},
|
|
|
|
|
+ {"output", 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());
|
|
|
|
|
+
|
|
|
|
|
+ // 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);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
void WorkflowEngine::cancelExecution(const std::string& execution_id) {
|
|
void WorkflowEngine::cancelExecution(const std::string& execution_id) {
|
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
|
cancelled_executions_.insert(execution_id);
|
|
cancelled_executions_.insert(execution_id);
|
|
@@ -902,6 +1171,7 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
engine::ScriptContext ctx;
|
|
engine::ScriptContext ctx;
|
|
|
ctx.execution_id = execution_id;
|
|
ctx.execution_id = execution_id;
|
|
|
ctx.node_id = node.id;
|
|
ctx.node_id = node.id;
|
|
|
|
|
+ ctx.workflow_id = workflow.id;
|
|
|
ctx.input = input;
|
|
ctx.input = input;
|
|
|
ctx.config = node.config;
|
|
ctx.config = node.config;
|
|
|
ctx.log_handler = [&node](const std::string& level, const std::string& msg) {
|
|
ctx.log_handler = [&node](const std::string& level, const std::string& msg) {
|
|
@@ -1287,8 +1557,68 @@ nlohmann::json WorkflowEngine::collectNodeInput(
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
void WorkflowEngine::storeExecution(const ExecutionResult& result) {
|
|
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. An expiresAt of 0 means "no
|
|
|
|
|
+ // deadline", not "no need to extend" - it is given the 30-day ceiling the
|
|
|
|
|
+ // approval node caps at, so a permanent approval does not fall back to
|
|
|
|
|
+ // the ordinary seven-day log TTL and evaporate with no trace.
|
|
|
|
|
+ int64_t ttl_ms = 7 * 24 * 60 * 60 * 1000;
|
|
|
|
|
+ if (result.status == ExecutionStatus::Waiting) {
|
|
|
|
|
+ const int64_t grace = 24 * 60 * 60 * 1000;
|
|
|
|
|
+ const int64_t never_ttl = 30LL * 24 * 60 * 60 * 1000;
|
|
|
|
|
+ int64_t wanted = never_ttl;
|
|
|
|
|
+ if (result.pause_expires_at > 0) {
|
|
|
|
|
+ wanted = result.pause_expires_at - TimeUtils::nowMs() + grace;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (wanted > ttl_ms) {
|
|
|
|
|
+ ttl_ms = wanted;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // A resumed execution reuses the id of the execution that paused, which
|
|
|
|
|
+ // already exists in the database. insert on an id that already exists is
|
|
|
|
|
+ // not guaranteed to behave like an upsert - if it rejects the duplicate,
|
|
|
|
|
+ // logging and returning here would leave the record at status "waiting"
|
|
|
|
|
+ // with its original pauseToken, letting the same resume be replayed
|
|
|
|
|
+ // indefinitely while this function still reports the run as Completed. So
|
|
|
|
|
+ // the record's existence is checked first and the write goes through
|
|
|
|
|
+ // 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() 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, upserted.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());
|
|
|
|
|
+ } else {
|
|
|
|
|
+ LOG_DEBUG("Updated execution {} successfully", result.execution_id);
|
|
|
|
|
+ }
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto insert_result = storage_.insert("executions", result.toJson(), result.execution_id, ttl_ms);
|
|
|
|
|
|
|
|
if (insert_result.failed()) {
|
|
if (insert_result.failed()) {
|
|
|
LOG_ERROR("Failed to store execution {}: {}", result.execution_id, insert_result.error().message());
|
|
LOG_ERROR("Failed to store execution {}: {}", result.execution_id, insert_result.error().message());
|
|
@@ -1459,6 +1789,23 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
LOG_INFO(" - Body node: {}", nid);
|
|
LOG_INFO(" - Body node: {}", nid);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // A completed body result - whether freshly executed or replayed from a
|
|
|
|
|
+ // pinned cache entry - that still carries a pause marker must not be
|
|
|
|
|
+ // replayed as Completed: the iteration state that would let it resume
|
|
|
|
|
+ // does not exist in this execution's stored data. Both the cache path and
|
|
|
|
|
+ // the fresh-execution path route through this single check so the
|
|
|
|
|
+ // failure and its message cannot drift apart between them.
|
|
|
|
|
+ auto rejectPauseInLoopBody = [](NodeExecutionResult& body_result) {
|
|
|
|
|
+ if (body_result.status == NodeStatus::Completed &&
|
|
|
|
|
+ body_result.output.contains("_pause")) {
|
|
|
|
|
+ body_result.status = NodeStatus::Failed;
|
|
|
|
|
+ body_result.error = "A node cannot pause inside a Loop body, because a "
|
|
|
|
|
+ "paused loop iteration cannot be resumed. Collect the "
|
|
|
|
|
+ "items first, approve once, then loop";
|
|
|
|
|
+ body_result.output.erase("_pause");
|
|
|
|
|
+ }
|
|
|
|
|
+ };
|
|
|
|
|
+
|
|
|
// Execute body for each item
|
|
// Execute body for each item
|
|
|
for (size_t i = 0; i < ctx.items.size(); ++i) {
|
|
for (size_t i = 0; i < ctx.items.size(); ++i) {
|
|
|
ctx.current_index = i;
|
|
ctx.current_index = i;
|
|
@@ -1617,21 +1964,49 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
cached_result.finished_at = cached_result.started_at;
|
|
cached_result.finished_at = cached_result.started_at;
|
|
|
cached_result.from_cache = true;
|
|
cached_result.from_cache = true;
|
|
|
|
|
|
|
|
|
|
+ rejectPauseInLoopBody(cached_result);
|
|
|
|
|
+
|
|
|
iteration_results[body_node_id] = cached_result;
|
|
iteration_results[body_node_id] = cached_result;
|
|
|
result.node_results[body_node_id + "_iter_" + std::to_string(i)] = cached_result;
|
|
result.node_results[body_node_id + "_iter_" + std::to_string(i)] = cached_result;
|
|
|
|
|
|
|
|
|
|
+ if (cached_result.status == NodeStatus::Completed &&
|
|
|
|
|
+ cached_result.output.contains("_webhookResponse")) {
|
|
|
|
|
+ if (!result.webhook_response.is_null()) {
|
|
|
|
|
+ LOG_WARN("Node {} in a loop body overrides a webhook response already set",
|
|
|
|
|
+ body_node_id);
|
|
|
|
|
+ }
|
|
|
|
|
+ result.webhook_response = cached_result.output["_webhookResponse"];
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
if (callback) {
|
|
if (callback) {
|
|
|
- callback("loop.node.completed", {
|
|
|
|
|
|
|
+ nlohmann::json cache_event_data = {
|
|
|
{"executionId", result.execution_id},
|
|
{"executionId", result.execution_id},
|
|
|
{"nodeId", body_node_id},
|
|
{"nodeId", body_node_id},
|
|
|
{"iteration", i},
|
|
{"iteration", i},
|
|
|
- {"status", "completed"},
|
|
|
|
|
|
|
+ {"status", nodeStatusToString(cached_result.status)},
|
|
|
{"output", cached_result.output},
|
|
{"output", cached_result.output},
|
|
|
{"fromCache", true}
|
|
{"fromCache", true}
|
|
|
- });
|
|
|
|
|
|
|
+ };
|
|
|
|
|
+ if (!cached_result.error.empty()) {
|
|
|
|
|
+ cache_event_data["error"] = cached_result.error;
|
|
|
|
|
+ }
|
|
|
|
|
+ callback("loop.node." + nodeStatusToString(cached_result.status), cache_event_data);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
LOG_INFO("Using pinned output for body node {}", body_node_id);
|
|
LOG_INFO("Using pinned output for body node {}", body_node_id);
|
|
|
|
|
+
|
|
|
|
|
+ // A pinned output rejected above by rejectPauseInLoopBody must go
|
|
|
|
|
+ // through the same continueOnError decision a freshly executed
|
|
|
|
|
+ // failure does, rather than silently moving on to the next body
|
|
|
|
|
+ // node as a plain cache hit would.
|
|
|
|
|
+ if (cached_result.status == NodeStatus::Failed) {
|
|
|
|
|
+ iteration_failed = true;
|
|
|
|
|
+ all_succeeded = false;
|
|
|
|
|
+ if (!ctx.continue_on_error) {
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -1674,12 +2049,24 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
} else {
|
|
} else {
|
|
|
body_result = executeNode(evaluated_body_node, node_input, result.execution_id, workflow);
|
|
body_result = executeNode(evaluated_body_node, node_input, result.execution_id, workflow);
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
|
|
+ rejectPauseInLoopBody(body_result);
|
|
|
|
|
+
|
|
|
iteration_results[body_node_id] = body_result;
|
|
iteration_results[body_node_id] = body_result;
|
|
|
|
|
|
|
|
// Store in main results with iteration suffix
|
|
// Store in main results with iteration suffix
|
|
|
std::string result_key = body_node_id + "_iter_" + std::to_string(i);
|
|
std::string result_key = body_node_id + "_iter_" + std::to_string(i);
|
|
|
result.node_results[result_key] = body_result;
|
|
result.node_results[result_key] = body_result;
|
|
|
|
|
|
|
|
|
|
+ if (body_result.status == NodeStatus::Completed &&
|
|
|
|
|
+ body_result.output.contains("_webhookResponse")) {
|
|
|
|
|
+ if (!result.webhook_response.is_null()) {
|
|
|
|
|
+ LOG_WARN("Node {} in a loop body overrides a webhook response already set",
|
|
|
|
|
+ body_node_id);
|
|
|
|
|
+ }
|
|
|
|
|
+ result.webhook_response = body_result.output["_webhookResponse"];
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
if (callback) {
|
|
if (callback) {
|
|
|
nlohmann::json event_data = {
|
|
nlohmann::json event_data = {
|
|
|
{"executionId", result.execution_id},
|
|
{"executionId", result.execution_id},
|