Skip to content

Commit

Permalink
Added subscriber info in unsubscribe request
Browse files Browse the repository at this point in the history
  • Loading branch information
neelam-kushwah committed Jul 23, 2024
1 parent c29da91 commit b4dac12
Show file tree
Hide file tree
Showing 2 changed files with 6 additions and 1 deletion.
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,7 @@ async def test_subscribe_when_we_try_to_subscribe_to_the_same_topic_twice_with_d
self.assertEqual(self.transport.get_source.call_count, 2)

async def test_unsubscribe_using_mock_rpcclient_and_simplernotifier(self):
self.transport.get_source.return_value = self.source
self.transport.register_listener.return_value = UStatus(code=UCode.OK)
self.transport.unregister_listener.return_value = UStatus(code=UCode.OK)
self.rpc_client.invoke_method.return_value = UPayload.pack(UnsubscribeResponse())
Expand All @@ -289,6 +290,7 @@ async def test_unsubscribe_using_mock_rpcclient_and_simplernotifier(self):
self.transport.unregister_listener.assert_called_once()

async def test_unsubscribe_when_invokemethod_return_an_exception(self):
self.transport.get_source.return_value = self.source
self.transport.register_listener.return_value = UStatus(code=UCode.OK)
self.transport.unregister_listener.return_value = UStatus(code=UCode.OK)
self.rpc_client.invoke_method.return_value = UStatusError.from_code_message(
Expand All @@ -306,6 +308,7 @@ async def test_unsubscribe_when_invokemethod_return_an_exception(self):
self.transport.unregister_listener.assert_not_called()

async def test_unsubscribe_when_invokemethod_returned_ok_but_we_failed_to_unregister_the_listener(self):
self.transport.get_source.return_value = self.source
self.transport.register_listener.return_value = UStatus(code=UCode.OK)
self.transport.unregister_listener.return_value = UStatusError.from_code_message(UCode.ABORTED, "aborted")
self.rpc_client.invoke_method.return_value = UPayload.pack(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -218,7 +218,9 @@ async def unsubscribe(
raise ValueError("Listener missing")
if not options:
raise ValueError("CallOptions missing")
unsubscribe_request = UnsubscribeRequest(topic=topic)
unsubscribe_request = UnsubscribeRequest(
topic=topic, subscriber=SubscriberInfo(uri=self.transport.get_source())
)
future_result = self.rpc_client.invoke_method(self.unsubscribe_uri, UPayload.pack(unsubscribe_request), options)

response = await RpcMapper.map_response_to_result(future_result, UnsubscribeResponse)
Expand Down

0 comments on commit b4dac12

Please sign in to comment.