Skip to content
Open
Show file tree
Hide file tree
Changes from 16 commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
215929d
UN-1924 [FIX] Reject unsupported files by sniffing MIME in API storag…
Deepak-Kesavan Sep 1, 2026
7af4732
UN-1924 [FIX] Terminalise an API execution when every file is rejected
Deepak-Kesavan Sep 1, 2026
63d5237
UN-1924 [FIX] Address review findings on unsupported-file rejection
Deepak-Kesavan Sep 2, 2026
4b436a8
UN-1924 [FIX] Release API rate limit slots by org id, not model instance
Deepak-Kesavan Sep 2, 2026
9929f88
UN-1924 [FIX] Fail one undetectable upload instead of the whole request
Deepak-Kesavan Sep 2, 2026
0efd943
UN-1924 [FIX] Acknowledge results and notify subscribers on an all-re…
Deepak-Kesavan Sep 7, 2026
b27ef4c
UN-1924 [FIX] Keep acknowledgement and notification independent
Deepak-Kesavan Sep 7, 2026
2bf3bb7
UN-1924 [MISC] Drop formatter churn on two untouched assertions
Deepak-Kesavan Sep 7, 2026
cb97bf6
UN-1924 [FIX] Mirror LLMWhisperer's file type gate
Deepak-Kesavan Sep 18, 2026
938a957
Merge remote-tracking branch 'origin/main' into UN-1924-reject-unsupp…
Deepak-Kesavan Sep 18, 2026
6d9f450
UN-1924 [FIX] Recognise a PDF that does not start at offset 0
Deepak-Kesavan Sep 21, 2026
80867ea
Merge remote-tracking branch 'origin/main' into UN-1924-reject-unsupp…
Deepak-Kesavan Sep 21, 2026
5cdce85
UN-1924 [FIX] Harden rejection reporting and the PDF rescue
Deepak-Kesavan Sep 21, 2026
61985aa
UN-1924 [FIX] Require a PDF version header before rescuing a file
Deepak-Kesavan Sep 21, 2026
1b274b5
UN-1924 [FIX] Resolve zip containers, which libmagic will not name fr…
Deepak-Kesavan Sep 21, 2026
048e94f
UN-1924 [REFACTOR] Move the shared MIME gate logic into unstract/core
Deepak-Kesavan Sep 21, 2026
85d3ae5
UN-1924 [FIX] Bound the ODF mimetype read so a zip bomb cannot expand…
Deepak-Kesavan Sep 21, 2026
6c953a9
UN-1924 [FIX] Stop an archive nominating its own file type
Deepak-Kesavan Sep 21, 2026
3b50d5e
UN-1924 [MISC] Leave one throwing call inside the raises block
Deepak-Kesavan Sep 21, 2026
9c8e9ff
UN-1924 [MISC] Pin the zip-carrying-a-PDF-marker case
Deepak-Kesavan Sep 21, 2026
81a9e6a
UN-1924 [FIX] Three small inconsistencies in the new terminal paths
Deepak-Kesavan Sep 21, 2026
f059e50
UN-1924 [FIX] Count and keep files rejected before dispatch
Deepak-Kesavan Sep 21, 2026
2d70dfd
UN-1924 [FIX] Keep two rejections apart when the bytes are identical
Deepak-Kesavan Sep 21, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 79 additions & 2 deletions backend/api_v2/deployment_helper.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@
from utils.constants import Account, CeleryQueue
from utils.local_context import StateStore
from workflow_manager.endpoint_v2.destination import DestinationConnector
from workflow_manager.endpoint_v2.result_cache_utils import ResultCacheUtils
from workflow_manager.endpoint_v2.source import SourceConnector
from workflow_manager.utils.pipeline_utils import PipelineUtils
from workflow_manager.workflow_v2.dto import ExecutionResponse
from workflow_manager.workflow_v2.enums import ExecutionStatus
from workflow_manager.workflow_v2.execution import WorkflowExecutionServiceHelper
Expand Down Expand Up @@ -293,7 +295,14 @@ def execute_workflow(
logger.exception(f"Failed to mark execution {execution_id} as ERROR")

# Async job never started — release the rate limit slot and clean up.
APIDeploymentRateLimiter.release_slot(api.organization, str(execution_id))
# str(...organization_id), NOT the model instance: release_slot formats
# its argument into the Redis key, and acquire_slot built that key from
# str(organization.organization_id). Passing the instance ZREMs a
# non-member — it returns 0 and raises nothing, so the slot silently
# stays held for the full TTL. Same trap as undispatched_sweep.py:245.
APIDeploymentRateLimiter.release_slot(
str(api.organization.organization_id), str(execution_id)
)
DestinationConnector.delete_api_storage_dir(
workflow_id=workflow_id, execution_id=execution_id
)
Expand All @@ -306,6 +315,72 @@ def execute_workflow(
)
).data

# Staging rejected every file, so there is nothing to dispatch. The worker
# short-circuits an empty file set without writing a status back, which
# would strand this execution in PENDING — terminalise it here instead.
if not hash_values_of_files:
Comment thread
Deepak-Kesavan marked this conversation as resolved.
# Isolate the DB write the way the staging-failure path above does, so
# the rate limit slot and staging dir are released even if it raises.
execution = None
try:
execution = WorkflowExecutionServiceHelper.update_execution_completed(
str(execution_id),
total_files=len(file_objs),
failed_files=len(file_objs),
)
except Exception:
logger.exception(f"Failed to mark execution {execution_id} as COMPLETED")
Comment thread
Deepak-Kesavan marked this conversation as resolved.

APIDeploymentRateLimiter.release_slot(
str(api.organization.organization_id), str(execution_id)
)
DestinationConnector.delete_api_storage_dir(
workflow_id=workflow_id, execution_id=execution_id
)
Comment thread
Deepak-Kesavan marked this conversation as resolved.
api_results = ResultCacheUtils.get_api_results(
workflow_id=str(workflow_id), execution_id=str(execution_id)
)
if execution is not None:
# Separate try blocks on purpose: these are independent obligations,
# and sharing one would let a failed acknowledgement silence the
# notification that subscribers depend on.
try:
# This response carries the results, so mark them consumed the way
# the synchronous dispatch path does — otherwise a follow-up
# GET /status serves them a second time with 200 instead of 406.
WorkflowHelper.set_result_acknowledge(execution)
except Exception:
logger.exception(
f"Failed to acknowledge results for execution {execution_id}"
)

try:
# Terminalising here bypasses WorkflowHelper, which is what
# normally notifies API deployment subscribers on a terminal
# status. Without this an all-rejected run alerts nobody.
PipelineUtils.update_pipeline_status(
pipeline_id=pipeline_id, workflow_execution=execution
)
Comment thread
greptile-apps[bot] marked this conversation as resolved.
except Exception:
logger.exception(
f"Failed to notify subscribers for execution {execution_id}"
)

# Report the stored status rather than asserting COMPLETED: the row may
# be missing, or the terminal guard may have refused the change. Claiming
# success here would only hide the stranded execution behind a 200 that a
# follow-up GET /status then contradicts.
return APIExecutionResponseSerializer(
ExecutionResponse(
workflow_id=workflow_id,
execution_id=execution_id,
execution_status=(
execution.status if execution else ExecutionStatus.ERROR.value
),
result=api_results,
)
).data
Comment thread
Deepak-Kesavan marked this conversation as resolved.

try:
result = WorkflowHelper.execute_workflow_async(
workflow_id=workflow_id,
Expand Down Expand Up @@ -352,7 +427,9 @@ def execute_workflow(
# Dispatch failures are marked ERROR internally by execute_workflow_async;
# post-dispatch failures (enrichment/config) must not overwrite a running
# execution's status, so only release the slot and clean up storage here.
APIDeploymentRateLimiter.release_slot(api.organization, str(execution_id))
APIDeploymentRateLimiter.release_slot(
str(api.organization.organization_id), str(execution_id)
)

# Clean up storage
DestinationConnector.delete_api_storage_dir(
Expand Down
193 changes: 193 additions & 0 deletions backend/api_v2/tests/test_deployment_helper.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ def _api() -> MagicMock:
api = MagicMock()
api.workflow.id = "wf-1"
api.id = "pipe-1"
api.organization.organization_id = "org-uuid-1"
return api


Expand Down Expand Up @@ -89,3 +90,195 @@ def test_staging_failure_cleanup_survives_db_marking_error(collaborators) -> Non
# Cleanup still runs even though error-marking raised.
collaborators["APIDeploymentRateLimiter"].release_slot.assert_called_once()
collaborators["DestinationConnector"].delete_api_storage_dir.assert_called_once()


@pytest.fixture
def staging_rejects_everything():
"""Patch execute_workflow's collaborators; staging returns no dispatchable files."""
with mock.patch.multiple(
dh,
WorkflowExecutionServiceHelper=mock.DEFAULT,
SourceConnector=mock.DEFAULT,
DestinationConnector=mock.DEFAULT,
APIDeploymentRateLimiter=mock.DEFAULT,
WorkflowHelper=mock.DEFAULT,
ResultCacheUtils=mock.DEFAULT,
PipelineUtils=mock.DEFAULT,
Tag=mock.DEFAULT,
logger=mock.DEFAULT,
) as mocks:
execution_row = MagicMock()
execution_row.id = "exec-123"
mocks[
"WorkflowExecutionServiceHelper"
].create_workflow_execution.return_value = execution_row
mocks["SourceConnector"].add_input_file_to_api_storage.return_value = {}
mocks["ResultCacheUtils"].get_api_results.return_value = [
{"file": "evil.pdf", "status": "Failed", "error": "unsupported MIME type"}
]
completed_row = MagicMock()
completed_row.status = "COMPLETED"
mocks[
"WorkflowExecutionServiceHelper"
].update_execution_completed.return_value = completed_row
yield mocks


def test_all_files_rejected_completes_without_dispatch(
staging_rejects_everything,
) -> None:
"""A request whose every file is rejected must reach a terminal status.

The worker short-circuits an empty file set without writing a status back, so
dispatching one strands the execution in PENDING and the caller polls forever.
"""
mocks = staging_rejects_everything
# A non-empty upload whose staging result is empty. Passing [] instead would
# leave the branch satisfied by `not file_objs` too, and the original bug -
# dispatching a request whose files were all rejected - would pass this test.
response = dh.DeploymentHelper.execute_workflow(
organization_name="org",
api=_api(),
file_objs=[MagicMock()],
timeout=-1,
)
Comment thread
Deepak-Kesavan marked this conversation as resolved.

# Nothing is dispatched...
mocks["WorkflowHelper"].execute_workflow_async.assert_not_called()
# ...the row is terminalised here instead of being left PENDING, and the
# counters are written so the run does not read back as a clean success...
mocks[
"WorkflowExecutionServiceHelper"
].update_execution_completed.assert_called_once_with(
"exec-123", total_files=1, failed_files=1
)
# ...the slot and staging dir are released. The slot must be released by org
# id string: release_slot formats its argument into the Redis key, so passing
# the model instance removes a non-member and silently holds the slot.
mocks["APIDeploymentRateLimiter"].release_slot.assert_called_once_with(
"org-uuid-1", "exec-123"
)
mocks["DestinationConnector"].delete_api_storage_dir.assert_called_once()
# ...and the caller still sees why each file failed.
assert response["execution_status"] == "COMPLETED"
assert response["result"][0]["file"] == "evil.pdf"
assert response["result"][0]["status"] == "Failed"


def test_all_files_rejected_acknowledges_and_notifies(
staging_rejects_everything,
) -> None:
"""The early return owes the caller what the dispatch path would have done.

It hands back the results in its own response and reaches a terminal status
without going through WorkflowHelper, so both the acknowledgement and the
subscriber notification have to happen here or they happen nowhere.
"""
mocks = staging_rejects_everything
completed_row = mocks[
"WorkflowExecutionServiceHelper"
].update_execution_completed.return_value

dh.DeploymentHelper.execute_workflow(
organization_name="org",
api=_api(),
file_objs=[MagicMock()],
timeout=-1,
)

# Results were served in this response, so a later GET /status must 406
# rather than serve them again.
mocks["WorkflowHelper"].set_result_acknowledge.assert_called_once_with(completed_row)
# PipelineUtils is the only dispatcher of API deployment notifications.
mocks["PipelineUtils"].update_pipeline_status.assert_called_once_with(
pipeline_id="pipe-1", workflow_execution=completed_row
)


def test_notification_survives_a_failed_acknowledgement(
staging_rejects_everything,
) -> None:
"""Acknowledgement and notification are independent obligations.

Sharing one try block would let a failed acknowledgement silence the
notification that subscribers depend on.
"""
mocks = staging_rejects_everything
mocks["WorkflowHelper"].set_result_acknowledge.side_effect = Exception("db is down")

response = dh.DeploymentHelper.execute_workflow(
organization_name="org",
api=_api(),
file_objs=[MagicMock()],
timeout=-1,
)

mocks["PipelineUtils"].update_pipeline_status.assert_called_once()
assert response["execution_status"] == "COMPLETED"


def test_all_files_rejected_still_responds_if_notification_fails(
staging_rejects_everything,
) -> None:
"""A failing webhook must not turn a handled rejection into a 500."""
mocks = staging_rejects_everything
mocks["PipelineUtils"].update_pipeline_status.side_effect = Exception("webhook down")

response = dh.DeploymentHelper.execute_workflow(
organization_name="org",
api=_api(),
file_objs=[MagicMock()],
timeout=-1,
)

assert response["execution_status"] == "COMPLETED"
assert response["result"][0]["file"] == "evil.pdf"


def test_files_staged_successfully_are_dispatched(staging_rejects_everything) -> None:
"""The short-circuit must not fire when staging did return files.

Sibling to the test above: together they pin the branch to the staging result
rather than to the upload list.
"""
mocks = staging_rejects_everything
mocks["SourceConnector"].add_input_file_to_api_storage.return_value = {
"good.pdf": MagicMock()
}

dh.DeploymentHelper.execute_workflow(
organization_name="org",
api=_api(),
file_objs=[MagicMock()],
timeout=-1,
)

mocks["WorkflowHelper"].execute_workflow_async.assert_called_once()
mocks["WorkflowExecutionServiceHelper"].update_execution_completed.assert_not_called()


def test_all_files_rejected_cleanup_survives_db_marking_error(
staging_rejects_everything,
) -> None:
"""A failing status write must not strand the slot or the staging dir.

update_execution_completed only catches DoesNotExist, so a lock timeout or a
dropped connection propagates; without isolation the org's rate limit slot
stays held for its full TTL and throttles every other call for that org.
"""
mocks = staging_rejects_everything
mocks[
"WorkflowExecutionServiceHelper"
].update_execution_completed.side_effect = Exception("db is down")

response = dh.DeploymentHelper.execute_workflow(
organization_name="org",
api=_api(),
file_objs=[MagicMock()],
timeout=-1,
)

mocks["APIDeploymentRateLimiter"].release_slot.assert_called_once()
mocks["DestinationConnector"].delete_api_storage_dir.assert_called_once()
# The row never reached COMPLETED, so the response must not claim it did.
assert response["execution_status"] == "ERROR"
25 changes: 23 additions & 2 deletions backend/workflow_manager/endpoint_v2/enums.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,27 @@
from enum import Enum

from unstract.core.mime_gate import ( # noqa: F401 re-exported
INCONCLUSIVE_MIME_TYPES,
resolve_inconclusive_mime_type,
)


class FileProcessingStatus(Enum):
SUCCESS = "SUCCESS"
ERROR = "ERROR"


class AllowedFileTypes(Enum):
"""MIME types accepted into a workflow.

Mirrors LLMWhisperer's own gate, `Util.is_valid_file_type_from_path`
(unstract-llm-whisperer `backend/app/util/base.py`), member for member.
Accepting what it cannot read only defers the failure to extraction, and
rejecting what it can read loses a file that would have worked - so the two
sets have to move together. Keep in step with
workers/shared/enums/file_types.py.
"""

PLAIN_TEXT = "text/plain"
PDF = "application/pdf"
JPEG = "image/jpeg"
Expand All @@ -19,7 +34,6 @@ class AllowedFileTypes(Enum):
DOC = "application/msword"
XLSX = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"
XLS = "application/vnd.ms-excel"
XLSM = "application/vnd.ms-excel.sheet.macroenabled.12"
PPTX = "application/vnd.openxmlformats-officedocument.presentationml.presentation"
PPT = "application/vnd.ms-powerpoint"
ODT = "application/vnd.oasis.opendocument.text"
Expand All @@ -28,8 +42,15 @@ class AllowedFileTypes(Enum):
CDFV2 = "application/CDFV2"
JSON = "application/json"
CSV = "text/csv"
Comment thread
Deepak-Kesavan marked this conversation as resolved.
OCTET_STREAM = "application/octet-stream"

@classmethod
def is_allowed(cls, mime_type: str) -> bool:
"""Whether LLMWhisperer would accept this type.

Any `text/*` passes, which is how html, xml, tsv, rtf and markdown are
handled - they reach the extractor as text rather than as a listed
format. The enum covers everything else.
"""
if mime_type.startswith("text/"):
return True
return mime_type in cls._value2member_map_
Loading
Loading