This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new 5e6d925ee8 test(pyamber): close the schema, table, proxy and
network-sender gaps (#8067)
5e6d925ee8 is described below
commit 5e6d925ee875b7286ff591e87c461e13d8cb701e
Author: Xinyuan Lin <[email protected]>
AuthorDate: Sat Aug 29 07:13:22 2026 +0000
test(pyamber): close the schema, table, proxy and network-sender gaps
(#8067)
### What changes were proposed in this PR?
Seven pyamber files, five existing specs extended and two new ones
added. **+15 fully-covered lines; six of the seven files reach 100%.**
| File | Missed before | After |
|---|---|---|
| `handlers/control/open_executor_handler.py` | 2 (67%) | **0 — 100%** |
| `models/schema/schema.py` | 3 (94%) | **0 — 100%** |
| `models/schema/attribute_type.py` | 2 (96%) | **0 — 100%** |
| `models/table.py` | 2 (95%) | **0 — 100%** |
| `proxy/proxy_client.py` | 2 (94%) | **0 — 100%** |
| `runnables/network_sender.py` | 2 (96%) | **0 — 100%** |
| `storage/iceberg/iceberg_utils.py` | 4 local / 2 CI | 2 left |
**No file here is worth a PR alone** — each carries two or three lines.
The bundle total is the point, and it is stated that way rather than
dressed up per file.
`iceberg_utils.py` closes lines 239 and 275; its remaining 302–303 are
**already green in CI** and worth zero. The local figure of 4 missed is
a local-environment artefact, which is why the honest CI-equivalent gain
is +15 and not +17.
Tests across the seven files go **37 → 89**.
### A measurement trap worth recording
This machine cannot run the Iceberg-backed tests — the postgres `texera`
role cannot read `iceberg_tables` — so a raw local coverage run
**under-reports**. An earlier assessment of one of these files measured
12 missed lines locally where Codecov showed 2, wrong by 6×. Every
target here was sized from Codecov and the local numbers were used only
for the delta, never for the ceiling.
### Dead production code found, and deliberately not pinned
`attribute_type.py:96-97` normalises a trailing `Z` before calling
`datetime.fromisoformat`. **Deleting the whole block leaves the suite
byte-identical.** That is not a coverage gap — `fromisoformat` has
accepted `Z` since Python 3.11, and the CI matrix is 3.11/3.12/3.13, so
the block is dead on every interpreter the project runs. A 10-case
differential probe confirmed it.
No test was written to pin it. The test that had been named after that
branch was **renamed and annotated**, so nothing in the suite claims to
constrain a branch that cannot be reached.
### Verification
**27 mutations, 26 killed, 1 equivalent** — the equivalent one being the
dead block above. Each was applied one at a time against a tree verified
hash-equal to a snapshot, reverted from that snapshot rather than from
git.
Two adversarial reviewers returned nine findings against the first
draft; all repaired, adding six further tests.
Measured with `--cov=src/main/python` under the exact marker filter CI
uses (`-m "not integration"`), no name filter, identical on both sides.
The before-state restored each tracked spec **by exact single file
path** via `git show HEAD:<path>` and moved the two new specs aside; all
seven were then restored from a scratch snapshot and md5-verified.
Full pyamber suite: `5 failed, 1269 passed, 7 errors`, with the
`FAILED`/`ERROR` set **identical by identity** to main's 12 entries — no
regression. `ruff check` and `ruff format --check` pass on CI's exact
scope (`src/main/python src/test/python`). Both new spec files carry the
Apache licence header.
No production file is touched; `git diff -- 'amber/src/main/*'` is
empty.
### Any related issues, documentation, discussions?
Closes #8065
### How was this PR tested?
```
cd amber && python -m pytest -m "not integration" -q
src/test/python/core/models src/test/python/core/proxy
src/test/python/core/runnables/test_network_sender.py
src/test/python/core/architecture/handlers/control/test_open_executor_handler.py
src/test/python/core/storage/iceberg/test_iceberg_utils_catalog.py
```
```
89 passed, 1 warning in 1.73s
```
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Opus 5)
---------
Signed-off-by: Xinyuan Lin <[email protected]>
Co-authored-by: Copilot Autofix powered by AI
<[email protected]>
---
.../handlers/control/test_open_executor_handler.py | 51 ++++++++++
.../core/models/schema/test_attribute_type.py | 94 +++++++++++++++++
.../test/python/core/models/schema/test_schema.py | 68 +++++++++++++
amber/src/test/python/core/models/test_table.py | 80 +++++++++++++++
.../test/python/core/proxy/test_proxy_client.py | 110 ++++++++++++++++++++
.../python/core/runnables/test_network_sender.py | 87 +++++++++++++++-
.../storage/iceberg/test_iceberg_utils_catalog.py | 113 ++++++++++++++++++++-
7 files changed, 600 insertions(+), 3 deletions(-)
diff --git
a/amber/src/test/python/core/architecture/handlers/control/test_open_executor_handler.py
b/amber/src/test/python/core/architecture/handlers/control/test_open_executor_handler.py
new file mode 100644
index 0000000000..5f6e75da22
--- /dev/null
+++
b/amber/src/test/python/core/architecture/handlers/control/test_open_executor_handler.py
@@ -0,0 +1,51 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+import asyncio
+from types import SimpleNamespace
+from unittest.mock import MagicMock
+
+from core.architecture.handlers.control.open_executor_handler import (
+ OpenExecutorHandler,
+)
+from proto.org.apache.texera.amber.engine.architecture.rpc import (
+ EmptyReturn,
+ EmptyRequest,
+)
+
+
+def make_handler() -> OpenExecutorHandler:
+ """Wire a handler with a SimpleNamespace context exposing
executor_manager."""
+ executor_manager = MagicMock()
+ context = SimpleNamespace(executor_manager=executor_manager)
+ return OpenExecutorHandler(context)
+
+
+class TestOpenExecutorHandler:
+ def test_opens_the_current_executor(self):
+ handler = make_handler()
+ asyncio.run(handler.open_executor(EmptyRequest()))
+ executor = handler.context.executor_manager.executor
+ executor.open.assert_called_once_with()
+ # `open` is the only lifecycle call this handler is allowed to make;
+ # pin that it does not also close or otherwise disturb the executor.
+ executor.close.assert_not_called()
+
+ def test_returns_empty_return(self):
+ handler = make_handler()
+ result = asyncio.run(handler.open_executor(EmptyRequest()))
+ assert isinstance(result, EmptyReturn)
diff --git a/amber/src/test/python/core/models/schema/test_attribute_type.py
b/amber/src/test/python/core/models/schema/test_attribute_type.py
new file mode 100644
index 0000000000..3f737dbb2c
--- /dev/null
+++ b/amber/src/test/python/core/models/schema/test_attribute_type.py
@@ -0,0 +1,94 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+import datetime
+
+import pytest
+
+from core.models.schema.attribute_type import (
+ AttributeType,
+ FROM_STRING_PARSER_MAPPING,
+)
+
+# Go through the dispatch table rather than the private helpers: this is the
+# entry point production code uses when it materializes a string column into
+# a typed field, so the tests pin the reachable behaviour rather than an
+# implementation detail.
+parse_bool = FROM_STRING_PARSER_MAPPING[AttributeType.BOOL]
+parse_timestamp = FROM_STRING_PARSER_MAPPING[AttributeType.TIMESTAMP]
+
+EPOCH = datetime.datetime(1970, 1, 1, tzinfo=datetime.timezone.utc)
+
+
+class TestParseBool:
+ @pytest.mark.parametrize("empty", [None, "", " ", "\t\n"])
+ def test_an_absent_value_parses_as_false(self, empty):
+ # Every "empty" spelling short-circuits to False before the
+ # true/false/numeric ladder is consulted. Covering several spellings
+ # keeps the assertion from being dodged by a narrower emptiness test.
+ assert parse_bool(empty) is False
+
+ @pytest.mark.parametrize(
+ "text, expected",
+ [("true", True), ("TRUE", True), (" True ", True), ("false", False)],
+ )
+ def test_literal_spellings_are_case_and_space_insensitive(self, text,
expected):
+ assert parse_bool(text) is expected
+
+ @pytest.mark.parametrize("text, expected", [("1", True), ("0", False)])
+ def test_a_numeric_value_is_compared_against_zero(self, text, expected):
+ assert parse_bool(text) is expected
+
+ def test_a_non_numeric_non_literal_value_is_rejected(self):
+ with pytest.raises(ValueError):
+ parse_bool("maybe")
+
+
+class TestParseTimestamp:
+ @pytest.mark.parametrize("empty", [None, "", " ", "\t\n"])
+ def test_an_absent_value_parses_as_the_utc_epoch(self, empty):
+ parsed = parse_timestamp(empty)
+ # Assert the exact instant *and* the tzinfo: a naive 1970-01-01 would
+ # compare unequal here, so dropping the timezone is caught too.
+ assert parsed == EPOCH
+ assert parsed.tzinfo == datetime.timezone.utc
+ assert (parsed.year, parsed.month, parsed.day) == (1970, 1, 1)
+
+ def test_a_zulu_suffix_yields_a_utc_aware_instant(self):
+ # Named for the observable outcome, not for the code that produces it.
+ # `_parse_timestamp` rewrites a trailing "Z" into "+00:00" before
+ # calling `fromisoformat`, but `fromisoformat` has accepted "Z" itself
+ # since Python 3.11 and the pyamber CI matrix is 3.11/3.12/3.13 -- so
+ # deleting that rewrite leaves the entire suite green, and this test
+ # must not be credited with pinning it. What it does pin is the
+ # resulting instant and its offset (and, uniquely in this file, the
+ # "+00:00" constant, should the rewrite ever run on an older runtime).
+ assert parse_timestamp("2024-05-06T07:08:09Z") == datetime.datetime(
+ 2024, 5, 6, 7, 8, 9, tzinfo=datetime.timezone.utc
+ )
+
+ def test_a_naive_value_is_assumed_to_be_utc(self):
+ assert parse_timestamp("2024-05-06T07:08:09") == datetime.datetime(
+ 2024, 5, 6, 7, 8, 9, tzinfo=datetime.timezone.utc
+ )
+
+ def test_an_explicit_offset_is_preserved(self):
+ parsed = parse_timestamp("2024-05-06T07:08:09+02:00")
+ assert parsed.utcoffset() == datetime.timedelta(hours=2)
+ assert parsed == datetime.datetime(
+ 2024, 5, 6, 5, 8, 9, tzinfo=datetime.timezone.utc
+ )
diff --git a/amber/src/test/python/core/models/schema/test_schema.py
b/amber/src/test/python/core/models/schema/test_schema.py
index 60e4c848a5..0842024224 100644
--- a/amber/src/test/python/core/models/schema/test_schema.py
+++ b/amber/src/test/python/core/models/schema/test_schema.py
@@ -154,3 +154,71 @@ class TestSchema:
assert round_trip_schema.get_attr_type("field1") ==
AttributeType.STRING
assert round_trip_schema.get_attr_type("field2") ==
AttributeType.LARGE_BINARY
assert round_trip_schema.get_attr_type("field3") == AttributeType.INT
+
+ @pytest.mark.parametrize(
+ "other", ["not a schema", None, 42, {"field-1": "STRING"}, ["field-1"]]
+ )
+ def test_comparing_against_a_non_schema_is_false_not_an_error(self,
schema, other):
+ # `__eq__` guards on isinstance before touching `as_key_value_pairs`,
+ # so a foreign operand must compare unequal rather than raise. Assert
+ # the boolean explicitly (and both directions) so a guard that returned
+ # True would be caught.
+ assert (schema == other) is False
+ assert schema != other
+
+ def test_a_schema_equals_only_a_schema_with_the_same_ordered_pairs(self,
schema):
+ # `__eq__` is defined in terms of `as_key_value_pairs`, and both
+ # operands below are built by replaying that same accessor -- so a
+ # corrupted accessor moves both sides together and `schema == same`
+ # would stay green on its own. Guard each constructed operand against a
+ # literal expectation first, read back through the *independent*
+ # `get_attr_names` accessor, so a degenerate fixture fails here instead
+ # of sailing through the equality claim.
+ same = Schema()
+ for name, attr_type in schema.as_key_value_pairs():
+ same.add(name, attr_type)
+ assert same.get_attr_names() == [f"field-{i}" for i in range(1, 8)]
+ assert schema == same
+
+ reordered = Schema()
+ for name, attr_type in reversed(schema.as_key_value_pairs()):
+ reordered.add(name, attr_type)
+ # ... and this one really is a *reordering* (same names, reverse
order),
+ # not an empty or truncated schema that would compare unequal for a
+ # reason that has nothing to do with ordering.
+ assert reordered.get_attr_names() == [f"field-{i}" for i in range(7,
0, -1)]
+ assert (schema == reordered) is False
+
+ def test_a_partial_schema_keeps_the_order_of_the_requested_names(self,
schema):
+ # `get_partial_schema` documents that it preserves "the order specified
+ # by the attribute names", and Schema equality is order-sensitive (see
+ # the test above), yet nothing pinned that promise. It matters: the
sole
+ # caller, `Tuple.get_partial_tuple`, builds its field values in
+ # `attribute_names` order and takes its schema from here, so a
+ # reordering would make a Tuple's data and its schema silently
disagree.
+ #
+ # The requested names are deliberately out of the source order, and the
+ # assertion is on the ordered pair list rather than on membership --
+ # both are what make this non-vacuous.
+ partial = schema.get_partial_schema(["field-5", "field-2"])
+ assert partial.as_key_value_pairs() == [
+ ("field-5", AttributeType.BOOL),
+ ("field-2", AttributeType.INT),
+ ]
+
+ def test_str_renders_each_attribute_with_a_zero_based_index_and_type(self):
+ rendered = Schema(
+ raw_schema={"field-1": "STRING", "field-2": "INTEGER", "field-3":
"BOOLEAN"}
+ )
+ # Pin the whole rendering: the bracketed header/footer, the ",\n"
+ # separator, the zero-based index, and the name-then-type ordering.
+ assert str(rendered) == (
+ "Schema[\n"
+ "(0)'field-1' -> AttributeType.STRING,\n"
+ "(1)'field-2' -> AttributeType.INT,\n"
+ "(2)'field-3' -> AttributeType.BOOL\n"
+ "]"
+ )
+
+ def test_str_of_an_empty_schema_still_renders_the_brackets(self):
+ assert str(Schema()) == "Schema[\n\n]"
diff --git a/amber/src/test/python/core/models/test_table.py
b/amber/src/test/python/core/models/test_table.py
index 368220779d..80a985e5d6 100644
--- a/amber/src/test/python/core/models/test_table.py
+++ b/amber/src/test/python/core/models/test_table.py
@@ -16,9 +16,11 @@
# under the License.
import datetime
+import numpy
import pandas
import pickle
import pytest
+import re
from pandas import RangeIndex
from core.models import Table, Tuple
@@ -142,3 +144,81 @@ class TestTable:
def test_validation_of_schema(self):
with pytest.raises(AssertionError):
Table([{"text": "hello"}, {"book": "harry"}])
+
+ @pytest.mark.parametrize(
+ "table_like", [42, "hello", None, {"field1": [1, 2]}, (1, 2), b"bytes"]
+ )
+ def test_an_unsupported_tablelike_is_rejected(self, table_like):
+ # Only Table / DataFrame / list reach a constructor; anything else must
+ # be refused with a message naming the offending type, rather than
+ # falling through into `super().__init__` with an unbound frame.
+ #
+ # Match the *whole* rendered message, interpolation included: a prefix
+ # match would leave `{type(table_like)}` -- the only non-constant part
+ # of that line -- unpinned, and would also make all six parametrized
+ # cases assert the identical string.
+ expected = (
+ "^" + re.escape(f"unsupported tablelike type {type(table_like)}")
+ "$"
+ )
+ with pytest.raises(TypeError, match=expected):
+ Table(table_like)
+
+ @pytest.fixture
+ def comparable_frame(self):
+ # Deliberately free of all-None columns: elementwise `None == None`
+ # is False, which would make the comparison below fail for reasons
+ # unrelated to the branch under test.
+ return pandas.DataFrame(
+ {"field1": [1, 2], "field2": ["hello", "world"], "field3": [2.3,
0.0]},
+ columns=["field1", "field2", "field3"],
+ )
+
+ def test_comparing_to_an_equal_data_frame_reports_equal(self,
comparable_frame):
+ # A non-Table operand takes the `super().__eq__` branch. Reduce with
+ # `numpy.all` rather than asserting the result's *shape*: today that
+ # branch yields a per-column Series (see the characterization test
+ # below), and hard-asserting that here would cement the very mismatch
+ # between the code and its `-> bool` annotation.
+ table = Table(comparable_frame)
+ assert numpy.all(table == comparable_frame)
+
+ def test_comparing_to_a_differing_data_frame_reports_unequal(
+ self, comparable_frame
+ ):
+ table = Table(comparable_frame)
+ differing = comparable_frame.copy()
+ differing.loc[0, "field2"] = "goodbye"
+ # One cell of one column differs. Reducing the whole comparison must
+ # therefore be falsy -- which is also what distinguishes the real
+ # `.all()` reduction from an `.any()` one (under `.any()` every column
+ # has at least one matching row, so the reduction would come back
True).
+ assert not numpy.all(table == differing)
+
+ def test_comparing_to_a_data_frame_currently_yields_a_per_column_series(
+ self, comparable_frame
+ ):
+ # CHARACTERIZATION, not a contract. `Table.__eq__` is annotated
+ # `-> bool`, but its non-Table branch returns
+ # `super().__eq__(other).all()`, which for a DataFrame operand reduces
+ # only over rows and leaves a Series indexed by column name. A bare
+ # `assert table == frame` therefore raises "truth value of a Series is
+ # ambiguous". This test records today's shape so that narrowing the
+ # branch to a real bool surfaces here deliberately, with a name that
+ # says so, rather than silently through the behavioural tests above.
+ table = Table(comparable_frame)
+ comparison = table == comparable_frame
+ assert isinstance(comparison, pandas.Series)
+ assert list(comparison.index) == ["field1", "field2", "field3"]
+
+ def test_two_tables_with_differing_rows_are_not_equal(self,
comparable_frame):
+ # The Table-vs-Table arm had only positive coverage: every pre-existing
+ # test in this file compares two *equal* Tables, so replacing the arm's
+ # body with `return True` survived the whole suite. One negative case
+ # closes that.
+ #
+ # Deliberately NOT pinned here: `zip` truncates to the shorter operand,
+ # so `Table(frame.head(1)) == Table(frame)` is True today. That is a
+ # defect, not a contract, and asserting it would cement it.
+ differing = comparable_frame.copy()
+ differing.loc[0, "field2"] = "goodbye"
+ assert (Table(comparable_frame) == Table(differing)) is False
diff --git a/amber/src/test/python/core/proxy/test_proxy_client.py
b/amber/src/test/python/core/proxy/test_proxy_client.py
index 891c2fda75..0c50393c74 100644
--- a/amber/src/test/python/core/proxy/test_proxy_client.py
+++ b/amber/src/test/python/core/proxy/test_proxy_client.py
@@ -20,6 +20,7 @@ from pandas import DataFrame
from pyarrow import ArrowNotImplementedError, Table
from queue import Queue
+import core.proxy.proxy_client as proxy_client_module
from core.proxy.proxy_client import ProxyClient
from core.proxy.proxy_server import ProxyServer
@@ -153,3 +154,112 @@ class TestProxyClient:
assert data_queue.qsize() == 4
for i, row in data_table.to_pandas().iterrows():
assert data_queue.get().equals(row)
+
+ def test_a_handshake_port_is_announced_to_the_server_on_connect(self,
server):
+ # The Java side learns which port the Python proxy *server* listens on
+ # through a "handshake" action issued while the client connects.
+ # ProxyServer does not register "handshake" itself, so the test
+ # registers a capturing stand-in for the Java handler.
+ received = []
+ server.register("handshake", lambda payload: received.append(payload)
or "ok")
+
+ client = ProxyClient(handshake_port=6789)
+ try:
+ # The payload must be the handshake port rendered as UTF-8 digits,
+ # not the client's timeout or the data port it dialled. (The action
+ # *name* is pinned separately, by the test below -- see why there.)
+ assert received == [b"6789"]
+ finally:
+ client.close()
+
+ def test_the_handshake_call_names_the_handshake_action(self, monkeypatch):
+ # The end-to-end test above cannot pin the action *name*. Substituting
+ # any other registered name makes construction blow up server-side (the
+ # built-in `heartbeat`/`shutdown` handlers take no argument, while
+ # `control`/`actor` try to deserialize the payload), so the test dies
+ # during `ProxyClient(...)` and its `received` assertion never runs --
+ # the name is pinned by an accident of the fixture, not by an
+ # assertion. Stub the call at the client boundary instead, so the name
+ # reaches an assertion rather than a server-side crash.
+ seen = []
+ monkeypatch.setattr(
+ ProxyClient,
+ "call_action",
+ lambda self, name, payload=bytes(), options=None: seen.append(
+ (name, payload)
+ ),
+ )
+
+ client = ProxyClient(handshake_port=6789)
+ try:
+ assert seen == [("handshake", b"6789")]
+ finally:
+ client.close()
+
+ def test_no_handshake_is_sent_when_no_handshake_port_is_given(self,
server):
+ # Companion negative case: without it, unconditionally calling
+ # `_handshake` would still satisfy the test above.
+ received = []
+ server.register("handshake", lambda payload: received.append(payload)
or "ok")
+
+ client = ProxyClient()
+ try:
+ assert received == []
+ # ... and the client is still usable, i.e. skipping the handshake
+ # is a real branch rather than a failed connect.
+ assert client.call_action("heartbeat") == b"ack"
+ finally:
+ client.close()
+
+ @staticmethod
+ def _record_call_options(monkeypatch):
+ """
+ Record every `FlightCallOptions(...)` `call_action` manufactures.
+
+ A plain factory function rather than a subclass: `FlightCallOptions` is
+ a Cython cdef class. It still delegates to the real constructor, so the
+ RPC underneath stays a real one.
+ """
+ real = proxy_client_module.FlightCallOptions
+ captured = []
+
+ def recording_factory(**kwargs):
+ captured.append(kwargs)
+ return real(**kwargs)
+
+ monkeypatch.setattr(proxy_client_module, "FlightCallOptions",
recording_factory)
+ return captured, real
+
+ def test_the_configured_timeout_is_applied_when_no_options_are_supplied(
+ self, server, monkeypatch
+ ):
+ # This is the only place the client's configured timeout ever reaches
+ # an RPC, and nothing constrained it: dropping the argument entirely
+ # (`FlightCallOptions()`) survived the whole suite. Without a timeout a
+ # hung Java-side server would block the Python worker forever instead
+ # of raising.
+ captured, _ = self._record_call_options(monkeypatch)
+
+ # A non-default timeout, so the assertion pins the value *flowing from
+ # the constructor* rather than a hard-coded 1000.
+ client = ProxyClient(timeout=1234)
+ try:
+ assert client.call_action("heartbeat") == b"ack"
+ finally:
+ client.close()
+
+ assert captured == [{"timeout": 1234}]
+
+ def test_caller_supplied_options_are_used_as_is(self, server, monkeypatch):
+ # The companion arm: when the caller hands in an options object,
+ # `call_action` must not manufacture its own and discard it.
+ captured, real = self._record_call_options(monkeypatch)
+
+ client = ProxyClient(timeout=1234)
+ try:
+ options = real(timeout=99)
+ assert client.call_action("heartbeat", options=options) == b"ack"
+ finally:
+ client.close()
+
+ assert captured == []
diff --git a/amber/src/test/python/core/runnables/test_network_sender.py
b/amber/src/test/python/core/runnables/test_network_sender.py
index 529cd19d33..bca9c52a58 100644
--- a/amber/src/test/python/core/runnables/test_network_sender.py
+++ b/amber/src/test/python/core/runnables/test_network_sender.py
@@ -17,11 +17,39 @@
import pytest
import threading
+from contextlib import contextmanager
from time import sleep
-from core.models.internal_queue import InternalQueue
+from loguru import logger
+
+from core.models.internal_queue import InternalQueue, InternalQueueElement
+from core.models.payload import DataFrame
from core.runnables.network_receiver import NetworkReceiver
from core.runnables.network_sender import NetworkSender
+from proto.org.apache.texera.amber.core import (
+ ActorVirtualIdentity,
+ ChannelIdentity,
+)
+
+
+@contextmanager
+def muted_catch_logs():
+ """
+ `NetworkSender._send_data` is wrapped in `@logger.catch(reraise=True)`,
+ which logs the whole traceback at ERROR before re-raising. CI runs pytest
+ with `-s` and `LOGURU_LEVEL=WARNING`, so an expected failure would
+ otherwise dump a traceback into the build log and read as a real error.
+
+ `logger.catch` attributes the record to the *caller's* module rather than
+ to the decorated function's module, so the name to silence is this test
+ module's own `__name__` (which varies with pytest's import mode -- hence
+ `__name__` rather than a literal).
+ """
+ logger.disable(__name__)
+ try:
+ yield
+ finally:
+ logger.enable(__name__)
class TestNetworkSender:
@@ -67,3 +95,60 @@ class TestNetworkSender:
assert not network_sender_thread.is_alive()
network_receiver_thread.join()
network_sender_thread.join()
+
+ @pytest.fixture
+ def channel_id(self):
+ worker_id = ActorVirtualIdentity(name="test")
+ return ChannelIdentity(worker_id, worker_id, False)
+
+ @pytest.mark.timeout(5)
+ def test_receive_rejects_an_element_that_is_neither_data_control_nor_ecm(
+ self, network_sender, channel_id
+ ):
+ # A plain subclass of InternalQueueElement is neither DataElement nor
+ # DCMElement nor ECMElement, so it walks the whole dispatch chain and
+ # falls off the end. The sender must refuse it loudly rather than drop
+ # it silently.
+ #
+ # The stable `__repr__` lets the matcher pin the *interpolated* entry
+ # as well, mirroring the payload test below: matching only the constant
+ # prefix would leave `{next_entry}` -- the sole non-constant part of
+ # that line -- unconstrained.
+ class UnknownEntry(InternalQueueElement):
+ def __repr__(self):
+ return "<unknown-entry>"
+
+ unknown = UnknownEntry(tag=channel_id)
+ with pytest.raises(TypeError, match="Unexpected entry
<unknown-entry>"):
+ network_sender.receive(unknown)
+
+ @pytest.mark.timeout(5)
+ def
test_send_data_rejects_a_payload_that_is_neither_dataframe_nor_stateframe(
+ self, network_sender, channel_id
+ ):
+ class NotAPayload:
+ def __repr__(self):
+ return "<not-a-payload>"
+
+ with muted_catch_logs():
+ with pytest.raises(TypeError, match="Unexpected payload
<not-a-payload>"):
+ network_sender._send_data(channel_id, NotAPayload())
+
+ @pytest.mark.timeout(5)
+ def test_receive_routes_a_data_element_to_send_data(
+ self, network_sender, channel_id, monkeypatch
+ ):
+ # Guards the dispatch chain above: without a positive case, swapping
+ # the DataElement arm for the else-arm would still leave the negative
+ # test green.
+ from core.models.internal_queue import DataElement
+
+ seen = []
+ monkeypatch.setattr(
+ network_sender,
+ "_send_data",
+ lambda to, payload: seen.append((to, payload)),
+ )
+ payload = DataFrame(frame=None)
+ network_sender.receive(DataElement(tag=channel_id, payload=payload))
+ assert seen == [(channel_id, payload)]
diff --git
a/amber/src/test/python/core/storage/iceberg/test_iceberg_utils_catalog.py
b/amber/src/test/python/core/storage/iceberg/test_iceberg_utils_catalog.py
index a387a51597..c3fb6e3541 100644
--- a/amber/src/test/python/core/storage/iceberg/test_iceberg_utils_catalog.py
+++ b/amber/src/test/python/core/storage/iceberg/test_iceberg_utils_catalog.py
@@ -15,10 +15,16 @@
# specific language governing permissions and limitations
# under the License.
-from unittest.mock import patch
+from unittest.mock import MagicMock, patch
+
+from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC
from core.storage.iceberg import iceberg_utils
-from core.storage.iceberg.iceberg_utils import create_postgres_catalog
+from core.storage.iceberg.iceberg_utils import (
+ create_postgres_catalog,
+ create_rest_catalog,
+ create_table,
+)
class TestCreatePostgresCatalog:
@@ -255,3 +261,106 @@ class TestCreatePostgresCatalogWindowsLocal:
assert kwargs["warehouse"] == ""
assert "py-io-impl" not in kwargs
+
+
+class TestCreateRestCatalog:
+ """
+ `create_rest_catalog` is the REST counterpart of `create_postgres_catalog`.
+
+ Unlike the postgres path it deliberately performs *no* Windows-local
+ normalization: `warehouse_name` is a logical warehouse identifier the REST
+ server resolves, and its I/O goes through S3FileIO rather than the local
+ filesystem. `load_catalog` is patched so these run without a REST server.
+ """
+
+ def _make(self, catalog_name, warehouse_name, rest_uri):
+ with patch.object(iceberg_utils, "load_catalog") as mock_load_catalog:
+ catalog = create_rest_catalog(
+ catalog_name=catalog_name,
+ warehouse_name=warehouse_name,
+ rest_uri=rest_uri,
+ )
+ assert mock_load_catalog.call_count == 1
+ assert catalog is mock_load_catalog.return_value
+ return mock_load_catalog.call_args
+
+ def test_catalog_name_uri_and_warehouse_are_forwarded_by_value(self):
+ args, kwargs = self._make(
+ "texera_iceberg", "texera-warehouse",
"http://lakekeeper:8181/catalog"
+ )
+ # Assert every value, not just `type`: swapping `uri` and `warehouse`
+ # would otherwise go unnoticed.
+ assert args == ("texera_iceberg",)
+ assert kwargs == {
+ "type": "rest",
+ "uri": "http://lakekeeper:8181/catalog",
+ "warehouse": "texera-warehouse",
+ }
+
+ def test_a_windows_style_warehouse_name_is_not_normalized(self):
+ """
+ The postgres path rewrites `C:\\...` into a `file:///` URI; the REST
+ path must not, because the identifier is resolved server-side.
+ """
+ _, kwargs = self._make(
+ "texera_iceberg", "C:\\Users\\texera\\warehouse",
"http://localhost:8181"
+ )
+ assert kwargs["warehouse"] == "C:\\Users\\texera\\warehouse"
+ assert "py-io-impl" not in kwargs
+
+
+class TestCreateTable:
+ """
+ `create_table` against a mocked `Catalog`: no postgres, no REST server,
+ no object store.
+ """
+
+ NAMESPACE = "ns"
+ NAME = "tbl"
+ IDENTIFIER = "ns.tbl"
+
+ def _catalog(self, table_exists):
+ catalog = MagicMock()
+ catalog.table_exists.return_value = table_exists
+ return catalog
+
+ def _create(self, catalog, override_if_exists):
+ schema = MagicMock(name="schema")
+ table = create_table(
+ catalog=catalog,
+ table_namespace=self.NAMESPACE,
+ table_name=self.NAME,
+ table_schema=schema,
+ override_if_exists=override_if_exists,
+ )
+ return table, schema
+
+ def test_an_existing_table_is_dropped_when_override_is_requested(self):
+ catalog = self._catalog(table_exists=True)
+ table, schema = self._create(catalog, override_if_exists=True)
+
+
catalog.create_namespace_if_not_exists.assert_called_once_with(self.NAMESPACE)
+ catalog.table_exists.assert_called_once_with(self.IDENTIFIER)
+ catalog.drop_table.assert_called_once_with(self.IDENTIFIER)
+ catalog.create_table.assert_called_once_with(
+ identifier=self.IDENTIFIER,
+ schema=schema,
+ partition_spec=UNPARTITIONED_PARTITION_SPEC,
+ )
+ assert table is catalog.create_table.return_value
+
+ def test_an_existing_table_is_kept_when_override_is_not_requested(self):
+ # Companion to the case above: without it, relaxing the guard from
+ # `and` to `or` would survive.
+ catalog = self._catalog(table_exists=True)
+ self._create(catalog, override_if_exists=False)
+
+ catalog.drop_table.assert_not_called()
+ catalog.create_table.assert_called_once()
+
+ def
test_a_missing_table_is_never_dropped_even_when_override_is_requested(self):
+ catalog = self._catalog(table_exists=False)
+ self._create(catalog, override_if_exists=True)
+
+ catalog.drop_table.assert_not_called()
+ catalog.create_table.assert_called_once()