|
@@ -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";
|
|
@@ -266,16 +280,30 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
for (auto it = resume_seed.begin(); it != resume_seed.end(); ++it) {
|
|
for (auto it = resume_seed.begin(); it != resume_seed.end(); ++it) {
|
|
|
NodeExecutionResult seeded;
|
|
NodeExecutionResult seeded;
|
|
|
seeded.node_id = it.key();
|
|
seeded.node_id = it.key();
|
|
|
- seeded.status = NodeStatus::Completed;
|
|
|
|
|
- seeded.output = it.value();
|
|
|
|
|
|
|
+ // 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.started_at = TimeUtils::nowMs();
|
|
|
seeded.finished_at = seeded.started_at;
|
|
seeded.finished_at = seeded.started_at;
|
|
|
- seeded.from_resume = true;
|
|
|
|
|
result.node_results[it.key()] = seeded;
|
|
result.node_results[it.key()] = seeded;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- // Create workflow snapshot for pinning feature
|
|
|
|
|
|
|
+ // 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;
|
|
@@ -504,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;
|
|
@@ -545,16 +575,6 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
input = collectNodeInput(node_id, workflow, result.node_results);
|
|
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)
|
|
// Check if node should be skipped (inactive branch)
|
|
|
if (input.contains("_skip") && input["_skip"].get<bool>()) {
|
|
if (input.contains("_skip") && input["_skip"].get<bool>()) {
|
|
|
NodeExecutionResult skip_result;
|
|
NodeExecutionResult skip_result;
|
|
@@ -846,15 +866,39 @@ common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execut
|
|
|
"Execution " + execution_id + " is waiting but records no paused node");
|
|
"Execution " + execution_id + " is waiting but records no paused node");
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- if (!record.contains("workflowSnapshot")) {
|
|
|
|
|
|
|
+ if (!record.contains("workflowSnapshot") || !record["workflowSnapshot"].is_object()) {
|
|
|
return common::Error(common::ErrorCode::Internal,
|
|
return common::Error(common::ErrorCode::Internal,
|
|
|
- "Execution " + execution_id + " has no workflow snapshot to resume against");
|
|
|
|
|
|
|
+ "Execution " + execution_id + " has no usable workflow snapshot to resume against");
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
Workflow workflow = Workflow::fromJson(record["workflowSnapshot"]);
|
|
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
|
|
// Seed everything that ran before the pause, including the paused node
|
|
|
- // itself, whose output becomes the answer that was given.
|
|
|
|
|
|
|
+ // 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();
|
|
nlohmann::json seed = nlohmann::json::object();
|
|
|
if (record.contains("nodeExecutions") && record["nodeExecutions"].is_array()) {
|
|
if (record.contains("nodeExecutions") && record["nodeExecutions"].is_array()) {
|
|
|
for (const auto& entry : record["nodeExecutions"]) {
|
|
for (const auto& entry : record["nodeExecutions"]) {
|
|
@@ -862,13 +906,19 @@ common::Result<ExecutionResult> WorkflowEngine::resume(const std::string& execut
|
|
|
if (node_id.empty()) {
|
|
if (node_id.empty()) {
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
|
- seed[node_id] = entry.value("output", nlohmann::json::object());
|
|
|
|
|
|
|
+ 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();
|
|
nlohmann::json answer = payload.is_object() ? payload : nlohmann::json::object();
|
|
|
answer["answeredAt"] = TimeUtils::nowMs();
|
|
answer["answeredAt"] = TimeUtils::nowMs();
|
|
|
- seed[paused_node_id] = answer;
|
|
|
|
|
|
|
+ seed[paused_node_id] = {
|
|
|
|
|
+ {"status", "completed"},
|
|
|
|
|
+ {"output", answer}
|
|
|
|
|
+ };
|
|
|
|
|
|
|
|
nlohmann::json trigger_data = record.value("triggerData", nlohmann::json::object());
|
|
nlohmann::json trigger_data = record.value("triggerData", nlohmann::json::object());
|
|
|
trigger_data["_resumeSeed"] = seed;
|
|
trigger_data["_resumeSeed"] = seed;
|
|
@@ -1475,6 +1525,25 @@ void WorkflowEngine::storeExecution(const ExecutionResult& result) {
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // 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()) {
|
|
|
|
|
+ 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);
|
|
auto insert_result = storage_.insert("executions", result.toJson(), result.execution_id, ttl_ms);
|
|
|
|
|
|
|
|
if (insert_result.failed()) {
|
|
if (insert_result.failed()) {
|