From c2f8850280c67da4ef1e85c6634da600272a408e Mon Sep 17 00:00:00 2001 From: Evan Hicks Date: Tue, 22 Sep 2026 16:14:09 -0400 Subject: [PATCH] fix(worker): Send a portless :authority from the push client grpc puts the whole host:port target in :authority, so a push worker dialing task--broker-grpc:50051 through the mesh matches no route -- HttpRoute hostnames are RFC1123 names and cannot carry a port -- and envoy's 404 comes back as UNIMPLEMENTED. Setting grpc.default_authority to the target without its port makes the worker look like the broker, which reaches its worker over a portless URL and has always matched. --- .../taskbroker_client/worker/push_clients.py | 15 +++++++++- .../python/tests/worker/test_push_clients.py | 30 +++++++++++++++++++ 2 files changed, 44 insertions(+), 1 deletion(-) create mode 100644 clients/python/tests/worker/test_push_clients.py diff --git a/clients/python/src/taskbroker_client/worker/push_clients.py b/clients/python/src/taskbroker_client/worker/push_clients.py index f41fb36c..31d7d5d8 100644 --- a/clients/python/src/taskbroker_client/worker/push_clients.py +++ b/clients/python/src/taskbroker_client/worker/push_clients.py @@ -27,6 +27,18 @@ logger = logging.getLogger(__name__) +def service_authority(service: str) -> str: + """The ``:authority`` to send for a push broker service target. + + grpc puts the whole ``host:port`` target in ``:authority``. The sidecar envoy matches + that against the mesh route's hostnames, which are RFC1123 names and cannot carry a + port, so leaving the port on it matches no route and the worker gets a 404 back as + UNIMPLEMENTED. + """ + host, sep, port = service.rpartition(":") + return host if sep and port.isdigit() else service + + class PushTaskbrokerClient: """ Taskworker RPC client wrapper @@ -52,7 +64,8 @@ def __init__( self._processing_pool_name = processing_pool_name or "unknown" self._grpc_options: list[tuple[str, Any]] = [ - ("grpc.max_receive_message_length", MAX_ACTIVATION_SIZE) + ("grpc.max_receive_message_length", MAX_ACTIVATION_SIZE), + ("grpc.default_authority", service_authority(service)), ] if grpc_config: self._grpc_options.append(("grpc.service_config", grpc_config)) diff --git a/clients/python/tests/worker/test_push_clients.py b/clients/python/tests/worker/test_push_clients.py new file mode 100644 index 00000000..9a8a2e16 --- /dev/null +++ b/clients/python/tests/worker/test_push_clients.py @@ -0,0 +1,30 @@ +import pytest + +from taskbroker_client.metrics import NoOpMetricsBackend +from taskbroker_client.worker.push_clients import PushTaskbrokerClient, service_authority + + +@pytest.mark.parametrize( + "service,expected", + [ + ("task-ingest-push-broker-grpc:50051", "task-ingest-push-broker-grpc"), + ("task-ingest-push-broker-grpc:80", "task-ingest-push-broker-grpc"), + ("localhost:50051", "localhost"), + ("[::1]:50051", "[::1]"), + ("task-ingest-push-broker-grpc", "task-ingest-push-broker-grpc"), + ], +) +def test_service_authority_strips_the_port(service: str, expected: str) -> None: + assert service_authority(service) == expected + + +def test_client_sets_default_authority() -> None: + client = PushTaskbrokerClient( + service="task-ingest-push-broker-grpc:50051", + application="tests", + metrics=NoOpMetricsBackend(), + ) + assert ( + "grpc.default_authority", + "task-ingest-push-broker-grpc", + ) in client._grpc_options