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 464aa8c19e chore(pyamber): remove the unused worker-to-worker RPC 
proxy (#7794)
464aa8c19e is described below

commit 464aa8c19e10cedcc727c913443619603b90157d
Author: Xinyuan Lin <[email protected]>
AuthorDate: Wed Sep 2 07:51:50 2026 +0000

    chore(pyamber): remove the unused worker-to-worker RPC proxy (#7794)
    
    ### What changes were proposed in this PR?
    
    Removes `AsyncRPCClient.get_worker_interface` and the `_create_proxy`
    helper it is the sole caller of. **109 lines deleted, 0 added.**
    
    Python workers talk only to the coordinator. The one production
    construction of the Python `AsyncRPCClient` is `main_loop.py:114`, and
    it uses `coordinator_stub()` (lines 271, 651, 671, 817) and `receive`
    (line 475). Worker-to-worker RPC initiated from Python is unused.
    
    | Grep | Result |
    |---|---|
    | `get_worker_interface` | its definition, plus three call sites — all
    in `test_async_rpc_client.py` |
    | `WorkerServiceStub` | `async_rpc_client.py` only: the import, the
    return annotation, one argument. No `.scala` hit |
    | `_create_proxy` | its definition, and one call inside
    `get_worker_interface` — its sole caller |
    | `_stream_unary` / `_unary_stream` / `_stream_stream` | only the
    `Proxy` blockers inside `_create_proxy`, plus those three tests |
    
    There is no `getattr`-style dynamic access anywhere. Scala's
    `workerInterface` (`AsyncRPCClient.scala:140`, `WorkerServiceFs2Grpc`)
    is a separate class and unrelated to this Python method.
    
    `_create_proxy` is removed whole rather than just its `_unary_unary`
    body: with its only caller gone, leaving a `Proxy` class holding three
    stream blockers and no reachable constructor would just relocate the
    dead code. The now-orphaned `WorkerServiceStub` import goes too — it was
    the only orphan, since every other import is still used by
    `_assign_context` / `_create_future` / `receive`.
    
    ### Three tests were removed, deliberately
    
    `test_stream_unary_blocked`, `test_unary_stream_blocked` and
    `test_stream_stream_blocked` each begin with
    `client.get_worker_interface(...)` and then assert on blockers that
    exist only inside `_create_proxy`'s local `Proxy`. They exercise nothing
    else. Their removal orphaned `import pytest` — the file's only `pytest.`
    uses were those three `pytest.raises` calls — so that went as well. The
    coordinator path is untouched: it goes through `_assign_context`, never
    had the blockers, and keeps its own tests.
    
    ### Scope: one item was assessed and deliberately kept
    
    The audit that produced this also flagged `IcebergDocument`'s read side
    (`get_uri`, `get_range`, `get_after`, `get_count`) as having no
    production callers — only `get()` and `writer()` are used
    (`input_port_materialization_reader_runnable.py:161,170`,
    `main_loop.py:157,216`, `output_manager.py:148`).
    
    **None of it is removed here**, because unlike the RPC proxy these are
    working implementations of `@abstractmethod`s on
    `ReadonlyVirtualDocument`, i.e. the read contract `IcebergDocument`
    exists to fulfil:
    
    - `get_range`, `get_after` and `get_count` have live passing tests
    asserting real Iceberg behaviour, and `get_range`/`get_after` are the
    only public entry points to the ranged form of
    `_get_using_file_sequence_order`. Removing them would strand
    `num_of_skipped_records`, `total_records_to_return`, `_skip_records` and
    the skip loop in `IcebergIterator`.
    - `get_uri` has zero references of any kind, so it *is* removable — but
    it is a working override, and deleting it silently changes behaviour
    from "returns the table location" to "raises `NotImplementedError`" via
    `VirtualDocument`'s fallback, and drops parity with Scala's
    `IcebergDocument.getURI`. That is an API decision, not a cleanup, so it
    is left for a maintainer.
    
    ### Verification
    
    | Run | Result |
    |---|---|
    | `core/architecture/rpc/` | **38 passed** |
    | `core/architecture/rpc/` + `core/storage/model/` | 49 passed |
    | `core/architecture` + `core/storage` + `test_main_loop.py`, `-m "not
    integration"` | 448 passed, 1 deselected |
    
    The deselected test is `test_iceberg_rest_catalog_integration.py`,
    excluded by its own `pytestmark`.
    
    `ruff format --check src/main/python src/test/python` → exit 0, "209
    files already formatted". `ruff check` → exit 0, "All checks passed!".
    `git diff --numstat` shows 0 insertions against 109 deletions, so there
    are no incidental rewrites.
    
    ### History
    
    | | |
    | --- | --- |
    | **Introduced by** | #2950 (2024-10-31) — "Migrate scala control
    messages to protobuf", which built the Python `AsyncRPCClient` with both
    a coordinator stub and this worker-to-worker proxy |
    | **Usage removed by** | **never** — no non-test Python file has
    referenced `get_worker_interface` at any point in the history |
    
    Dead on arrival. Its only callers are the three tests added by #4744
    (2026-05-03), a coverage PR that pinned the already-dead path — which is
    why it looks live.
    
    ### Any related issues, documentation, discussions?
    
    Closes #7793
    
    ### How was this PR tested?
    
    ```
    python -m pytest src/test/python/core/architecture/rpc/ -q
    ```
    
    ```
    38 passed, 1 warning
    ```
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Claude Code (Opus 5)
---
 .../core/architecture/rpc/async_rpc_client.py      | 87 ----------------------
 .../core/architecture/rpc/test_async_rpc_client.py | 22 ------
 2 files changed, 109 deletions(-)

diff --git a/amber/src/main/python/core/architecture/rpc/async_rpc_client.py 
b/amber/src/main/python/core/architecture/rpc/async_rpc_client.py
index 017bc35c56..98104dbc42 100644
--- a/amber/src/main/python/core/architecture/rpc/async_rpc_client.py
+++ b/amber/src/main/python/core/architecture/rpc/async_rpc_client.py
@@ -33,7 +33,6 @@ from proto.org.apache.texera.amber.engine.architecture.rpc 
import (
     ControlReturn,
     ControlInvocation,
     CoordinatorServiceStub,
-    WorkerServiceStub,
     ControlRequest,
 )
 from proto.org.apache.texera.amber.engine.common import 
DirectControlMessagePayloadV2
@@ -114,92 +113,6 @@ class AsyncRPCClient:
         """
         return self._coordinator_service_stub
 
-    def get_worker_interface(self, target_worker) -> WorkerServiceStub:
-        """
-        Returns a proxy for interacting with a worker interface.
-
-        :param target_worker: The identifier for the target worker.
-        """
-        return self._create_proxy(
-            WorkerServiceStub, ActorVirtualIdentity(target_worker)
-        )
-
-    def _create_proxy(self, service_class, target_worker: 
ActorVirtualIdentity):
-        """
-        Creates a dynamic proxy for the given service class, allowing
-        asynchronous RPC communication with the specified target actor.
-
-        :param service_class: The service class to be proxied.
-        :param target: The target actor's identity.
-        :return: An instance of the proxy class.
-        """
-        rpc_client = self  # to distinguish outer and inner self
-
-        class Proxy(service_class):
-            def __init__(self, target_actor: ActorVirtualIdentity):
-                self.target_actor = target_actor
-
-            async def _unary_unary(
-                self, route: str, request, response_type, *, timeout, 
deadline, metadata
-            ):
-                """
-                Handles unary-unary RPC calls by creating a ControlInvocation 
command
-                and sending it to the target actor.
-
-                :param route: The RPC route name.
-                :param request: The request message to be sent.
-                :param response_type: The expected response type (unused here).
-                :param timeout: The RPC call timeout (unused here).
-                :param deadline: The RPC call deadline (unused here).
-                :param metadata: Metadata for the RPC call (unused here).
-                :return: A future representing the RPC response.
-                """
-                rpc_context: AsyncRpcContext = AsyncRpcContext(
-                    ActorVirtualIdentity(rpc_client._context.worker_id),
-                    self.target_actor,
-                )
-                to = rpc_context.receiver
-                control_command = ControlInvocation(
-                    # to align with java side, only use the method name
-                    method_name=route.split("/")[-1],
-                    command=set_one_of(ControlRequest, request),
-                    context=rpc_context,
-                    command_id=rpc_client._send_sequences[to],
-                )
-                payload = set_one_of(
-                    DirectControlMessagePayloadV2,
-                    control_command,
-                )
-                rpc_client._output_queue.put(
-                    DCMElement(
-                        tag=ChannelIdentity(
-                            rpc_context.sender, rpc_context.receiver, True
-                        ),
-                        payload=payload,
-                    )
-                )
-                return rpc_client._create_future(to)
-
-            def _stream_unary(self, *args, **kwargs):
-                """Block the _stream_unary method."""
-                raise NotImplementedError(
-                    "Rpc call invokes _stream_unary, which is not supported."
-                )
-
-            def _unary_stream(self, *args, **kwargs):
-                """Block the _unary_stream method."""
-                raise NotImplementedError(
-                    "Rpc call invokes _unary_stream, which is not supported."
-                )
-
-            def _stream_stream(self, *args, **kwargs):
-                """Block the _stream_stream method."""
-                raise NotImplementedError(
-                    "Rpc call invokes _stream_stream, which is not supported."
-                )
-
-        return Proxy(target_worker)
-
     def _create_future(self, to: ActorVirtualIdentity) -> Future:
         """
         Create a promise for the target actor, recording the 
CommandInvocations sent
diff --git 
a/amber/src/test/python/core/architecture/rpc/test_async_rpc_client.py 
b/amber/src/test/python/core/architecture/rpc/test_async_rpc_client.py
index afa048b998..8d2336f493 100644
--- a/amber/src/test/python/core/architecture/rpc/test_async_rpc_client.py
+++ b/amber/src/test/python/core/architecture/rpc/test_async_rpc_client.py
@@ -21,8 +21,6 @@ from concurrent.futures import Future
 from types import SimpleNamespace
 from unittest.mock import MagicMock
 
-import pytest
-
 from core.architecture.rpc import async_rpc_client as async_rpc_client_module
 from core.architecture.rpc.async_rpc_client import AsyncRPCClient, async_run
 from proto.org.apache.texera.amber.core import (
@@ -185,26 +183,6 @@ class TestReceive:
         assert fut.done() and fut.result() is ret
 
 
-class TestProxyStreamBlockers:
-    def test_stream_unary_blocked(self):
-        client = _make_client()
-        proxy = client.get_worker_interface("worker-X")
-        with pytest.raises(NotImplementedError, match="_stream_unary"):
-            proxy._stream_unary()
-
-    def test_unary_stream_blocked(self):
-        client = _make_client()
-        proxy = client.get_worker_interface("worker-X")
-        with pytest.raises(NotImplementedError, match="_unary_stream"):
-            proxy._unary_stream()
-
-    def test_stream_stream_blocked(self):
-        client = _make_client()
-        proxy = client.get_worker_interface("worker-X")
-        with pytest.raises(NotImplementedError, match="_stream_stream"):
-            proxy._stream_stream()
-
-
 class TestCoordinatorStub:
     def test_coordinator_stub_returns_configured_stub(self):
         client = _make_client()

Reply via email to