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