From 1ef79f5fbbce0417e80b71bc9b3309b2427a2cfa Mon Sep 17 00:00:00 2001 From: Sebastian Schildt Date: Tue, 25 Aug 2026 20:49:14 +0200 Subject: [PATCH 1/2] Fix wildcard subscribe Signed-off-by: Sebastian Schildt --- kuksa-client/kuksa_client/cli_backend/grpc.py | 14 +- kuksa-client/tests/test_cli_backend.py | 163 ++++++++++++++++++ 2 files changed, 174 insertions(+), 3 deletions(-) create mode 100644 kuksa-client/tests/test_cli_backend.py diff --git a/kuksa-client/kuksa_client/cli_backend/grpc.py b/kuksa-client/kuksa_client/cli_backend/grpc.py index 36b82c9..69ac94a 100644 --- a/kuksa-client/kuksa_client/cli_backend/grpc.py +++ b/kuksa-client/kuksa_client/cli_backend/grpc.py @@ -270,10 +270,18 @@ async def _grpcHandler(self, vss_client: kuksa_client.grpc.aio.VSSClient): subscriber_response_stream = vss_client.v2_subscribe(paths=paths) resp = await subscriber_manager.add_subscriber(subscriber_response_stream, callback) except kuksa_client.grpc.VSSClientError as exc: - if exc.error["code"] != grpc.StatusCode.UNIMPLEMENTED.value[0]: + if exc.error["code"] == grpc.StatusCode.NOT_FOUND.value[0]: + logger.debug( + "v2 Subscribe returned NOT_FOUND; expanding branch paths via ListMetadata" + ) + expanded = await vss_client._expand_v2_branch_paths(paths) + subscriber_response_stream = vss_client.v2_subscribe(paths=expanded) + resp = await subscriber_manager.add_subscriber(subscriber_response_stream, callback) + elif exc.error["code"] != grpc.StatusCode.UNIMPLEMENTED.value[0]: raise - subscriber_response_stream = vss_client.subscribe(entries=entries) - resp = await subscriber_manager.add_subscriber(subscriber_response_stream, callback) + else: + subscriber_response_stream = vss_client.subscribe(entries=entries) + resp = await subscriber_manager.add_subscriber(subscriber_response_stream, callback) resp = {"subscriptionId": str(resp)} elif call == "unsubscribe": resp = await subscriber_manager.remove_subscriber(**requestArgs) diff --git a/kuksa-client/tests/test_cli_backend.py b/kuksa-client/tests/test_cli_backend.py new file mode 100644 index 0000000..168e51a --- /dev/null +++ b/kuksa-client/tests/test_cli_backend.py @@ -0,0 +1,163 @@ +######################################################################## +# Copyright (c) 2025 Robert Bosch GmbH +# +# Licensed 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. +# +# SPDX-License-Identifier: Apache-2.0 +######################################################################## + +import asyncio +import queue +import uuid + +import grpc +import pytest + +import kuksa_client.grpc.aio as aio_mod +from kuksa_client.cli_backend import grpc as grpc_backend +from kuksa_client.grpc import Field +from kuksa_client.grpc import SubscribeEntry +from kuksa_client.grpc import View +from kuksa_client.grpc import VSSClientError + + +def _not_found_error(): + return VSSClientError( + error={ + "code": grpc.StatusCode.NOT_FOUND.value[0], + "reason": grpc.StatusCode.NOT_FOUND.value[1], + "message": "Path not found", + }, + errors=[], + ) + + +def _unimplemented_error(): + return VSSClientError( + error={ + "code": grpc.StatusCode.UNIMPLEMENTED.value[0], + "reason": grpc.StatusCode.UNIMPLEMENTED.value[1], + "message": "Not implemented", + }, + errors=[], + ) + + +class TestGrpcCliBackend: + @staticmethod + def _backend(): + return grpc_backend.Backend( + { + "protocol": "grpc", + "ip": "127.0.0.1", + "port": 55555, + "insecure": True, + } + ) + + @staticmethod + def _subscribe_request(paths, entries): + return { + "paths": paths, + "entries": entries, + "callback": lambda updates: None, + } + + async def _process(self, backend, vss_client, request): + response_queue = queue.Queue(maxsize=1) + backend.sendMsgQueue.put(("subscribe", request, response_queue)) + task = asyncio.create_task(backend._grpcHandler(vss_client)) + try: + for _ in range(100): + try: + return response_queue.get_nowait() + except queue.Empty: + await asyncio.sleep(0.01) + pytest.fail("No response received from _grpcHandler") + finally: + backend.run = False + await task + + @pytest.mark.asyncio + async def test_subscribe_leaf_path(self, mocker): + backend = self._backend() + vss_client = mocker.MagicMock() + vss_client.v2_subscribe.return_value = "v2_stream" + subscriber_manager = mocker.MagicMock() + sub_id = uuid.uuid4() + subscriber_manager.add_subscriber = mocker.AsyncMock(return_value=sub_id) + mocker.patch.object(aio_mod, "SubscriberManager", return_value=subscriber_manager) + + entries = [SubscribeEntry("Vehicle.Speed", View.CURRENT_VALUE, (Field.VALUE,))] + resp, error = await self._process( + backend, vss_client, self._subscribe_request(["Vehicle.Speed"], entries) + ) + + assert error is None + assert resp == {"subscriptionId": str(sub_id)} + vss_client.v2_subscribe.assert_called_once_with(paths=["Vehicle.Speed"]) + vss_client._expand_v2_branch_paths.assert_not_called() + vss_client.subscribe.assert_not_called() + + @pytest.mark.asyncio + async def test_subscribe_expands_branch_paths_on_not_found(self, mocker): + backend = self._backend() + vss_client = mocker.MagicMock() + vss_client.v2_subscribe.side_effect = ["v2_stream_first", "v2_stream_second"] + vss_client._expand_v2_branch_paths = mocker.AsyncMock( + return_value=["Vehicle.Speed", "Vehicle.ADAS.ABS.IsActive"] + ) + subscriber_manager = mocker.MagicMock() + sub_id = uuid.uuid4() + subscriber_manager.add_subscriber = mocker.AsyncMock( + side_effect=[_not_found_error(), sub_id] + ) + mocker.patch.object(aio_mod, "SubscriberManager", return_value=subscriber_manager) + + entries = [SubscribeEntry("Vehicle", View.CURRENT_VALUE, (Field.VALUE,))] + resp, error = await self._process( + backend, vss_client, self._subscribe_request(["Vehicle"], entries) + ) + + assert error is None + assert resp == {"subscriptionId": str(sub_id)} + assert vss_client.v2_subscribe.call_args_list == [ + mocker.call(paths=["Vehicle"]), + mocker.call(paths=["Vehicle.Speed", "Vehicle.ADAS.ABS.IsActive"]), + ] + vss_client._expand_v2_branch_paths.assert_called_once_with(["Vehicle"]) + vss_client.subscribe.assert_not_called() + + @pytest.mark.asyncio + async def test_subscribe_falls_back_to_v1_on_unimplemented(self, mocker): + backend = self._backend() + vss_client = mocker.MagicMock() + vss_client.v2_subscribe.return_value = "v2_stream" + vss_client.subscribe.return_value = "v1_stream" + subscriber_manager = mocker.MagicMock() + sub_id = uuid.uuid4() + subscriber_manager.add_subscriber = mocker.AsyncMock( + side_effect=[_unimplemented_error(), sub_id] + ) + mocker.patch.object(aio_mod, "SubscriberManager", return_value=subscriber_manager) + + entries = [SubscribeEntry("Vehicle.Speed", View.CURRENT_VALUE, (Field.VALUE,))] + resp, error = await self._process( + backend, vss_client, self._subscribe_request(["Vehicle.Speed"], entries) + ) + + assert error is None + assert resp == {"subscriptionId": str(sub_id)} + vss_client.v2_subscribe.assert_called_once_with(paths=["Vehicle.Speed"]) + vss_client.subscribe.assert_called_once_with(entries=entries) + vss_client._expand_v2_branch_paths.assert_not_called() From 6ae573a0389270b429ae33b2d63229baa88a8709 Mon Sep 17 00:00:00 2001 From: Sebastian Schildt Date: Tue, 25 Aug 2026 21:13:59 +0200 Subject: [PATCH 2/2] Remove broken subscribeMultiple and add feature to normal subscribe call Signed-off-by: Sebastian Schildt --- docs/cli.md | 3 +- kuksa-client/kuksa_client/__main__.py | 51 +++------------------ kuksa-client/kuksa_client/cli_backend/ws.py | 7 +-- 3 files changed, 11 insertions(+), 50 deletions(-) diff --git a/docs/cli.md b/docs/cli.md index 1e3f22a..881ceaa 100644 --- a/docs/cli.md +++ b/docs/cli.md @@ -155,8 +155,7 @@ setTargetValue Set the target value of a path setTargetValues Set the target value of given paths setValue Set the value of a path setValues Set the value of given paths -subscribe Subscribe the value of a path -subscribeMultiple Subscribe to updates of given paths +subscribe Subscribe to updates of given paths unsubscribe Unsubscribe an existing subscription updateMetaData Update MetaData of a given path updateVSSTree Update VSS Tree Entry diff --git a/kuksa-client/kuksa_client/__main__.py b/kuksa-client/kuksa_client/__main__.py index 6c2c009..9f87b31 100755 --- a/kuksa-client/kuksa_client/__main__.py +++ b/kuksa-client/kuksa_client/__main__.py @@ -242,7 +242,7 @@ def subscriptionIdCompleter(self, text, line, begidx, endidx): ap_subscribe = Cmd2ArgumentParser() ap_subscribe.add_argument( - "Path", help="Path to subscribe to", completer=path_completer + "Path", help="Path to subscribe to", nargs="+", completer=path_completer ) ap_subscribe.add_argument( "-a", "--attribute", help="Attribute to subscribe to", default="value" @@ -255,20 +255,6 @@ def subscriptionIdCompleter(self, text, line, begidx, endidx): action="store_true", ) - ap_subscribeMultiple = Cmd2ArgumentParser() - ap_subscribeMultiple.add_argument( - "Path", help="Path to subscribe to", nargs="+", completer=path_completer - ) - ap_subscribeMultiple.add_argument( - "-a", "--attribute", help="Attribute to subscribe to", default="value" - ) - ap_subscribeMultiple.add_argument( - "-f", - "--output-to-file", - help="Redirect the subscription output to file", - action="store_true", - ) - ap_unsubscribe = Cmd2ArgumentParser() ap_unsubscribe.add_argument( "SubscribeId", @@ -443,43 +429,18 @@ def do_getTargetValues(self, args): @with_category(VSS_COMMANDS) @with_argparser(ap_subscribe) def do_subscribe(self, args): - """Subscribe the value of a path""" - if self.connection_established(): - if args.output_to_file: - logPath = ( - pathlib.Path.cwd() - / f"log_{args.Path.replace('/', '.')}_{args.attribute}_{str(time.time())}" - ) - callback = functools.partial(self.subscribeCallback, logPath) - else: - callback = functools.partial(self.subscribeCallback, None) - - resp = self.commThread.subscribe(args.Path, callback, args.attribute) - resJson = json.loads(resp) - if "subscriptionId" in resJson: - self.subscribeIds.add(resJson["subscriptionId"]) - if args.output_to_file: - logPath.touch() - print(f"Subscription log available at {logPath}") - print(highlight(resp, lexers.JsonLexer(), formatters.TerminalFormatter())) - self.pathCompletionItems = [] - - @with_category(VSS_COMMANDS) - @with_argparser(ap_subscribeMultiple) - def do_subscribeMultiple(self, args): """Subscribe to updates of given paths""" if self.connection_established(): if args.output_to_file: logPath = ( pathlib.Path.cwd() - / f"subscribeMultiple_{args.attribute}_{str(time.time())}.log" + / f"log_{'_'.join(args.Path).replace('/', '.')}_{args.attribute}_{str(time.time())}" ) callback = functools.partial(self.subscribeCallback, logPath) else: callback = functools.partial(self.subscribeCallback, None) - resp = self.commThread.subscribeMultiple( - args.Path, callback, args.attribute - ) + + resp = self.commThread.subscribeMultiple(args.Path, callback, args.attribute) resJson = json.loads(resp) if "subscriptionId" in resJson: self.subscribeIds.add(resJson["subscriptionId"]) @@ -591,8 +552,8 @@ def connect(self): # Explain were we are connecting to: print( - f"Connecting to VSS server at {config['ip'] } port {config['port'] } \ -using {'KUKSA GRPC' if config['protocol'] == 'grpc' else 'VISS' } protocol." + f"Connecting to VSS server at {config['ip']} port {config['port']} \ +using {'KUKSA GRPC' if config['protocol'] == 'grpc' else 'VISS'} protocol." ) print(f"TLS will {'not be' if config['insecure'] else 'be'} used.") diff --git a/kuksa-client/kuksa_client/cli_backend/ws.py b/kuksa-client/kuksa_client/cli_backend/ws.py index 76b754b..8bc2744 100644 --- a/kuksa-client/kuksa_client/cli_backend/ws.py +++ b/kuksa-client/kuksa_client/cli_backend/ws.py @@ -260,9 +260,10 @@ def subscribe(self, path, callback, attribute="value", timeout=5): return res def subscribeMultiple(self, paths, callback, attribute="value", timeout=5): - raise Exception("Not supported by VISSv2. " - "Try using `subscribe` if you meant to use the " - "`subscribe` function of VISSv2") + responses = [] + for path in paths: + responses.append(json.loads(self.subscribe(path, callback, attribute, timeout))) + return json.dumps(responses[0] if len(responses) == 1 else responses) # Unsubscribe value changes of to a given path. # The subscription id from the response of the corresponding subscription request will be required