|
@@ -925,7 +925,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
node_result = *cached;
|
|
node_result = *cached;
|
|
|
|
|
|
|
|
marker_outcome = applyNodeMarkers(node_result, node_id, node_id, result,
|
|
marker_outcome = applyNodeMarkers(node_result, node_id, node_id, result,
|
|
|
- nlohmann::json::object(), {}, -1);
|
|
|
|
|
|
|
+ nlohmann::json::object(), {}, -1, callback);
|
|
|
|
|
|
|
|
if (callback) {
|
|
if (callback) {
|
|
|
callback("node.completed", {
|
|
callback("node.completed", {
|
|
@@ -983,7 +983,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
marker_outcome = applyNodeMarkers(node_result, node_id, node_id, result,
|
|
marker_outcome = applyNodeMarkers(node_result, node_id, node_id, result,
|
|
|
- overlay, overlay_ignored, -1);
|
|
|
|
|
|
|
+ overlay, overlay_ignored, -1, callback);
|
|
|
|
|
|
|
|
if (callback) {
|
|
if (callback) {
|
|
|
nlohmann::json event_data = {
|
|
nlohmann::json event_data = {
|
|
@@ -2220,7 +2220,8 @@ WorkflowEngine::MarkerOutcome WorkflowEngine::applyNodeMarkers(
|
|
|
ExecutionResult& result,
|
|
ExecutionResult& result,
|
|
|
const nlohmann::json& overlay,
|
|
const nlohmann::json& overlay,
|
|
|
const std::vector<std::string>& overlay_ignored,
|
|
const std::vector<std::string>& overlay_ignored,
|
|
|
- int iteration) {
|
|
|
|
|
|
|
+ int iteration,
|
|
|
|
|
+ const ExecutionCallback& callback) {
|
|
|
|
|
|
|
|
const bool in_loop = iteration >= 0;
|
|
const bool in_loop = iteration >= 0;
|
|
|
const std::string where = in_loop
|
|
const std::string where = in_loop
|
|
@@ -2238,7 +2239,7 @@ WorkflowEngine::MarkerOutcome WorkflowEngine::applyNodeMarkers(
|
|
|
// has no way to reach the engine, and this is the point where the result
|
|
// has no way to reach the engine, and this is the point where the result
|
|
|
// can still become its output.
|
|
// can still become its output.
|
|
|
if (node_result.output.contains("_callWorkflow")) {
|
|
if (node_result.output.contains("_callWorkflow")) {
|
|
|
- runSubWorkflow(node_result, result.call_depth, result.execution_id);
|
|
|
|
|
|
|
+ runSubWorkflow(node_result, result.call_depth, result.execution_id, callback);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// A pause cannot work inside a loop: there is no way to resume one item of
|
|
// A pause cannot work inside a loop: there is no way to resume one item of
|
|
@@ -2397,7 +2398,8 @@ WorkflowEngine::MarkerOutcome WorkflowEngine::applyNodeMarkers(
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
bool WorkflowEngine::runSubWorkflow(NodeExecutionResult& node_result, int call_depth,
|
|
bool WorkflowEngine::runSubWorkflow(NodeExecutionResult& node_result, int call_depth,
|
|
|
- const std::string& parent_execution_id) {
|
|
|
|
|
|
|
+ const std::string& parent_execution_id,
|
|
|
|
|
+ const ExecutionCallback& parent_callback) {
|
|
|
const auto call = node_result.output["_callWorkflow"];
|
|
const auto call = node_result.output["_callWorkflow"];
|
|
|
const std::string workflow_id = call.value("workflowId", "");
|
|
const std::string workflow_id = call.value("workflowId", "");
|
|
|
node_result.output.erase("_callWorkflow");
|
|
node_result.output.erase("_callWorkflow");
|
|
@@ -2460,9 +2462,22 @@ bool WorkflowEngine::runSubWorkflow(NodeExecutionResult& node_result, int call_d
|
|
|
LOG_INFO("Execution {} calls workflow {} at depth {}", parent_execution_id, workflow_id,
|
|
LOG_INFO("Execution {} calls workflow {} at depth {}", parent_execution_id, workflow_id,
|
|
|
call_depth + 1);
|
|
call_depth + 1);
|
|
|
|
|
|
|
|
- // No callback: the sub-workflow reports its own progress against its own
|
|
|
|
|
- // execution, and forwarding its node events to the parent's subscribers
|
|
|
|
|
- // would make the parent's canvas light up nodes it does not have.
|
|
|
|
|
|
|
+ // The sub-workflow reports against ITS OWN execution and workflow, not the
|
|
|
|
|
+ // caller's.
|
|
|
|
|
+ //
|
|
|
|
|
+ // It used to be given no callback at all, on the reasoning that its node
|
|
|
|
|
+ // events would light up nodes the parent's canvas does not have. They would
|
|
|
|
|
+ // not: an event is broadcast on a channel keyed by execution id, and the
|
|
|
|
|
+ // child has its own. What the old arrangement actually cost was everything
|
|
|
|
|
+ // the child never reported - its error workflow never ran, because that is
|
|
|
|
|
+ // driven by execution.failed reaching the webserver; its runs never
|
|
|
|
|
+ // appeared live in the executions list; and its own canvas showed nothing
|
|
|
|
|
+ // while it worked. A sub-workflow that failed was silent by construction.
|
|
|
|
|
+ //
|
|
|
|
|
+ // The one thing that does have to be corrected is attribution: the runner's
|
|
|
|
|
+ // own wrapper stamps the workflow id it was built with, which is the
|
|
|
|
|
+ // caller's. This re-stamps the child's before handing the event on, and the
|
|
|
|
|
+ // wrapper now leaves an id that is already set alone.
|
|
|
// The id is generated inside execute(), so a caller cannot know it in
|
|
// The id is generated inside execute(), so a caller cannot know it in
|
|
|
// advance. It is handed in on the trigger data instead, and recorded
|
|
// advance. It is handed in on the trigger data instead, and recorded
|
|
|
// against the caller as soon as it is known - which is what lets a cancel
|
|
// against the caller as soon as it is known - which is what lets a cancel
|
|
@@ -2474,7 +2489,18 @@ bool WorkflowEngine::runSubWorkflow(NodeExecutionResult& node_result, int call_d
|
|
|
child_executions_[parent_execution_id].insert(child_execution_id);
|
|
child_executions_[parent_execution_id].insert(child_execution_id);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- auto outcome = execute(sub_workflow, "workflow", sub_trigger, nullptr);
|
|
|
|
|
|
|
+ ExecutionCallback child_callback;
|
|
|
|
|
+ if (parent_callback) {
|
|
|
|
|
+ const std::string child_workflow_id = sub_workflow.id;
|
|
|
|
|
+ child_callback = [parent_callback, child_workflow_id](const std::string& event_type,
|
|
|
|
|
+ const nlohmann::json& data) {
|
|
|
|
|
+ nlohmann::json event_data = data;
|
|
|
|
|
+ event_data["workflowId"] = child_workflow_id;
|
|
|
|
|
+ parent_callback(event_type, event_data);
|
|
|
|
|
+ };
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto outcome = execute(sub_workflow, "workflow", sub_trigger, child_callback);
|
|
|
if (outcome.failed()) {
|
|
if (outcome.failed()) {
|
|
|
node_result.status = NodeStatus::Failed;
|
|
node_result.status = NodeStatus::Failed;
|
|
|
node_result.error = "Call Workflow: " + outcome.error().message();
|
|
node_result.error = "Call Workflow: " + outcome.error().message();
|
|
@@ -3381,7 +3407,7 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
const MarkerOutcome cached_outcome = applyNodeMarkers(
|
|
const MarkerOutcome cached_outcome = applyNodeMarkers(
|
|
|
cached_result, body_node_id,
|
|
cached_result, body_node_id,
|
|
|
key_prefix + body_node_id + "_iter_" + std::to_string(i), result,
|
|
key_prefix + body_node_id + "_iter_" + std::to_string(i), result,
|
|
|
- nlohmann::json::object(), {}, static_cast<int>(i));
|
|
|
|
|
|
|
+ nlohmann::json::object(), {}, static_cast<int>(i), callback);
|
|
|
|
|
|
|
|
iteration_results[body_node_id] = cached_result;
|
|
iteration_results[body_node_id] = cached_result;
|
|
|
|
|
|
|
@@ -3430,7 +3456,22 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// Execute the body node
|
|
// Execute the body node
|
|
|
- LOG_INFO("Executing body node {} with input: {}", body_node_id, node_input.dump());
|
|
|
|
|
|
|
+ // Truncated twice over: truncateLargeValues replaces a long string
|
|
|
|
|
+ // with a note of its size, and the whole line is then capped.
|
|
|
|
|
+ //
|
|
|
|
|
+ // This used to dump the input whole. A loop body carrying an image
|
|
|
|
|
+ // wrote the entire base64 payload to the journal on every node of
|
|
|
|
|
+ // every iteration - one execution ran to 170KB of embedded JPEG,
|
|
|
|
|
+ // which is slow to write, expensive to keep, and makes the log
|
|
|
|
|
+ // unreadable exactly when something has gone wrong and you need it.
|
|
|
|
|
+ {
|
|
|
|
|
+ std::string shown = truncateLargeValues(node_input).dump();
|
|
|
|
|
+ if (shown.size() > 400) {
|
|
|
|
|
+ shown = shown.substr(0, 400) + "... (" + std::to_string(shown.size()) +
|
|
|
|
|
+ " bytes)";
|
|
|
|
|
+ }
|
|
|
|
|
+ LOG_INFO("Executing body node {} with input: {}", body_node_id, shown);
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
// Send running event before execution
|
|
// Send running event before execution
|
|
|
if (callback) {
|
|
if (callback) {
|
|
@@ -3556,7 +3597,7 @@ bool WorkflowEngine::executeLoopBody(
|
|
|
|
|
|
|
|
const MarkerOutcome body_outcome = applyNodeMarkers(
|
|
const MarkerOutcome body_outcome = applyNodeMarkers(
|
|
|
body_result, body_node_id, result_key, result,
|
|
body_result, body_node_id, result_key, result,
|
|
|
- body_overlay, body_overlay_ignored, static_cast<int>(i));
|
|
|
|
|
|
|
+ body_overlay, body_overlay_ignored, static_cast<int>(i), callback);
|
|
|
|
|
|
|
|
if (inner_loop_left_early) {
|
|
if (inner_loop_left_early) {
|
|
|
iteration_failed = true;
|
|
iteration_failed = true;
|