Skip to content

Commit ddaa948

Browse files
committed
Preserve nested failure metadata
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 26daf565-a1f4-4bf0-9540-b57f31a4a12b
1 parent 25d060b commit ddaa948

5 files changed

Lines changed: 102 additions & 17 deletions

File tree

‎durabletask/extensions/history_export/serialization.py‎

Lines changed: 15 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,20 @@ def event_to_dict(event: history.HistoryEvent) -> dict[str, Any]:
5252
return {"event_type": type(event).__name__, **payload}
5353

5454

55+
def failure_details_to_dict(failure: task.FailureDetails) -> dict[str, Any]:
56+
"""Convert recursively nested failure details into a JSON-safe dict."""
57+
return {
58+
"message": failure.message,
59+
"error_type": failure.error_type,
60+
"stack_trace": failure.stack_trace,
61+
"inner_failure": (
62+
failure_details_to_dict(failure.inner_failure)
63+
if failure.inner_failure is not None else None
64+
),
65+
"properties": failure.properties,
66+
}
67+
68+
5569
def orchestration_state_to_dict(
5670
state: client_module.OrchestrationState,
5771
) -> dict[str, Any]:
@@ -61,20 +75,7 @@ def orchestration_state_to_dict(
6175
class names or module paths appear in the resulting dict.
6276
"""
6377
failure = state.failure_details
64-
failure_dict: dict[str, Any] | None = None
65-
if failure is not None:
66-
failure_dict = {
67-
"message": failure.message,
68-
"error_type": failure.error_type,
69-
"stack_trace": failure.stack_trace,
70-
}
71-
inner = getattr(failure, "inner_failure", None)
72-
if isinstance(inner, task.FailureDetails):
73-
failure_dict["inner_failure"] = {
74-
"message": inner.message,
75-
"error_type": inner.error_type,
76-
"stack_trace": inner.stack_trace,
77-
}
78+
failure_dict = failure_details_to_dict(failure) if failure is not None else None
7879
return {
7980
"instance_id": state.instance_id,
8081
"name": state.name,

‎durabletask/task.py‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -579,12 +579,18 @@ def __init__(self, message: str, details: pb.TaskFailureDetails | Exception):
579579
details.innerFailure.CopyFrom(nested_failure)
580580
for key, value in nested_failure.properties.items():
581581
details.properties[key].CopyFrom(value)
582+
self._failure_details = details
582583
self._details = pbh.failure_details_from_protobuf(details)
583584

584585
@property
585586
def details(self) -> FailureDetails:
586587
return self._details
587588

589+
@property
590+
def failure_details(self) -> pb.TaskFailureDetails:
591+
"""Return the underlying protobuf failure details."""
592+
return self._failure_details
593+
588594

589595
class NonDeterminismError(Exception):
590596
pass

‎durabletask/worker.py‎

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1624,12 +1624,21 @@ def set_failed(self, ex: Exception | pb.TaskFailureDetails):
16241624
# self._pending_actions.clear() # Cancel any pending actions
16251625
self._completion_status = pb.ORCHESTRATION_STATUS_FAILED
16261626

1627+
if isinstance(ex, task.TaskFailedError):
1628+
failure_details = ph.new_failure_details(
1629+
ex, self._exception_properties_provider, self._logger)
1630+
failure_details.innerFailure.CopyFrom(ex.failure_details)
1631+
elif isinstance(ex, Exception):
1632+
failure_details = ph.new_failure_details(
1633+
ex, self._exception_properties_provider, self._logger)
1634+
else:
1635+
failure_details = ex
1636+
16271637
action = ph.new_complete_orchestration_action(
16281638
self.next_sequence_number(),
16291639
pb.ORCHESTRATION_STATUS_FAILED,
16301640
None,
1631-
ph.new_failure_details(
1632-
ex, self._exception_properties_provider, self._logger) if isinstance(ex, Exception) else ex,
1641+
failure_details,
16331642
)
16341643
self._pending_actions[action.id] = action
16351644

‎tests/durabletask/extensions/history_export/test_serialization.py‎

Lines changed: 30 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -184,13 +184,42 @@ def test_state_to_dict_has_no_python_type_metadata(self) -> None:
184184
def test_state_with_failure_details(self) -> None:
185185
st = self._state()
186186
st.failure_details = task.FailureDetails(
187-
message="boom", error_type="RuntimeError", stack_trace="trace",
187+
message="boom",
188+
error_type="RuntimeError",
189+
stack_trace="trace",
190+
properties={"code": "OUTER"},
191+
inner_failure=task.FailureDetails(
192+
message="middle",
193+
error_type="ValueError",
194+
stack_trace=None,
195+
properties={"code": "MIDDLE"},
196+
inner_failure=task.FailureDetails(
197+
message="inner",
198+
error_type="KeyError",
199+
stack_trace=None,
200+
properties={"code": "INNER"},
201+
),
202+
),
188203
)
189204
d = orchestration_state_to_dict(st)
190205
assert d["failure_details"] == {
191206
"message": "boom",
192207
"error_type": "RuntimeError",
193208
"stack_trace": "trace",
209+
"properties": {"code": "OUTER"},
210+
"inner_failure": {
211+
"message": "middle",
212+
"error_type": "ValueError",
213+
"stack_trace": None,
214+
"properties": {"code": "MIDDLE"},
215+
"inner_failure": {
216+
"message": "inner",
217+
"error_type": "KeyError",
218+
"stack_trace": None,
219+
"properties": {"code": "INNER"},
220+
"inner_failure": None,
221+
},
222+
},
194223
}
195224

196225
def test_metadata_embedded_in_json(self) -> None:

‎tests/durabletask/test_orchestration_executor.py‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -718,6 +718,46 @@ def orchestrator(ctx: task.OrchestrationContext, orchestrator_input):
718718
assert user_code_statement in complete_action.failureDetails.stackTrace.value
719719

720720

721+
def test_unhandled_activity_failure_preserves_failure_details():
722+
"""An unhandled task failure keeps its original metadata as an inner failure."""
723+
def dummy_activity(ctx, _):
724+
pass
725+
726+
def orchestrator(ctx: task.OrchestrationContext, orchestrator_input):
727+
yield ctx.call_activity(dummy_activity, input=orchestrator_input)
728+
729+
class PropertiesProvider:
730+
def get_exception_properties(self, exception: Exception):
731+
if isinstance(exception, ValueError):
732+
return {"code": "VALUE"}
733+
return None
734+
735+
registry = worker._Registry()
736+
name = registry.add_orchestrator(orchestrator)
737+
old_events = [
738+
helpers.new_orchestrator_started_event(),
739+
helpers.new_execution_started_event(name, TEST_INSTANCE_ID, encoded_input=None),
740+
helpers.new_task_scheduled_event(1, task.get_name(dummy_activity)),
741+
]
742+
executor = worker._OrchestrationExecutor(
743+
registry,
744+
TEST_LOGGER,
745+
JsonDataConverter(),
746+
exception_properties_provider=PropertiesProvider(),
747+
)
748+
activity_failure = ValueError("boom")
749+
failed_event = helpers.new_task_failed_event(1, activity_failure)
750+
failed_event.taskFailed.failureDetails.CopyFrom(
751+
helpers.new_failure_details(activity_failure, PropertiesProvider()))
752+
753+
result = executor.execute(TEST_INSTANCE_ID, old_events, [failed_event])
754+
755+
failure = get_and_validate_complete_orchestration_action_list(1, result.actions).failureDetails
756+
assert failure.errorType == "durabletask.task.TaskFailedError"
757+
assert failure.innerFailure.errorType == "builtins.ValueError"
758+
assert failure.innerFailure.properties["code"].string_value == "VALUE"
759+
760+
721761
def test_activity_retry_policies():
722762
"""Tests the retry policy logic for activity tasks"""
723763

0 commit comments

Comments
 (0)