Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,9 @@
## 0.1.13

### Fixes

- **Bound gzip decompression memory**: gzip uploads are decompressed incrementally into a spooled temporary file instead of being materialized as a complete in-memory `bytes` value. Expanded outputs larger than 1 MiB spill to temporary disk, and the request still closes the decompressed file when it finishes. Once `unstructured` includes [unstructured#4419](https://github.com/Unstructured-IO/unstructured/pull/4419) (expected in 0.27.12), MIME detection and partitioning read that file in place instead of copying it into memory. No size limit is introduced.

## 0.1.12

### Improvements
Expand Down
2 changes: 1 addition & 1 deletion prepline_general/api/__version__.py
Original file line number Diff line number Diff line change
@@ -1 +1 @@
__version__ = "0.1.12" # pragma: no cover
__version__ = "0.1.13" # pragma: no cover
37 changes: 27 additions & 10 deletions prepline_general/api/general.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,12 @@
import mimetypes
import os
import secrets
import shutil
import tempfile
from base64 import b64encode
from concurrent.futures import ThreadPoolExecutor
from functools import partial
from typing import IO, Any, Dict, List, Mapping, Optional, Sequence, Tuple, Union, cast
from typing import IO, Any, BinaryIO, Dict, List, Mapping, Optional, Sequence, Tuple, Union, cast

import backoff
import pandas as pd
Expand All @@ -31,6 +33,7 @@
from starlette.datastructures import Headers
from starlette.types import Send

from prepline_general.api import __version__ as api_version
from prepline_general.api.filetypes import get_validated_mimetype
from prepline_general.api.models.form_params import GeneralFormParams
from unstructured.documents.elements import Element
Expand All @@ -41,11 +44,13 @@
elements_from_json,
)
from unstructured_inference.models.base import UnknownModelException
from prepline_general.api import __version__ as api_version

app = FastAPI()
router = APIRouter()

_GZIP_COPY_CHUNK_SIZE = 1024 * 1024
_GZIP_SPOOL_MAX_MEMORY_BYTES = 1024 * 1024


def is_compatible_response_type(media_type: str, response_type: type) -> bool:
"""True when `response_type` can be converted to `media_type` for HTTP Response."""
Expand Down Expand Up @@ -613,13 +618,25 @@ def return_content_type(filename: str):
if filename.endswith(".gz"):
filename = filename[:-3]

gzip_file = gzip.open(file.file).read()
return UploadFile(
file=io.BytesIO(gzip_file),
size=len(gzip_file),
filename=filename,
headers=Headers({"content-type": return_content_type(filename)}),
)
output_file = tempfile.SpooledTemporaryFile(max_size=_GZIP_SPOOL_MAX_MEMORY_BYTES)
try:
with gzip.open(file.file) as gzip_file:
shutil.copyfileobj(gzip_file, output_file, length=_GZIP_COPY_CHUNK_SIZE)
Comment thread
CyMule marked this conversation as resolved.
uncompressed_size = output_file.tell()
output_file.seek(0)
except Exception:
output_file.close()
raise

# Reuse the request-owned UploadFile so FastAPI closes the decompressed spool when the
# request finishes. A newly-created UploadFile would not belong to the parsed FormData and
# would therefore remain open after the response.
file.file.close()
file.file = cast(BinaryIO, output_file)
file.size = uncompressed_size
file.filename = filename
file.headers = Headers({"content-type": return_content_type(filename)})
return file


@router.get("/general/v0/general", include_in_schema=False)
Expand Down Expand Up @@ -747,7 +764,7 @@ def join_responses(
frames = [
pd.read_csv(io.BytesIO(response.body)) # pyright: ignore[reportUnknownMemberType]
for response in responses
if response.body.strip()
if cast(bytes, response.body).strip()
]
if not frames:
return PlainTextResponse(responses[0].body)
Expand Down
71 changes: 71 additions & 0 deletions test_general/api/test_gzip.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,13 +9,84 @@
import pandas as pd
import pytest
from deepdiff import DeepDiff
from fastapi import UploadFile
from fastapi.testclient import TestClient

from prepline_general.api import general
from prepline_general.api.app import app
from prepline_general.api.general import _GZIP_SPOOL_MAX_MEMORY_BYTES, ungz_file
from unstructured.partition.common.common import convert_to_bytes

MAIN_API_ROUTE = "general/v0/general"


def _gzip_upload(content: bytes, filename: str = "sample.txt.gz") -> UploadFile:
compressed = io.BytesIO(gzip.compress(content))
return UploadFile(file=compressed, filename=filename)


@pytest.mark.parametrize(
("content", "spills_to_disk"),
[
pytest.param(b"small gzip payload", False, id="small-stays-in-memory"),
pytest.param(b"x" * (_GZIP_SPOOL_MAX_MEMORY_BYTES + 1), True, id="large-spills-to-disk"),
],
)
def test_ungz_file_bounds_decompression_memory(content: bytes, spills_to_disk: bool):
upload = _gzip_upload(content)
compressed_file = upload.file

result = ungz_file(upload)

try:
assert result is upload
assert compressed_file.closed
assert result.filename == "sample.txt"
assert result.content_type == "text/plain"
assert result.size == len(content)
assert result.file.read() == content
assert result.file._rolled is spills_to_disk
assert isinstance(result.file, tempfile.SpooledTemporaryFile)
finally:
result.file.close()
assert result.file.closed


def test_ungz_file_output_that_spills_to_disk_is_readable_by_unstructured():
"""PDF OCR and text-encoding detection read the whole upload through `convert_to_bytes()`."""
content = b"x" * (_GZIP_SPOOL_MAX_MEMORY_BYTES + 1)
result = ungz_file(_gzip_upload(content, filename="sample.pdf.gz"))

try:
assert result.file._rolled
assert convert_to_bytes(result.file) == content
assert result.file.tell() == 0
finally:
result.file.close()


def test_gzipped_upload_is_partitioned_and_closed_after_the_request(monkeypatch):
content = b"Gzipped uploads are partitioned like their uncompressed content."
decompressed_uploads: list[UploadFile] = []

def capture_ungz_file(*args, **kwargs):
upload = ungz_file(*args, **kwargs)
decompressed_uploads.append(upload)
return upload

monkeypatch.setattr(general, "ungz_file", capture_ungz_file)

response = TestClient(app).post(
MAIN_API_ROUTE,
files=[("files", ("sample.txt.gz", gzip.compress(content), "application/gzip"))],
)

assert response.status_code == 200
assert [element["text"] for element in response.json()] == [content.decode()]
assert len(decompressed_uploads) == 1
assert decompressed_uploads[0].file.closed


@pytest.mark.xfail(reason="The outputs are different as of unstructured==0.13.5")
@pytest.mark.parametrize("output_format", ["application/json", "text/csv"])
@pytest.mark.parametrize(
Expand Down