diff --git a/CHANGELOG.md b/CHANGELOG.md index f289a01d1d..00de0e49f2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,9 @@ +## 0.27.15 + +### Fixes + +- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. A lone `"\r"` line ending is also converted to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read one followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. `"\r\n"` is left as is. The file is streamed to Pandas rather than read into memory whole, and `partition_tsv()` still decompresses a compressed filename (e.g. `.tsv.gz`) and accepts a stream that cannot seek. When the delimiter is sniffed, it is now sniffed once, from the first non-blank line, and passed to Pandas, and a file with no usable delimiter is read as one column. + ## 0.27.14 ### Fixes diff --git a/test_unstructured/partition/test_csv.py b/test_unstructured/partition/test_csv.py index b648134d08..15c4c6b6d4 100644 --- a/test_unstructured/partition/test_csv.py +++ b/test_unstructured/partition/test_csv.py @@ -3,7 +3,9 @@ from __future__ import annotations import io +from pathlib import Path +import pandas as pd import pytest from pytest_mock import MockFixture @@ -28,7 +30,15 @@ from unstructured.chunking.title import chunk_by_title from unstructured.cleaners.core import clean_extra_whitespace from unstructured.documents.elements import Table -from unstructured.partition.csv import _CsvPartitioningContext, partition_csv +from unstructured.errors import UnprocessableEntityError +from unstructured.partition import csv as csv_module +from unstructured.partition.csv import ( + _CsvPartitioningContext, + _first_record_width, + check_cell_count, + partition_csv, + read_delimited_text, +) from unstructured.partition.utils.constants import UNSTRUCTURED_INCLUDE_DEBUG_METADATA EXPECTED_FILETYPE = "text/csv" @@ -211,6 +221,144 @@ def test_partition_csv_header(): assert table.metadata.text_as_html is not None +# -- cell-count limit ---------------------------------------------------------------------------- + + +def test_partition_csv_rejects_a_wide_first_line_before_pandas_reads_it( + tmp_path: Path, mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.delenv("CSV_MAX_CELLS", raising=False) # -- default limit of 5M cells -- + # -- Pandas pads every row out to the first line's 5,000 fields: 25M cells from 15KB -- + file_path = tmp_path / "ragged.csv" + file_path.write_text("h" + "," * 4999 + "\n" + "a\n" * 5000) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="rows x 5,000 columns"): + partition_csv(str(file_path)) + + read_csv_.assert_not_called() + + +@pytest.mark.parametrize("from_file", [False, True]) +@pytest.mark.parametrize( + ("content", "n_measured_cells"), + [ + # -- the context's restricted sniffer gives up on this one, so the delimiter is sniffed + # -- from the first line without restricting the candidates -- + ("h" + "," * 99 + "\n" + "a\n" * 99, 100 * 100), + # -- every line ending counts as a row, so the blank line makes this 4 rows (Pandas + # -- reads 3): the measure is an upper bound -- + ("a;b;c\n1;2\n\n4\n", 3 * 4), + # -- single-column file; the sniffer finds "a" as the delimiter, giving 2 columns -- + ("a\nb\nc\nd\n", 4 * 2), + ('"x,\ny",z\n1\n', 2 * 2), # -- quoted delimiter and newline start no field or row -- + ("a,b\r\nc,d\re,f", 3 * 2), # -- "\r\n", "\r" and an unterminated last line -- + ], +) +def test_partition_csv_limits_the_cells_the_file_spans( + content: str, + n_measured_cells: int, + from_file: bool, + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +): + file_path = tmp_path / "table.csv" + file_path.write_bytes(content.encode()) + + def partition(): + if from_file: + with open(file_path, "rb") as f: + return partition_csv(file=f) + return partition_csv(str(file_path)) + + monkeypatch.setenv("CSV_MAX_CELLS", str(n_measured_cells)) + assert len(partition()) == 1 + + monkeypatch.setenv("CSV_MAX_CELLS", str(n_measured_cells - 1)) + with pytest.raises(UnprocessableEntityError, match="CSV_MAX_CELLS"): + partition() + + +@pytest.mark.parametrize( + "first_line", + [ + " \t\n", # -- a whitespace-only line is skipped, so the wide line is the first record -- + '"x\n"', # -- a quoted newline does not end the first record -- + "\x0c", # -- form-feed is not a line ending, so the commas are on the first line -- + 'a"b', # -- a quote inside a field does not open a quoted field -- + "\ufeff\n", # -- a byte-order mark is dropped, leaving a blank line that is skipped -- + ], +) +def test_check_cell_count_measures_the_record_pandas_sizes_the_data_frame_by( + first_line: str, monkeypatch: pytest.MonkeyPatch +): + # -- Pandas reads 4 rows x 101 columns from each of these -- + file = io.BytesIO((first_line + "," * 100 + "\n" + "a,b\n" * 3).encode()) + monkeypatch.setenv("CSV_MAX_CELLS", "403") + + with pytest.raises(UnprocessableEntityError, match="101 columns"): + check_cell_count(file, ",", None) + + +def test_partition_csv_is_not_tricked_into_millions_of_rows_by_a_carriage_return(): + # -- Pandas 2.x's C tokenizer reads this as 262,145 rows -- + elements = partition_csv(file=io.BytesIO(b"a,b\n\r ,c\n")) + + assert elements[0].metadata.text_as_html == ( + "
ab
c
" + ) + + +@pytest.mark.parametrize(("python_engine", "expected_width"), [(True, 11), (False, 1)]) +def test_first_record_width_applies_each_pandas_engines_blank_line_rule( + python_engine: bool, expected_width: int +): + # -- the Python engine skips a record whose one value is whitespace, even quoted; the C engine + # -- skips only an unquoted line of spaces and tabs -- + chunks = iter(['" "\n' + ";" * 10 + "\n"]) + + width, _ = _first_record_width(chunks, ";", python_engine, 10**9, Mock()) + + assert width == expected_width + + +def test_partition_csv_sniffs_the_delimiter_from_the_first_non_blank_line(): + # -- the context's sniffer only tries ",;|", so Pandas' delimiter is sniffed here -- + elements = partition_csv(file=io.BytesIO(b"\n\na\tb\tc\n1\t2\t3\n")) + + assert elements[0].metadata.text_as_html == ( + "" + "
abc
123
" + ) + + +def test_read_delimited_text_streams_the_file_in_chunks(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1) + read_sizes: list[int] = [] + + class RecordingBytesIO(io.BytesIO): + def read(self, size: int | None = -1) -> bytes: + read_sizes.append(-1 if size is None else size) + return super().read(size) + + # -- with 1-byte chunks every "\r" is held back to see whether a "\n" follows it -- + dataframe = read_delimited_text( + RecordingBytesIO(b'a,"x\r\ny"\r\nc,"d\re"\r'), sep=",", header=None, encoding=None + ) + + assert dataframe.values.tolist() == [["a", "x\r\ny"], ["c", "d\ne"]] + assert set(read_sizes) == {1} + + +def test_partition_csv_reads_a_file_with_no_usable_delimiter_as_one_column(): + # -- the sniffer picks the quote as the delimiter, which cannot delimit fields -- + elements = partition_csv(file=io.BytesIO(b'"a"\n"b"\n')) + + assert ( + elements[0].metadata.text_as_html == "
a
b
" + ) + + # ================================================================================================ # UNIT-TESTS # ================================================================================================ @@ -327,3 +475,143 @@ def it_raises_when_neither_file_path_nor_file_is_provided(self): @pytest.fixture() def get_last_modified_date_(self, request: FixtureRequest) -> Mock: return function_mock(request, "unstructured.partition.csv.get_last_modified_date") + + +# -- streaming and boundaries -------------------------------------------------------------------- + + +def test_partition_csv_reads_past_a_long_blank_prefix_in_linear_time( + monkeypatch: pytest.MonkeyPatch, +): + # -- no delimiter in the context's sample, so the delimiter is sniffed past the blank lines and + # -- every blank line is replayed to Pandas' Python engine -- + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1024) + data = b"\n" * 200_000 + b"a\tb\n1\t2\n" + + elements = partition_csv(file=io.BytesIO(data)) + + assert elements[0].metadata.text_as_html == ( + "
ab
12
" + ) + + +def test_lone_carriage_return_reader_copies_each_character_a_bounded_number_of_times( + monkeypatch: pytest.MonkeyPatch, +): + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1024) + data = b"\n" * 200_000 + b"a\tb\n1\t2\n" + reader = csv_module._LoneCarriageReturnReader(io.BytesIO(data), None) + + assert reader.peek_first_non_blank_line() == "a\tb\n" + lines = list(iter(reader.readline, "")) + + assert len(lines) == 200_002 + assert "".join(lines) == data.decode() + assert reader._chars_copied <= 3 * len(data) + + +def test_peek_first_non_blank_line_replays_the_original_chunks(): + chunks = ["\n" * 10, " \t\n", "a,", "b\n", "c,d\n"] + + line, replay = csv_module._peek_first_non_blank_line(iter(chunks)) + + assert line == "a,b\n" + assert all(a is b for a, b in zip(replay, chunks, strict=True)) + + +def test_partition_csv_counts_the_implicit_index_column_of_a_header( + monkeypatch: pytest.MonkeyPatch, +): + # -- data rows one field wider than the header: Pandas uses the extra field as the index -- + data = b"a,b\n1,2,3\n4,5,6\n" + monkeypatch.setenv("CSV_MAX_CELLS", str(3 * 3)) + assert len(partition_csv(file=io.BytesIO(data), include_header=True)) == 1 + + monkeypatch.setenv("CSV_MAX_CELLS", str(3 * 3 - 1)) + with pytest.raises(UnprocessableEntityError, match="3 columns"): + partition_csv(file=io.BytesIO(data), include_header=True) + + +@pytest.mark.parametrize( + "data", + [ + 'é,ü\r\nx,"a\r\nb"\r\n', # -- BOM, multibyte, CRLF inside and outside quotes -- + '日本,語\r1,"2\r3"\r\r\n', # -- lone "\r" inside and outside quotes -- + '\n"q""uote","é"\n', # -- repeated BOM, doubled quote -- + ], +) +@pytest.mark.parametrize("chunk_size", [1, 3]) +def test_reading_does_not_depend_on_where_chunks_split_the_file( + data: str, chunk_size: int, monkeypatch: pytest.MonkeyPatch +): + encoded = data.encode() + expected_text = csv_module._LoneCarriageReturnReader(io.BytesIO(encoded), None).read() + expected_frame = read_delimited_text(io.BytesIO(encoded), sep=",", header=None, encoding=None) + expected_cells = expected_frame.shape[0] * expected_frame.shape[1] + + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", chunk_size) + + assert csv_module._LoneCarriageReturnReader(io.BytesIO(encoded), None).read() == expected_text + frame = read_delimited_text(io.BytesIO(encoded), sep=",", header=None, encoding=None) + assert frame.equals(expected_frame) + monkeypatch.setenv("CSV_MAX_CELLS", str(expected_cells - 1)) + with pytest.raises(UnprocessableEntityError): + check_cell_count(io.BytesIO(encoded), ",", None) + + +@pytest.mark.parametrize( + "data", + [ + b"h" + + b"," * 100 + + b"\n" + + b"a\n" * 100_000, # -- the first record alone passes the limit -- + b"a,b\n" * 100_000, # -- the rows pass the limit -- + ], +) +def test_check_cell_count_stops_reading_once_the_limit_is_passed( + data: bytes, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 64) + monkeypatch.setenv("CSV_MAX_CELLS", "50") + n_bytes_read = 0 + + class CountingBytesIO(io.BytesIO): + def read(self, size: int | None = -1) -> bytes: + nonlocal n_bytes_read + chunk = super().read(size) + n_bytes_read += len(chunk) + return chunk + + with pytest.raises(UnprocessableEntityError): + check_cell_count(CountingBytesIO(data), ",", None) + + assert n_bytes_read <= 4 * 64 + + +def test_peek_first_non_blank_line_completes_a_line_across_a_chunk_of_only_a_held_carriage_return( + monkeypatch: pytest.MonkeyPatch, +): + # -- with 1-byte chunks the "\r" chunk queues nothing until the "\n" after it arrives -- + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1) + reader = csv_module._LoneCarriageReturnReader(io.BytesIO(b"a\r\nb\n"), None) + + assert reader.peek_first_non_blank_line() == "a\r\n" + assert reader.read() == "a\r\nb\n" + + +def test_partition_csv_counts_index_names_the_python_engine_infers( + mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + # -- the context's sniffer only tries ",;|", so this tab-delimited file is read by Pandas' + # -- Python engine. Its second data row is as wide as the first plus the header, so Pandas + # -- makes the first data row the index names: 3 columns, from a 2-field header -- + data = b"a\tb\nidx\ni\t1\t2\n" + b"j\t3\t4\n" * 10 + pandas_cells = 13 * 3 + monkeypatch.setenv("CSV_MAX_CELLS", str(pandas_cells - 1)) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="3 columns"): + partition_csv(file=io.BytesIO(data), include_header=True) + + read_csv_.assert_not_called() diff --git a/test_unstructured/partition/test_tsv.py b/test_unstructured/partition/test_tsv.py index 5276e6006d..59bf820f2a 100644 --- a/test_unstructured/partition/test_tsv.py +++ b/test_unstructured/partition/test_tsv.py @@ -2,6 +2,12 @@ from __future__ import annotations +import gzip +import io +from pathlib import Path +from typing import Any + +import pandas as pd import pytest from pytest_mock import MockFixture @@ -14,7 +20,9 @@ ) from test_unstructured.unit_utils import assert_round_trips_through_JSON, example_doc_path from unstructured.chunking.title import chunk_by_title +from unstructured.common.html_table import HtmlTable from unstructured.documents.elements import Table +from unstructured.errors import UnprocessableEntityError from unstructured.partition.tsv import partition_tsv EXPECTED_FILETYPE = "text/tsv" @@ -159,3 +167,122 @@ def test_partition_tsv_supports_chunking_strategy_while_partitioning(): # The same chunks are returned if chunking elements or chunking during partitioning. assert chunk_elements == chunks + + +# -- cell-count limit ---------------------------------------------------------------------------- + + +@pytest.mark.parametrize("from_file", [False, True]) +def test_partition_tsv_rejects_a_wide_first_line_before_pandas_reads_it( + from_file: bool, tmp_path: Path, mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.delenv("CSV_MAX_CELLS", raising=False) # -- default limit of 5M cells -- + # -- Pandas pads every row out to the first line's 5,000 fields: 25M cells from 15KB -- + file_path = tmp_path / "ragged.tsv" + file_path.write_text("h" + "\t" * 4999 + "\n" + "a\n" * 5000) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="rows x 5,000 columns"): + if from_file: + with open(file_path, "rb") as f: + partition_tsv(file=f) + else: + partition_tsv(str(file_path)) + + read_csv_.assert_not_called() + + +@pytest.mark.parametrize("from_file", [False, True]) +def test_partition_tsv_partitions_a_file_at_the_cell_limit( + from_file: bool, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.setenv("CSV_MAX_CELLS", "6") + file_path = tmp_path / "table.tsv" + file_path.write_text("a\tb\tc\n1\t2\n") + + if from_file: + with open(file_path, "rb") as f: + elements = partition_tsv(file=f) + else: + elements = partition_tsv(str(file_path)) + + assert [e.text for e in elements] == ["a b c 1 2"] + + +def test_partition_tsv_reads_a_field_larger_than_the_csv_module_field_limit(tmp_path: Path): + # -- the `csv` module's default field limit is 128 KiB; Pandas has none -- + file_path = tmp_path / "big-field.tsv" + file_path.write_text("a\t" + "x" * 200_000 + "\n") + + (table,) = partition_tsv(str(file_path)) + + assert table.text == "a " + "x" * 200_000 + + +def test_partition_tsv_decompresses_a_compressed_filename( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + file_path = tmp_path / "table.tsv.gz" + with gzip.open(file_path, "wb") as f: + f.write(b"a\tb\n1\t2\n") + + (table,) = partition_tsv(str(file_path)) + + assert table.text == "a b 1 2" + + # -- the size check measures the decompressed content -- + monkeypatch.setenv("CSV_MAX_CELLS", "3") + with pytest.raises(UnprocessableEntityError, match="CSV_MAX_CELLS"): + partition_tsv(str(file_path)) + + +def test_partition_tsv_reads_a_stream_that_cannot_seek(): + class Pipe(io.RawIOBase): + def __init__(self, data: bytes): + self._data = io.BytesIO(data) + + def readable(self) -> bool: + return True + + def seekable(self) -> bool: + return False + + def readinto(self, buffer: Any) -> int: + chunk = self._data.read(len(buffer)) + buffer[: len(chunk)] = chunk + return len(chunk) + + (table,) = partition_tsv(file=io.BufferedReader(Pipe(b"a\tb\n1\t2\n"))) + + assert table.text == "a b 1 2" + + +def test_partition_tsv_counts_every_implicit_index_column_before_pandas_reads( + tmp_path: Path, mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + # -- a 1-field header, a 100-field first data row (99 implicit-index columns) and a ragged + # -- tail: Pandas builds 1,002 rows x 100 columns -- + file_path = tmp_path / "index.tsv" + file_path.write_text("h\n" + "\t".join(f"v{i}" for i in range(100)) + "\n" + "a\n" * 1000) + monkeypatch.setenv("CSV_MAX_CELLS", str(1002 * 100 - 1)) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="100 columns"): + partition_tsv(str(file_path), include_header=True) + + read_csv_.assert_not_called() + + +def test_partition_tsv_with_implicit_index_columns_matches_pandas_within_the_limit( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.delenv("CSV_MAX_CELLS", raising=False) + file_path = tmp_path / "index.tsv" + file_path.write_text("h\n" + "x\ty\tz\n" + "a\n" * 3) + + (table,) = partition_tsv(str(file_path), include_header=True) + + expected = pd.read_csv(file_path, sep="\t", header=0).to_html( + index=False, header=True, na_rep="" + ) + assert table.text == HtmlTable.from_html_text(expected).text diff --git a/unstructured/__version__.py b/unstructured/__version__.py index 05d19bf1a2..106c6581b5 100644 --- a/unstructured/__version__.py +++ b/unstructured/__version__.py @@ -1 +1 @@ -__version__ = "0.27.14" # pragma: no cover +__version__ = "0.27.15" # pragma: no cover diff --git a/unstructured/partition/csv.py b/unstructured/partition/csv.py index 85c847d948..3c9f772167 100644 --- a/unstructured/partition/csv.py +++ b/unstructured/partition/csv.py @@ -1,17 +1,24 @@ from __future__ import annotations +import codecs +import collections import contextlib import csv +import io +import itertools +import re from functools import cached_property -from typing import IO, Any, Iterator +from typing import IO, Any, Callable, Iterator, NoReturn import pandas as pd from unstructured.chunking import add_chunking_strategy from unstructured.common.html_table import HtmlTable from unstructured.documents.elements import Element, ElementMetadata, Table +from unstructured.errors import UnprocessableEntityError from unstructured.file_utils.model import FileType from unstructured.partition.common.metadata import apply_metadata, get_last_modified_date +from unstructured.partition.utils.config import env_config from unstructured.telemetry import partition_runtime_telemetry from unstructured.utils import is_temp_file_path @@ -60,11 +67,11 @@ def partition_csv( csv.field_size_limit(CSV_FIELD_LIMIT) with ctx.open() as file: - read_kw: dict = {"header": ctx.header, "sep": ctx.delimiter, "encoding": ctx.encoding} - # sep=None is not supported by the C engine; use Python engine to avoid ParserWarning. - if ctx.delimiter is None: - read_kw["engine"] = "python" - dataframe = pd.read_csv(file, **read_kw) + check_cell_count(file, ctx.delimiter, ctx.encoding, header=ctx.header is not None) + with ctx.open() as file: + dataframe = read_delimited_text( + file, sep=ctx.delimiter, header=ctx.header, encoding=ctx.encoding + ) html_table = HtmlTable.from_html_text( dataframe.to_html(index=False, header=include_header, na_rep="") @@ -80,6 +87,403 @@ def partition_csv( return [Table(text=html_table.text, metadata=metadata, detection_origin=DETECTION_ORIGIN)] +def read_delimited_text( + file: IO[bytes], *, sep: str | None, header: int | None, encoding: str | None +) -> pd.DataFrame: + """Read delimited text from `file` into a data-frame, with each lone "\r" made a "\n". + + Pandas 2.x's C tokenizer mishandles a lone "\r" line ending followed by a whitespace-only line + while skipping blank lines: it emits 2^18 empty rows for each, so a few bytes become millions + of rows inside `pd.read_csv()`, before any limit can be checked. Pandas therefore reads the + file through `_LoneCarriageReturnReader`, which streams it with each "\r" not followed by "\n" + converted to "\n". "\r\n" is left alone, so text inside quoted fields of a "\r\n" file is + unchanged. + + When `sep` is `None` the delimiter is sniffed by `_sniff_delimiter()` from the first non-blank + line, the same way `check_cell_count()` measures it, rather than leaving Pandas to sniff its + own. + """ + reader = _LoneCarriageReturnReader(file, encoding) + if sep is None: + sep = _sniff_delimiter(_normalize_line_endings(reader.peek_first_non_blank_line())) + # -- the C engine does not sniff, so keep the Python engine this path has always used -- + return pd.read_csv(reader, sep=sep, header=header, engine="python") + return pd.read_csv(reader, sep=sep, header=header) + + +def check_cell_count( + file: IO[bytes], delimiter: str | None, encoding: str | None, *, header: bool = False +) -> None: + """Raise `UnprocessableEntityError` when `file` would span more than `CSV_MAX_CELLS` cells. + + Pandas sizes the data-frame by the first record and pads every shorter record out to that + width, so a tiny file whose first line is a long run of delimiters can span millions of cells. + The span is measured here before Pandas allocates anything, reading `file` in fixed-size chunks + from its current position, with each line ending counted as one "\n", and the scan stops as + soon as the limit is passed: + + - columns are the fields of the first record, counted with Pandas' quoting rules without + building the fields; + - rows are counted as line terminators after it. That is an upper bound, since blank lines and + newlines inside quoted fields count too, so the scan can over-count but never under-count. + + With `header`, the first record is the header, and Pandas makes the fields a data row has + beyond it an implicit index: its C engine takes their number from the first data row, and its + Python engine also from the second, which can make the first data row the index names. So the + columns are the widest of the header and the first two data records, each measured the same + way; any later, wider row is an error in Pandas. + + This is a conservative estimate of the cells Pandas will build, not a bound on memory or bytes + read: blank lines held while sniffing, or a stream spooled because it cannot seek, still take + space in proportion to the input. + """ + max_cells = env_config.CSV_MAX_CELLS + chunks = (_normalize_line_endings(c) for c in _iter_decoded_chunks(file, encoding)) + + def raise_limit_exceeded(n_rows: int, n_cols: int) -> NoReturn: + raise UnprocessableEntityError( + f"File exceeds the maximum of {max_cells:,} table cells (CSV_MAX_CELLS): it spans" + f" at least {n_rows:,} rows x {n_cols:,} columns." + ) + + # -- with no delimiter given, `read_delimited_text()` sniffs one and uses Pandas' Python + # -- engine, which skips more kinds of blank line -- + python_engine = delimiter is None + if delimiter is None: + first_line, chunks = _peek_first_non_blank_line(chunks) + delimiter = _sniff_delimiter(first_line) + + n_cols, chunks = _first_record_width( + chunks, delimiter, python_engine, max_cells, raise_limit_exceeded + ) + if n_cols == 0: + return + n_rows = 1 + + if header: + # -- the first two data records decide how many implicit-index columns Pandas adds -- + for _ in range(2): + width, chunks = _first_record_width( + chunks, delimiter, python_engine, max_cells, raise_limit_exceeded + ) + if width == 0: + return + n_rows, n_cols = n_rows + 1, max(n_cols, width) + if n_rows * n_cols > max_cells: + raise_limit_exceeded(n_rows, n_cols) + + # -- each later "\n" ends a row, and so does end-of-file when the last line is unterminated -- + unterminated = False + for chunk in chunks: + if n_terminators := chunk.count("\n"): + n_rows += n_terminators + unterminated = not chunk.endswith("\n") + elif chunk: + unterminated = True + if (n_rows + unterminated) * n_cols > max_cells: + raise_limit_exceeded(n_rows + unterminated, n_cols) + + +# -- the delimiter of a file with no usable one, which is then read as one column. It is the ASCII +# -- "unit separator", so a file that does contain it is split on it, the same way by both +# -- `check_cell_count()` and Pandas -- +_ONE_COLUMN_DELIMITER = "\x1f" + +_CSV_CHUNK_CHARS = 1 << 16 + + +class _LoneCarriageReturnReader(io.TextIOBase): + """Text stream of `file`, decoded in chunks, with each "\r" not followed by "\n" made a "\n". + + A "\r" ending a chunk is held back until the next chunk shows whether it starts a "\r\n". The + file is never held in memory whole: decoded chunks wait in a queue until they are read, and + only the lines `.peek_first_non_blank_line()` looks past stay queued after it returns. + + Each character is copied a bounded number of times however the text is read, so reading is + linear in the size of the file; `._chars_copied` counts the copies. + """ + + def __init__(self, file: IO[bytes], encoding: str | None): + self._file = file + self._decoder = codecs.getincrementaldecoder(_python_encoding(encoding))() + self._chunks: collections.deque[str] = collections.deque() + self._offset = 0 # -- read position in the first queued chunk -- + self._held_cr = False + self._at_start = True + self._eof = False + self._chars_copied = 0 + + def peek_first_non_blank_line(self) -> str: + """The first line that is not blank, or "" if there is none, without consuming it.""" + line_parts: list[str] = [] # -- the current line, kept only until it proves blank -- + idx, start = 0, self._offset + while True: + if idx == len(self._chunks): + if not self._fill(): + line = self._join(line_parts) + return line if line.strip() else "" + continue + chunk = self._chunks[idx] + end = chunk.find("\n", start) + 1 or len(chunk) + line_parts.append(self._slice(chunk, start, end)) + if not line_parts[-1].isspace() and line_parts[-1]: + # -- not blank; complete the line, then return it -- + while not line_parts[-1].endswith("\n"): + idx, start = idx + 1, 0 + # -- a fill can queue nothing (e.g. a chunk holding only a held-back "\r") -- + while idx == len(self._chunks): + if not self._fill(): + return self._join(line_parts) + chunk = self._chunks[idx] + end = chunk.find("\n") + 1 or len(chunk) + line_parts.append(self._slice(chunk, 0, end)) + return self._join(line_parts) + if line_parts[-1].endswith("\n"): + line_parts = [] + if end == len(chunk): + idx, start = idx + 1, 0 + else: + start = end + + def read(self, size: int | None = -1) -> str: + parts: list[str] = [] + remaining = -1 if size is None or size < 0 else size + while remaining != 0 and (self._chunks or self._fill()): + if not self._chunks: + continue + chunk = self._chunks[0] + end = len(chunk) if remaining < 0 else min(len(chunk), self._offset + remaining) + parts.append(self._slice(chunk, self._offset, end)) + if remaining > 0: + remaining -= end - self._offset + self._advance(end) + return self._join(parts) + + def readline(self, size: int | None = -1) -> str: + parts: list[str] = [] + remaining = -1 if size is None or size < 0 else size + while remaining != 0 and (self._chunks or self._fill()): + if not self._chunks: + continue + chunk = self._chunks[0] + end = chunk.find("\n", self._offset) + 1 or len(chunk) + if remaining > 0: + end = min(end, self._offset + remaining) + remaining -= end - self._offset + parts.append(self._slice(chunk, self._offset, end)) + self._advance(end) + if parts[-1].endswith("\n"): + break + return self._join(parts) + + def readable(self) -> bool: + return True + + def _advance(self, end: int) -> None: + """Move the read position to `end` in the first queued chunk, dropping it once read.""" + if end == len(self._chunks[0]): + self._chunks.popleft() + self._offset = 0 + else: + self._offset = end + + def _fill(self) -> bool: + """Queue the next decoded chunk; False when the file is exhausted.""" + if self._eof: + return False + chunk = self._file.read(_CSV_CHUNK_CHARS) + text = self._decoder.decode(chunk, final=not chunk) + if self._at_start: + text, self._at_start = _strip_leading_bom(text) + if self._held_cr: + text = "\r" + text + self._held_cr = False + if not chunk: + self._eof = True + elif text.endswith("\r"): + text, self._held_cr = text[:-1], True + if text: + self._chunks.append(_LONE_CARRIAGE_RETURN.sub("\n", text)) + return True + + def _join(self, parts: list[str]) -> str: + if len(parts) == 1: + return parts[0] + self._chars_copied += sum(map(len, parts)) + return "".join(parts) + + def _slice(self, chunk: str, start: int, end: int) -> str: + if start == 0 and end == len(chunk): + return chunk + self._chars_copied += end - start + return chunk[start:end] + + +_LONE_CARRIAGE_RETURN = re.compile("\r(?!\n)") + + +def _python_encoding(encoding: str | None) -> str: + """The codec to decode with: `encoding`, or UTF-8, dropping a UTF-8 byte-order mark as Pandas + does.""" + encoding = encoding or "utf-8" + return "utf-8-sig" if codecs.lookup(encoding).name == "utf-8" else encoding + + +def _strip_leading_bom(text: str) -> tuple[str, bool]: + """`text` without leading byte-order marks, and whether the start of the file is still ahead. + + Pandas drops a byte-order mark from the start of the text it is given, so every leading one is + dropped before Pandas or the scan sees the text, leaving none for Pandas to drop differently. + """ + text = text.lstrip("\ufeff") + return text, not text + + +def _normalize_line_endings(text: str) -> str: + """`text` with each "\r\n" and "\r" line ending replaced by "\n". + + Applied chunk by chunk, a "\r\n" split across chunks becomes two line endings, which only adds + a blank line. + """ + return text.replace("\r\n", "\n").replace("\r", "\n") + + +def _sniff_delimiter(first_line: str) -> str: + """The delimiter `csv.Sniffer` finds in `first_line`, or `_ONE_COLUMN_DELIMITER`.""" + try: + delimiter = csv.Sniffer().sniff(first_line).delimiter + except csv.Error: + return _ONE_COLUMN_DELIMITER + # -- a quote or line ending cannot delimit fields -- + return _ONE_COLUMN_DELIMITER if delimiter in ('"', "\n", "\r") else delimiter + + +def _iter_decoded_chunks(file: IO[bytes], encoding: str | None) -> Iterator[str]: + """Generate the text of `file` in chunks of up to `_CSV_CHUNK_CHARS` bytes, decoded.""" + # -- `errors="replace"` so the scan never fails on a file Pandas would read -- + decoder = codecs.getincrementaldecoder(_python_encoding(encoding))(errors="replace") + at_start = True + while chunk := file.read(_CSV_CHUNK_CHARS): + text = decoder.decode(chunk) + if at_start: + text, at_start = _strip_leading_bom(text) + if text: + yield text + if text := decoder.decode(b"", final=True): + yield text.lstrip("\ufeff") if at_start else text + + +def _peek_first_non_blank_line(chunks: Iterator[str]) -> tuple[str, Iterator[str]]: + """The first line of `chunks` that is not blank, and all of `chunks`, that line included. + + The line ends with its "\n"; it is "" when every line is blank. Each chunk read is queued + unchanged for the returned chunks, and only the non-blank line itself is copied. + """ + read: list[str] = [] + line_parts: list[str] = [] # -- the current line, kept only until it proves blank -- + for chunk in chunks: + read.append(chunk) + start = 0 + while start < len(chunk): + end = chunk.find("\n", start) + 1 or len(chunk) + segment = chunk[start:end] + line_parts.append(segment) + if segment and not segment.isspace(): + if not segment.endswith("\n"): + # -- complete the line from the chunks that follow -- + for more in chunks: + read.append(more) + end = more.find("\n") + 1 or len(more) + line_parts.append(more[:end]) + if line_parts[-1].endswith("\n"): + break + return "".join(line_parts), itertools.chain(read, chunks) + if segment.endswith("\n"): + line_parts = [] + start = end + line = "".join(line_parts) + return (line if line.strip() else ""), iter(read) + + +def _first_record_width( + chunks: Iterator[str], + delimiter: str, + python_engine: bool, + max_cells: int, + raise_limit_exceeded: Callable[[int, int], NoReturn], +) -> tuple[int, Iterator[str]]: + """The field count of the first non-blank record, and the chunks that follow it. + + Follows the quoting rules Pandas applies by default: a `"` opens a quoted field only at the + start of a field, `""` inside a quoted field is a literal quote, and delimiters and newlines + inside a quoted field are part of the field. A record with no delimiter is blank when Pandas + skips it: for its C engine when it has no quote and only spaces and tabs; for its Python engine + when its one field's value, quoted or not, is empty or only `str.isspace()` characters. The + count is 0 when there is no non-blank record. The limit is applied to the record alone. + """ + is_blank = str.isspace if python_engine else (lambda text: not text.strip(" \t")) + special = re.compile(f"[{re.escape(delimiter)}\n]") + in_quotes = quote_pending = has_content = False + at_field_start = True + n_delimiters = 0 + + for chunk in chunks: + i, n = 0, len(chunk) + while i < n: + if in_quotes: + if quote_pending: + quote_pending = False + if chunk[i] == '"': # -- `""` is a literal quote -- + has_content = True + i += 1 + else: # -- the quote closed the quoted part; the field continues -- + in_quotes = False + continue + j = chunk.find('"', i) + # -- the Python engine judges blankness by the value, quoted text included -- + quoted_text = chunk[i:] if j < 0 else chunk[i:j] + if quoted_text and not quoted_text.isspace(): + has_content = True + if j < 0: + break + quote_pending, i = True, j + 1 + continue + + if at_field_start and chunk[i] == '"': + # -- a quote alone makes a record non-blank for the C engine, which skips blank + # -- lines before parsing quotes -- + in_quotes, at_field_start = True, False + has_content = has_content or not python_engine + i += 1 + continue + + match = special.search(chunk, i) + j = match.start() if match else n + if j > i: + at_field_start = False + has_content = has_content or not is_blank(chunk[i:j]) + if match is None: + break + i = j + 1 + if chunk[j] == delimiter: + n_delimiters += 1 + has_content = at_field_start = True + if n_delimiters + 1 > max_cells: + raise_limit_exceeded(1, n_delimiters + 1) + elif has_content: # -- end of the first record -- + if n_delimiters + 1 > max_cells: + raise_limit_exceeded(1, n_delimiters + 1) + return n_delimiters + 1, itertools.chain([chunk[i:]], chunks) + else: # -- a blank line, which Pandas skips -- + at_field_start = True + + if not has_content: + return 0, iter(()) + if n_delimiters + 1 > max_cells: + raise_limit_exceeded(1, n_delimiters + 1) + return n_delimiters + 1, iter(()) + + class _CsvPartitioningContext: """Encapsulates the partitioning-run details. diff --git a/unstructured/partition/tsv.py b/unstructured/partition/tsv.py index 1a5844eea0..dd5c22911a 100644 --- a/unstructured/partition/tsv.py +++ b/unstructured/partition/tsv.py @@ -1,8 +1,11 @@ from __future__ import annotations -from typing import IO, Any, Optional +import contextlib +import shutil +import tempfile +from typing import IO, Any, Iterator, Optional, cast -import pandas as pd +from pandas.io.common import get_handle from unstructured.chunking import add_chunking_strategy from unstructured.common.html_table import HtmlTable @@ -13,6 +16,7 @@ spooled_to_bytes_io_if_needed, ) from unstructured.partition.common.metadata import apply_metadata, get_last_modified_date +from unstructured.partition.csv import check_cell_count, read_delimited_text from unstructured.telemetry import partition_runtime_telemetry DETECTION_ORIGIN: str = "tsv" @@ -44,13 +48,22 @@ def partition_tsv( header = 0 if include_header else None if filename: - dataframe = pd.read_csv(filename, sep="\t", header=header) + # -- like `pd.read_csv(filename)`, decompress a file named e.g. "x.tsv.gz"; the size check + # -- and the read each get their own decompressing handle on the same content -- + with _open_decompressed(filename) as f: + check_cell_count(f, "\t", None, header=include_header) + with _open_decompressed(filename) as f: + dataframe = read_delimited_text(f, sep="\t", header=header, encoding=None) else: assert file is not None # -- Note(scanny): `SpooledTemporaryFile` on Python<3.11 does not implement `.readable()` # -- which triggers an exception on `pd.DataFrame.read_csv()` call. f = spooled_to_bytes_io_if_needed(file) - dataframe = pd.read_csv(f, sep="\t", header=header) + with _rereadable(f) as f: + start = f.tell() + check_cell_count(f, "\t", None, header=include_header) + f.seek(start) + dataframe = read_delimited_text(f, sep="\t", header=header, encoding=None) html_table = HtmlTable.from_html_text( dataframe.to_html(index=False, header=include_header, na_rep="") @@ -64,3 +77,33 @@ def partition_tsv( metadata.detection_origin = DETECTION_ORIGIN return [Table(text=html_table.text, metadata=metadata)] + + +@contextlib.contextmanager +def _open_decompressed(filename: str) -> Iterator[IO[bytes]]: + """Open `filename` for reading bytes, decompressed as its extension (e.g. ".gz") implies.""" + with get_handle(filename, "rb", compression="infer", is_text=False) as handles: + yield cast(IO[bytes], handles.handle) + + +@contextlib.contextmanager +def _rereadable(file: IO[bytes]) -> Iterator[IO[bytes]]: + """`file` itself if it can seek, otherwise a spooled copy of the rest of it. + + The size check and the read each read the file, so a stream that cannot seek back (e.g. a + pipe or socket) is copied once to a temporary file, which spills to disk when large. + """ + try: + seekable = file.seekable() + except (AttributeError, ValueError): + seekable = False + if seekable: + yield file + return + with tempfile.SpooledTemporaryFile(max_size=_SPOOL_MAX_MEMORY_BYTES) as copy: + shutil.copyfileobj(file, copy) + copy.seek(0) + yield cast(IO[bytes], copy) + + +_SPOOL_MAX_MEMORY_BYTES = 64 * 1024 * 1024 diff --git a/unstructured/partition/utils/config.py b/unstructured/partition/utils/config.py index aeeb1ed034..3e684eef46 100644 --- a/unstructured/partition/utils/config.py +++ b/unstructured/partition/utils/config.py @@ -338,5 +338,10 @@ def DOCX_TABLE_MAX_CELLS(self) -> int: """Maximum layout-grid cells, summed across a DOCX document's tables, rendered as HTML""" return self._get_int("DOCX_TABLE_MAX_CELLS", 5_000_000) + @property + def CSV_MAX_CELLS(self) -> int: + """Maximum `rows x columns` cells a CSV or TSV file may span""" + return self._get_int("CSV_MAX_CELLS", 5_000_000) + env_config = ENVConfig()