This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-8048-8b3aefec9dadd2933fa63b53c05e6a152f435175
in repository https://gitbox.apache.org/repos/asf/texera.git

commit 761a0d8c1f97e6b025b4be5000ee6caad2819dfe
Author: Xinyuan Lin <[email protected]>
AuthorDate: Fri Aug 28 02:42:05 2026 +0000

    test(pyamber): cover IcebergDocument uri, count and clear paths (#8048)
    
    ### What changes were proposed in this PR?
    
    A new `TestIcebergDocumentWithMockCatalog` class in
    `test_iceberg_document.py`, taking the file from 9 collected items to
    22.
    
    | Metric | Before | After |
    |---|---|---|
    | Codecov fully-covered lines **today** | 118/124 | **124/124** |
    | Codecov fully-covered lines **after #8040** | 115/124 | **124/124** |
    | Branch arms | 21/26 | **26/26** |
    
    **`iceberg_document.py` reaches 100% on both lines and branches.** +6
    fully-covered lines today, +9 after #8040 — lines 90, 108 and 189 are
    `50% (1/2)` partials, invisible to the current branch-less report.
    
    The nine lines were exactly {79, 80, 81, 84, 90, 108, 109, 189, 190}:
    `get_uri`, `get_count`, `clear`, and the `get_after` skip guard.
    
    ### The measurement problem, and how it was handled
    
    The existing `TestIcebergDocument` needs a live postgres Iceberg
    catalog. CI supplies one; on this machine the `texera` role cannot read
    `iceberg_tables`, so **all 9 of those tests fail locally** and a raw
    local coverage run under-reports the file badly — enough to suggest ~50
    winnable lines instead of 9.
    
    A throwaway sqlite-backed `SqlCatalog` proxy was used to reproduce the
    same code paths, applied identically to both sides and then deleted. The
    honest gap is 9.
    
    **The new tests need none of that.** They patch
    `IcebergCatalogInstance.get_instance` rather than mutating it, so they
    touch no process-global state and need no postgres, filesystem or
    network. They run in 1.25s.
    
    ### Two hazards found in the existing fixture
    
    - **`StorageConfig.initialize` is a process-wide one-shot** that raises
    on a second call, and `test_iceberg_document.py:45` calls it **unguarded
    at module level**. That works only because this module happens to be the
    first module-level initializer imported — two sibling modules guard with
    `if not StorageConfig._initialized`. A new test module initializing at
    module level and sorting before it would break collection for the whole
    suite. Appending a class to the existing file avoids this; a separate
    file would have been a trap.
    - **`IcebergCatalogInstance._catalogs` is a class-level dict** and
    `replace_instance` mutates it for the rest of the session, poisoning
    every later suite in the shared process. The new tests sidestep it by
    patching `get_instance`.
    
    ### Honest discounts
    
    - **`get_uri` and `get_count` have zero callers anywhere in
    `src/main/python`**, and `clear` has none either — the Scala side has
    its own implementation. These are real statements with real logic, not
    `@abstractmethod` stubs, but the value is interface-contract conformance
    rather than live behaviour.
    - **Lines 189–190 need a negative `from_index`.**
    `num_of_skipped_records` is provably 0 when the guard first runs, so the
    only route is `get_after(-1)`. It is drivable through the public API,
    but no production caller produces a negative offset — a defensive-guard
    test. Counted strictly, the bundle is 7 lines and 4 arms.
    - **One mutant survives this class**: `sum(f.file.record_count ...)` to
    `sum(1 ...)` at line 110, killed only by the pre-existing postgres-gated
    `test_get_counts`. Line 110 is already covered in CI so it is not a
    coverage gap, but the kill is not this bundle's and one was not
    manufactured.
    
    ### Verification
    
    Two reviewers returned seven findings; all repaired. **The repair pass
    added zero coverage lines and zero arms** — the builder's bundle already
    reached 124/124. Its entire value was mutation strength: 14 mutants that
    passed the delivered spec now fail. Stated plainly rather than sold as
    more coverage.
    
    Full pyamber suite: `5 failed, 1172 passed, 7 errors` against main's
    `1159 passed`. Failure sets sorted and diffed — **identical by
    identity**, 12 entries. `ruff check` and `ruff format --check` pass on
    CI's scope. The throwaway proxy is deleted and `git status --porcelain`
    is clean.
    
    No production file is touched; `git diff -- 'amber/src/main/*'` is
    empty.
    
    ### Any related issues, documentation, discussions?
    
    Closes #8044
    
    ### How was this PR tested?
    
    ```
    cd amber && python -m pytest -m "not integration" -q 
src/test/python/core/storage/iceberg/test_iceberg_document.py -k MockCatalog
    ```
    
    ```
    13 passed, 1 warning in 1.25s
    ```
    
    The 9 pre-existing `TestIcebergDocument` tests still require the
    postgres catalog that CI provides.
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Claude Code (Opus 5)
---
 .../core/storage/iceberg/test_iceberg_document.py  | 191 +++++++++++++++++++++
 1 file changed, 191 insertions(+)

diff --git 
a/amber/src/test/python/core/storage/iceberg/test_iceberg_document.py 
b/amber/src/test/python/core/storage/iceberg/test_iceberg_document.py
index 5e61b4ed68..4e2df0f34c 100644
--- a/amber/src/test/python/core/storage/iceberg/test_iceberg_document.py
+++ b/amber/src/test/python/core/storage/iceberg/test_iceberg_document.py
@@ -22,10 +22,21 @@ import tempfile
 import uuid
 from concurrent.futures import as_completed
 from concurrent.futures.thread import ThreadPoolExecutor
+from unittest.mock import MagicMock, Mock, patch
+from urllib.parse import urlparse
+
+from pyiceberg import types as iceberg_types
+from pyiceberg.schema import Schema as IcebergSchema
 
 from core.models import Schema, Tuple
 from core.models.state import State
 from core.storage.document_factory import DocumentFactory
+from core.storage.iceberg import iceberg_document
+from core.storage.iceberg.iceberg_document import IcebergDocument
+from core.storage.iceberg.iceberg_utils import (
+    amber_tuples_to_arrow_table,
+    arrow_table_to_amber_tuples,
+)
 from core.storage.storage_config import StorageConfig
 from core.storage.vfs_uri_factory import VFSURIFactory
 from proto.org.apache.texera.amber.core import (
@@ -411,3 +422,183 @@ class TestIcebergDocument:
             key=lambda state: state["loop_counter"],
         )
         assert actual_states == states
+
+
+class TestIcebergDocumentWithMockCatalog:
+    """
+    The catalog-facing paths of IcebergDocument that need no catalog service:
+    the table-location lookup, the "table is absent" arms of clear() and
+    get_count(), the locks clear() and the read path take, what each read entry
+    point hands the iterator, and both sides of the iterator's seek guard. The
+    catalog is a mock, so unlike TestIcebergDocument above these run on any 
host
+    without a live Iceberg catalog, and they mutate no cached catalog state.
+    """
+
+    @pytest.fixture
+    def iceberg_schema(self):
+        return IcebergSchema(
+            iceberg_types.NestedField(
+                field_id=1,
+                name="col-int",
+                field_type=iceberg_types.IntegerType(),
+                required=False,
+            )
+        )
+
+    @pytest.fixture
+    def document(self, iceberg_schema):
+        """
+        An IcebergDocument for `ns.tbl` whose catalog is a mock. `get_instance`
+        is patched only for the duration of the construction, so the real
+        per-warehouse catalog cache is never touched.
+        """
+        catalog = Mock()
+        with patch.object(
+            iceberg_document.IcebergCatalogInstance,
+            "get_instance",
+            return_value=catalog,
+        ):
+            return IcebergDocument(
+                "ns",
+                "tbl",
+                iceberg_schema,
+                amber_tuples_to_arrow_table,
+                arrow_table_to_amber_tuples,
+            )
+
+    def test_get_uri_returns_the_parsed_table_location(self, document):
+        # The location is unique per run and the expectation is derived from 
it,
+        # so a get_uri that returned a fixed URI instead of parsing the loaded
+        # table's own location could not accidentally match.
+        location = f"file:///warehouse/{uuid.uuid4().hex}/ns.db/tbl"
+        table = Mock()
+        table.location.return_value = location
+
+        with patch.object(
+            iceberg_document, "load_table_metadata", return_value=table
+        ) as load_table_metadata:
+            uri = document.get_uri()
+
+        expected = urlparse(location)
+        assert (uri.scheme, uri.path) == (expected.scheme, expected.path)
+        table.location.assert_called_once_with()
+        assert load_table_metadata.call_args.args == (document.catalog, "ns", 
"tbl")
+
+    def test_get_uri_rejects_a_table_that_does_not_exist(self, document):
+        with patch.object(iceberg_document, "load_table_metadata", 
return_value=None):
+            with pytest.raises(Exception, match=r"table ns\.tbl doesn't 
exist\."):
+                document.get_uri()
+
+    def test_clear_drops_a_table_that_exists(self, document):
+        document.catalog.table_exists.return_value = True
+
+        document.clear()
+
+        document.catalog.drop_table.assert_called_once_with("ns.tbl")
+
+    def test_clear_leaves_an_absent_table_alone(self, document):
+        document.catalog.table_exists.return_value = False
+
+        document.clear()
+
+        document.catalog.table_exists.assert_called_once_with("ns.tbl")
+        document.catalog.drop_table.assert_not_called()
+
+    def test_get_count_is_zero_when_the_table_does_not_exist(self, document):
+        with patch.object(
+            iceberg_document, "load_table_metadata", return_value=None
+        ) as load_table_metadata:
+            assert document.get_count() == 0
+
+        assert load_table_metadata.call_args.args == (document.catalog, "ns", 
"tbl")
+
+    def test_a_negative_offset_is_rejected_rather_than_read_as_zero(self, 
document):
+        """
+        IcebergIterator guards its file seek against having already skipped 
past
+        `from_index`. The only way the guard can fire is a negative offset, 
since
+        the skip counter is still 0 when the seek generator first runs.
+
+        Neither the guard's message ("seek operation should not be called", 
which
+        describes a re-entrant seek rather than a bad argument) nor its 
exception
+        type is pinned: rejecting a negative offset as a ValueError would be 
the
+        better behaviour, so accepting either type here keeps that fix open 
while
+        still requiring that the offset is rejected rather than read as zero.
+        """
+        iterator = document.get_after(-1)
+
+        with pytest.raises((RuntimeError, ValueError)):
+            next(iterator)
+
+    @pytest.mark.parametrize("offset", [0, 5])
+    def test_a_non_negative_offset_does_not_trip_the_seek_guard(self, 
document, offset):
+        """
+        The other side of the seek guard's boundary: a legal offset must reach
+        the table lookup and then end the iteration cleanly, not raise. Without
+        this, the guard's comparison is unconstrained on hosts that cannot run
+        TestIcebergDocument above.
+        """
+        with patch.object(
+            iceberg_document, "load_table_metadata", return_value=None
+        ) as load_table_metadata:
+            with pytest.raises(StopIteration):
+                next(document.get_after(offset))
+
+        assert load_table_metadata.call_args.args == (document.catalog, "ns", 
"tbl")
+
+    @pytest.mark.parametrize(
+        "read, from_index, until_index, total",
+        [
+            (lambda document: document.get(), 0, None, float("inf")),
+            (lambda document: document.get_range(3, 7), 3, 7, 4),
+            (lambda document: document.get_after(4), 4, None, float("inf")),
+        ],
+    )
+    def test_the_read_entry_points_delegate_to_the_iterator(
+        self, document, read, from_index, until_index, total
+    ):
+        """
+        Each read entry point hands a specific [from, until) range -- and the
+        document's own catalog, table identity, schema and deserde -- to the
+        iterator. The seek generator's body does not run at construction, so
+        this needs no catalog.
+
+        `total` is spelled out per case rather than recomputed from the range,
+        so the expectation does not restate the production formula.
+        """
+        iterator = read(document)
+
+        assert (iterator.from_index, iterator.until_index) == (from_index, 
until_index)
+        assert iterator.total_records_to_return == total
+        assert (iterator.table_namespace, iterator.table_name) == ("ns", "tbl")
+        assert iterator.catalog is document.catalog
+        assert iterator.table_schema is document.table_schema
+        assert iterator.deserde is document.deserde
+
+    def test_the_read_path_takes_the_shared_read_lock(self, document):
+        """
+        The counterpart of test_clear_takes_the_write_lock: reads must take the
+        shared read lock, so that concurrent reads are not serialised behind
+        each other. MagicMock (not Mock) is required: the lock is used as a
+        context manager.
+        """
+        document.lock = MagicMock()
+
+        document.get()
+
+        document.lock.gen_rlock.assert_called_once_with()
+        document.lock.gen_wlock.assert_not_called()
+
+    def test_clear_takes_the_write_lock(self, document):
+        """
+        clear() drops the table, so it must hold the write lock rather than the
+        shared read lock the readers take. MagicMock (not Mock) is required: 
the
+        lock is used as a context manager.
+        """
+        document.catalog.table_exists.return_value = True
+        document.lock = MagicMock()
+
+        document.clear()
+
+        document.lock.gen_wlock.assert_called_once_with()
+        document.lock.gen_rlock.assert_not_called()
+        document.catalog.drop_table.assert_called_once_with("ns.tbl")

Reply via email to