Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions docs/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
51 changes: 6 additions & 45 deletions kuksa-client/kuksa_client/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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",
Expand Down Expand Up @@ -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"])
Expand Down Expand Up @@ -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.")

Expand Down
14 changes: 11 additions & 3 deletions kuksa-client/kuksa_client/cli_backend/grpc.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
7 changes: 4 additions & 3 deletions kuksa-client/kuksa_client/cli_backend/ws.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
163 changes: 163 additions & 0 deletions kuksa-client/tests/test_cli_backend.py
Original file line number Diff line number Diff line change
@@ -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()
Loading