Skip to content

Commit bac8eea

Browse files
authored
Merge branch 'main' into ct-requests-unixsocket
2 parents 7dc4c9e + e605c14 commit bac8eea

4 files changed

Lines changed: 39 additions & 12 deletions

File tree

‎azure-iot-device/azure/iot/device/common/mqtt_transport.py‎

Lines changed: 21 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,7 @@ def _create_mqtt_client(self):
140140
if self._websockets:
141141
logger.info("Creating client for connecting using MQTT over websockets")
142142
mqtt_client = mqtt.Client(
143+
callback_api_version=mqtt.CallbackAPIVersion.VERSION1,
143144
client_id=self._client_id,
144145
clean_session=False,
145146
protocol=mqtt.MQTTv311,
@@ -149,7 +150,10 @@ def _create_mqtt_client(self):
149150
else:
150151
logger.info("Creating client for connecting using MQTT over TCP")
151152
mqtt_client = mqtt.Client(
152-
client_id=self._client_id, clean_session=False, protocol=mqtt.MQTTv311
153+
callback_api_version=mqtt.CallbackAPIVersion.VERSION1,
154+
client_id=self._client_id,
155+
clean_session=False,
156+
protocol=mqtt.MQTTv311,
153157
)
154158

155159
if self._proxy_options:
@@ -436,7 +440,6 @@ def disconnect(self, clear_inflight=False):
436440
:raises: ConnectionDroppedError in unexpected cases.
437441
:raises: UnauthorizedError in unexpected cases.
438442
:raises: ConnectionFailedError in unexpected cases.
439-
:raises: NoConnectionError if the client isn't actually connected.
440443
"""
441444
logger.info("disconnecting MQTT client")
442445
try:
@@ -452,11 +455,22 @@ def disconnect(self, clear_inflight=False):
452455

453456
logger.debug("_mqtt_client.disconnect returned rc={}".format(rc))
454457
if rc:
455-
# This could result in ConnectionDroppedError or ProtocolClientError
456-
# No matter what, we always raise here to give upper layers a chance to respond
457-
# to this error.
458-
err = _create_error_from_rc_code(rc)
459-
raise err
458+
# Special case: MQTT_ERR_NO_CONN (rc=4) during disconnect means the socket
459+
# is already closed. In Paho 2.x, this can happen even after a successful
460+
# disconnect because the on_disconnect callback fires (with rc=0) before
461+
# disconnect() returns, and Paho's internal cleanup closes the socket.
462+
# Since we wanted to disconnect and we're disconnected, treat this as success.
463+
if rc == mqtt.MQTT_ERR_NO_CONN:
464+
logger.debug(
465+
"disconnect returned MQTT_ERR_NO_CONN - socket already closed, treating as success"
466+
)
467+
# Still clear inflight operations since we're effectively disconnected
468+
if clear_inflight:
469+
self._op_manager.cancel_all_operations()
470+
else:
471+
# This could result in ConnectionDroppedError or ProtocolClientError
472+
err = _create_error_from_rc_code(rc)
473+
raise err
460474
else:
461475
# Clear pending ops if instructed, but only if the disconnect was successful.
462476
# Technically the disconnect could still fail upon response, however that would then

‎setup.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@
7272
],
7373
install_requires=[
7474
"deprecation>=2.1.0,<3.0.0",
75-
"paho-mqtt>=1.6.1,<2.0.0",
75+
"paho-mqtt>=2.0.0,<3.0.0",
7676
"requests>=2.32.3,<3.0.0",
7777
"requests-unixsocket>=0.4.1",
7878
"janus",

‎tests/e2e/iothub_e2e/conftest.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,7 @@ def random_reported_props():
6767
all_objects_can_leak = []
6868
one_object_can_leak = [
6969
"<class 'azure.iot.device.common.alarm.Alarm'>",
70-
"<class 'paho.mqtt.client.WebsocketWrapper'>",
70+
"<class 'paho.mqtt.client._WebsocketWrapper'>",
7171
]
7272

7373

‎tests/unit/common/test_mqtt_transport.py‎

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -124,6 +124,12 @@
124124
},
125125
]
126126

127+
# For disconnect, MQTT_ERR_NO_CONN is treated as success (socket already closed)
128+
# so we exclude it from the error return codes for disconnect tests
129+
disconnect_operation_return_codes = [
130+
x for x in operation_return_codes if x["rc"] != mqtt.MQTT_ERR_NO_CONN
131+
]
132+
127133

128134
@pytest.fixture
129135
def mock_mqtt_client(mocker, fake_paho_thread):
@@ -202,7 +208,10 @@ def test_instantiates_mqtt_client(self, mocker):
202208

203209
assert mock_mqtt_client_constructor.call_count == 1
204210
assert mock_mqtt_client_constructor.call_args == mocker.call(
205-
client_id=fake_device_id, clean_session=False, protocol=mqtt.MQTTv311
211+
callback_api_version=mqtt.CallbackAPIVersion.VERSION1,
212+
client_id=fake_device_id,
213+
clean_session=False,
214+
protocol=mqtt.MQTTv311,
206215
)
207216

208217
@pytest.mark.it(
@@ -221,6 +230,7 @@ def test_configures_mqtt_websockets(self, mocker):
221230

222231
assert mock_mqtt_client_constructor.call_count == 1
223232
assert mock_mqtt_client_constructor.call_args == mocker.call(
233+
callback_api_version=mqtt.CallbackAPIVersion.VERSION1,
224234
client_id=fake_device_id,
225235
clean_session=False,
226236
protocol=mqtt.MQTTv311,
@@ -802,8 +812,11 @@ def test_client_raises_base_exception(
802812
@pytest.mark.it("Raises a custom Exception if Paho disconnect returns a failing rc code")
803813
@pytest.mark.parametrize(
804814
"error_params",
805-
operation_return_codes,
806-
ids=["{}->{}".format(x["name"], x["error"].__name__) for x in operation_return_codes],
815+
disconnect_operation_return_codes,
816+
ids=[
817+
"{}->{}".format(x["name"], x["error"].__name__)
818+
for x in disconnect_operation_return_codes
819+
],
807820
)
808821
def test_client_returns_failing_rc_code(
809822
self, mocker, mock_mqtt_client, transport, error_params

0 commit comments

Comments
 (0)