Skip to content
Prev Previous commit
Next Next commit
UN-3016 [CLEANUP] Simplify per /simplify; test the value, not the sou…
…rce shape

Behaviour-preserving cleanup, run once after the review verdict settled.

- views.py: drop the second cleanup handler added in the previous commit and
  stage inside the existing try instead. That handler's guard is exactly the
  staging condition, so one handler now covers both paths rather than two
  copies of the same contract.

- test_un3016_execution_error.py: import callback.tasks directly instead of
  extracting the helper with ast/exec. The docstring's claim that importing
  pulls in an unusable celery runtime is false — conftest loads .env.test
  before collection, and test_pg_callback_duplicate_guard.py already imports
  the module at module level. Verified by running the import under pytest.
  test_status_function_returns_a_reason is now behavioural: it calls
  _determine_execution_status_unified and asserts the reason is non-blank.
  The old version asserted only that every return was a 4-tuple, which would
  have passed with an always-None fourth element — i.e. it could not detect
  the very defect it was named for. Mutation-checked: forcing
  error_message = None now fails the test.

- _summarize_file_errors: errors is keyed by file name, so entries are
  distinct by construction and the `entry not in seen` dedup could never
  fire. Removed, along with the duplicate early return it guarded.

- source.py: skipped_files was a dict never used as a mapping; now a list of
  pre-formatted entries.

Tests: 1298 passed, 132 skipped. test_pg_reaper.py deselected — it needs a
live Postgres on 127.0.0.1:5432 and hangs identically on unmodified HEAD.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NQhqNqCwFcXZZ7cQU6HxUE
  • Loading branch information
hari-kuriakose and claude committed Aug 30, 2026
commit af64650ca96d28f40864b299b3d20c1868c9c2b3
14 changes: 6 additions & 8 deletions backend/workflow_manager/endpoint_v2/source.py
Original file line number Diff line number Diff line change
Expand Up @@ -1222,9 +1222,9 @@ def add_input_file_to_api_storage(
workflow: Workflow = Workflow.objects.get(id=workflow_id)
file_hashes: dict[str, FileHash] = {}
unique_file_hashes: set[str] = set()
# UN-3016: files rejected for an unsupported MIME type, kept only to
# report them; they are never staged and never handed to a worker.
skipped_files: dict[str, str] = {}
# UN-3016: files rejected for an unsupported MIME type, pre-formatted
# for the error below; they are never staged and never handed to a worker.
skipped_files: list[str] = []
connection_type = WorkflowEndpoint.ConnectionType.API
for file in file_objs:
file_name = file.name
Expand Down Expand Up @@ -1252,7 +1252,7 @@ def add_input_file_to_api_storage(
f"'{mime_type}'. It will not be processed."
)
workflow_log.log_error(logger=logger, message=log_message)
skipped_files[file_name] = mime_type
skipped_files.append(f"'{file_name}' ({mime_type})")
continue

file_system = FileSystem(FileStorageType.API_EXECUTION)
Expand Down Expand Up @@ -1292,11 +1292,9 @@ def add_input_file_to_api_storage(
# which would otherwise finish as a vacuous success and leave the user
# wondering why nothing happened.
if skipped_files and not file_hashes:
details = ", ".join(
f"'{name}' ({mime})" for name, mime in skipped_files.items()
)
raise UnsupportedMimeTypeError(
f"No files could be processed. Unsupported file type(s): {details}"
"No files could be processed. Unsupported file type(s): "
+ ", ".join(skipped_files)
)

return file_hashes
Expand Down
26 changes: 8 additions & 18 deletions backend/workflow_manager/workflow_v2/views.py
Original file line number Diff line number Diff line change
Expand Up @@ -258,31 +258,21 @@ def execute(
use_file_history: bool = True

hashes_of_files: dict[str, FileHash] = {}
if file_objs and execution_id and workflow_id:
has_uploads = bool(file_objs and execution_id and workflow_id)
if has_uploads:
use_file_history = False
# Staging sits outside the main try/except below, so it needs its own
# cleanup: a partial stage can leave already-written files behind.
# Mirrors the handler at the end of this method.
try:

try:
# Staged inside this try so the handler below cleans up after a
# partial stage: its guard is exactly this staging condition.
if has_uploads:
hashes_of_files = SourceConnector.add_input_file_to_api_storage(
pipeline_id=pipeline_guid,
workflow_id=workflow_id,
execution_id=execution_id,
file_objs=file_objs,
use_file_history=False,
)
except Exception as exception:
logger.error(
f"Error while staging files for execution {execution_id}: "
f"{exception}",
exc_info=True,
)
DestinationConnector.delete_api_storage_dir(
workflow_id=workflow_id, execution_id=execution_id
)
raise

try:
workflow = self.get_workflow_by_id(workflow_id=workflow_id)
Comment on lines +268 to 276

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[High] [Lens 13 — Testing] — The staging move is unpinned against dropping the staged hashes

Both tests in test_un3016_execute_staging_cleanup.py (:57, :77) raise before execute_workflow is reached, so nothing asserts the staged hashes still flow through.

Losing this assignment dispatches an execution with hash_values_of_files={} — the run reports success having processed zero uploaded files. A silent correctness failure on the primary upload flow.

Mutation executed. Changing :269 to drop only the assignment → 2 passed. No other test exercises WorkflowViewSet.execute().

Fix: a third test asserting execute_workflow.call_args.kwargs["hash_values_of_files"] is the sentinel returned by staging, and that use_file_history is False.

Confidence: High (mutation executed).

execution_response = self.execute_workflow(
workflow=workflow,
Expand All @@ -304,7 +294,7 @@ def execute(
)
except Exception as exception:
logger.error(f"Error while executing workflow: {exception}", exc_info=True)
if file_objs and execution_id and workflow_id:
if has_uploads:
DestinationConnector.delete_api_storage_dir(
workflow_id=workflow_id, execution_id=execution_id
)
Expand Down
30 changes: 11 additions & 19 deletions workers/callback/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -270,27 +270,19 @@ def _summarize_file_errors(aggregated_results: dict[str, Any], total_files: int)
errors are already aggregated as {file_name: error}; surface them here.
"""
errors: dict[str, Any] = aggregated_results.get("errors") or {}
if not errors:
return f"All {total_files} file(s) failed."

# Report the distinct reasons rather than repeating an identical message
# once per file; cap the detail so a large batch cannot bloat the column.
seen: list[str] = []
for file_name, error in errors.items():
detail = str(error).strip() if error else ""
if not detail:
continue
entry = f"{file_name}: {detail}"
if entry not in seen:
seen.append(entry)

if not seen:
# `errors` is keyed by file name, so every entry is distinct already; the
# cap is what keeps a large batch from bloating the column.
entries = [
f"{file_name}: {str(error).strip()}"
for file_name, error in errors.items()
if error and str(error).strip()
]
if not entries:
return f"All {total_files} file(s) failed."
Comment on lines +265 to +281

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[High] [Lens 1, 3, 10, 13] — The error summary is a tautology in production: aggregated_results["errors"] is always empty

_summarize_file_errors reads only aggregated_results["errors"]. That map admits an entry only when file_result.get("status") == "error" (lowercase). No producer ever emits that string:

  • API path — workers/file_processing/tasks.py:1844 wraps every result in FileExecutionResult, whose __post_init__ (unstract/core/src/unstract/core/worker_models.py:301-305) forces status to ApiDeploymentResultStatus.SUCCESS/FAILED = "Success" / "Failed" (worker_models.py:108-112).
  • ETL/TASK path — workers/file_processing/tasks.py:1073-1081 constructs BatchExecutionResult with no file_results= argument at all.

So errors == {} always, and :281 always returns the fallback. The Moody's incident row would now read "All 1 file(s) failed." — a restatement of status=ERROR + failed_files=1. The user still cannot see why. The blank was replaced by a near-blank.

A repo-wide grep finds no file_results producer emitting "error". The PR's own fixture hand-writes {"status": "error", ...} at workers/tests/test_un3016_execution_error.py:85-93, so test_status_function_returns_a_reason (:97) is green against a shape production cannot produce.

Fix: key on the presence of an error, not on a status string — time_utils.py:180 → if isinstance(file_result, dict) and file_result.get("error"):, keying the name off file_name or file. Separately file_processing/tasks.py:1073 must pass file_results= or the ETL path keeps summarising nothing. Then rebuild the fixture from a real FileExecutionResult(...).to_dict().

Confidence: High — verified independently against PR-head source.


summary = f"All {total_files} file(s) failed. "
shown = seen[:_MAX_ERRORS_IN_SUMMARY]
summary += " | ".join(shown)
remaining = len(seen) - len(shown)
shown = entries[:_MAX_ERRORS_IN_SUMMARY]
summary = f"All {total_files} file(s) failed. " + " | ".join(shown)
remaining = len(entries) - len(shown)
if remaining > 0:
summary += f" | (+{remaining} more)"

Expand Down
103 changes: 65 additions & 38 deletions workers/tests/test_un3016_execution_error.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,36 +2,27 @@

The execution row used to be written with status=ERROR and a blank
error_message, so a failed run gave the user no reason at all. These tests
pin the summary helper that now supplies that reason.
pin the summary helper that now supplies that reason, and assert the reason
survives the call that builds it.

The helper is extracted with `ast` rather than imported, because importing
workers.callback.tasks pulls in celery and the whole worker runtime.
conftest loads .env.test before collection, so importing callback.tasks here
works the same way it does in test_pg_callback_duplicate_guard.py.
"""

import ast
from pathlib import Path
from unittest.mock import MagicMock, patch

import pytest

_TASKS = Path(__file__).resolve().parents[1] / "callback" / "tasks.py"


def _load_helper():
"""Exec just the helper and its constants out of callback/tasks.py."""
tree = ast.parse(_TASKS.read_text())
ns: dict = {"Any": object}
for node in tree.body:
if isinstance(node, ast.Assign) and any(
getattr(t, "id", "").startswith(("_MAX_ERRORS", "_EXECUTION_ERROR"))
for t in node.targets
):
exec(compile(ast.Module([node], []), "<helper>", "exec"), ns)
elif isinstance(node, ast.FunctionDef) and node.name == "_summarize_file_errors":
exec(compile(ast.Module([node], []), "<helper>", "exec"), ns)
return ns["_summarize_file_errors"], ns["_EXECUTION_ERROR_MAX_LENGTH"]


summarize, MAX_LEN = _load_helper()
import callback.tasks as _tasks_module
from callback.tasks import (
_EXECUTION_ERROR_MAX_LENGTH as MAX_LEN,
)
from callback.tasks import (
_determine_execution_status_unified,
)
from callback.tasks import (
_summarize_file_errors as summarize,
)


def test_real_un3016_error_is_surfaced():
Expand Down Expand Up @@ -84,20 +75,48 @@ def test_fits_the_database_column():
assert result.endswith("...")


def _all_failed_batch():
"""One batch, one file, that file errored — the UN-3016 shape."""
return [
{
"total_files": 1,
"successful_files": 0,
"failed_files": 1,
"execution_time": 1.0,
"file_results": [
{
"status": "error",
"file_name": "Villa Bella.xlsm",
"error": "Workflow error: Execution: unstract/api/org_x/e/x.xlsm",
}
],
}
]


def test_status_function_returns_a_reason():
"""_determine_execution_status_unified must return a 4-tuple."""
tree = ast.parse(_TASKS.read_text())
fn = next(
n
for n in ast.walk(tree)
if isinstance(n, ast.FunctionDef)
and n.name == "_determine_execution_status_unified"
)
returns = [n for n in ast.walk(fn) if isinstance(n, ast.Return)]
assert returns, "function must return"
assert all(
isinstance(r.value, ast.Tuple) and len(r.value.elts) == 4 for r in returns
), "every return must carry (results, status, expected_files, error_message)"
"""The real call must hand back a non-blank reason, not just a 4-tuple.

This is what the defect actually was: the tuple gained a fourth slot but
an always-None fourth slot would still leave the execution row blank, so
assert on the value rather than on the shape.
"""
api_client = MagicMock()
with patch(
"callback.tasks.WallClockTimeCalculator.calculate_execution_time",
return_value=1.0,
):
_, final_status, _, error_message = _determine_execution_status_unified(
file_batch_results=_all_failed_batch(),
api_client=api_client,
execution_id="e-1",
organization_id="org-1",
)

assert final_status == "ERROR"
assert error_message, "an ERROR execution must carry a reason (UN-3016)"
assert "Villa Bella.xlsm" in error_message
assert len(error_message) <= MAX_LEN


def test_no_caller_passes_a_hardcoded_none_error():
Expand All @@ -108,8 +127,16 @@ def test_no_caller_passes_a_hardcoded_none_error():
elsewhere in the module cannot trip it. Detects the literal `error_message=None`
keyword only — a positional None, an indirected variable, or a `**kwargs`
splat would pass; all three defect sites were the literal form.

Source-shape rather than behavioural because it guards the *call sites*:
the behavioural cover for the value itself is
test_status_function_returns_a_reason above.
"""
tree = ast.parse(_TASKS.read_text())
import ast
from pathlib import Path

tasks_py = Path(_tasks_module.__file__)
tree = ast.parse(tasks_py.read_text())
callers = [
n
for n in ast.walk(tree)
Expand Down