|
@@ -174,35 +174,26 @@ grpc::Status RunnerServiceImpl::ExecuteWorkflow(grpc::ServerContext* context,
|
|
|
LOG_INFO("Workflow execution completed, execution_id: {}, status: {}",
|
|
LOG_INFO("Workflow execution completed, execution_id: {}, status: {}",
|
|
|
result.value().execution_id, executionStatusToString(result.value().status));
|
|
result.value().execution_id, executionStatusToString(result.value().status));
|
|
|
|
|
|
|
|
- // Emit final execution event based on status
|
|
|
|
|
- if (event_callback_) {
|
|
|
|
|
- auto& exec_result = result.value();
|
|
|
|
|
- if (exec_result.status == ExecutionStatus::Completed) {
|
|
|
|
|
- event_callback_("execution.completed", {
|
|
|
|
|
- {"executionId", exec_result.execution_id},
|
|
|
|
|
- {"workflowId", workflow.id},
|
|
|
|
|
- {"output", exec_result.final_output}
|
|
|
|
|
- });
|
|
|
|
|
- } else if (exec_result.status == ExecutionStatus::Failed) {
|
|
|
|
|
- // The execution already carries the message that explains the
|
|
|
|
|
- // failure, including which node produced it. Fall back to scanning
|
|
|
|
|
- // the node results only when it does not.
|
|
|
|
|
- std::string error_msg = exec_result.error;
|
|
|
|
|
- if (error_msg.empty()) {
|
|
|
|
|
- for (const auto& [node_id, node_result] : exec_result.node_results) {
|
|
|
|
|
- if (!node_result.error.empty()) {
|
|
|
|
|
- error_msg = node_result.error;
|
|
|
|
|
- break;
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
- event_callback_("execution.failed", {
|
|
|
|
|
- {"executionId", exec_result.execution_id},
|
|
|
|
|
- {"workflowId", workflow.id},
|
|
|
|
|
- {"error", error_msg}
|
|
|
|
|
- });
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ // No terminal event is emitted here. The engine already emits one for
|
|
|
|
|
+ // every finished run, on its own way out, and this was a second copy of
|
|
|
|
|
+ // the same event: every completed run broadcast execution.completed twice
|
|
|
|
|
+ // and every failed run broadcast execution.failed twice. Confirmed on the
|
|
|
|
|
+ // wire before removing - two broadcasts a millisecond apart, from one
|
|
|
|
|
+ // "Workflow execution completed" in the runner's log.
|
|
|
|
|
+ //
|
|
|
|
|
+ // The duplicate was not just noise. The webserver runs an error workflow
|
|
|
|
|
+ // on execution.failed, so an error handler ran twice per failure, and it
|
|
|
|
|
+ // releases the scheduler slot on the same event.
|
|
|
|
|
+ //
|
|
|
|
|
+ // The engine's is also the better of the two: it covers waiting and
|
|
|
|
|
+ // cancelled rather than only the two statuses handled here, it carries the
|
|
|
|
|
+ // status and error fields, it truncates a large final output instead of
|
|
|
|
|
+ // posting the whole thing over HTTP, and it fires on the resume path too -
|
|
|
|
|
+ // which is why a resumed run already emitted exactly once and made the
|
|
|
|
|
+ // asymmetry visible.
|
|
|
|
|
+ //
|
|
|
|
|
+ // The failure branch above stays: it covers execute() itself returning an
|
|
|
|
|
+ // error, where the engine produced no result and emitted nothing.
|
|
|
|
|
|
|
|
response->set_execution_id(result.value().execution_id);
|
|
response->set_execution_id(result.value().execution_id);
|
|
|
response->set_status(toProtoStatus(result.value().status));
|
|
response->set_status(toProtoStatus(result.value().status));
|