Closed
Conversation
This code (self._messages[topic] = []) in "subscribe" and "listen" will abandon some messages that would not return to scripts when multiple messages received and multiple “subscribe” or “listen” called in very short time. This modification suites to repair the deficiency above.
Owner
|
Hi @obnjed thanks for the PR. Can you provide a test that reproduces the issue? |
This was referenced Sep 25, 2026
randomsync
added a commit
that referenced
this pull request
Sep 27, 2026
… aliases (#56) * Run a background network loop per connection, with message queues and aliases Each Connect creates a _Connection (src/MQTTLibrary/connection.py) that owns the paho client, runs loop_start() until Disconnect, and keeps one queue per subscription filter plus an unclaimed queue fed by a single on_message installed before connecting. Keywords wait on acknowledgements by message id instead of calling loop(). One condition per connection guards the filter registry, dispatch and the unclaimed-queue drain. Keywords: - Connect takes keepalive= and alias=, uses reconnect_on_failure=False and connect_timeout equal to the loop timeout, and fails with the broker's reason and the host. Connecting again on an open alias disconnects the old connection with a warning. - New Switch Connection and Disconnect All. alias= on every keyword that uses a connection. - Publish waits for its acknowledgement; Subscribe waits for SUBACK and fails on a refused subscription; Unsubscribe removes only its filter. - Listen returns the oldest messages first and keeps the rest. - Subscribe And Validate reads the same queues; its error text is kept. - Publish Single and Publish Multiple accept protocol as a name or number. Tests: a pytest unit layer with a fake paho client, suites moved to tests/acceptance with named connections in the async cases, the limit test rewritten for the new semantics, and regression tests for #23, #24, #25, #28, #33, thread leaks, persistent sessions and aliases. CI runs both layers under coverage with a 90% floor. Fixes #25, #28, #33. * Address review of the connection layer - Overlapping filters: Mosquitto sends one copy per matching subscription, so each filter got every message twice. The first copy now goes to every matching queue and the identical copies that follow it directly are dropped, which is right for brokers that send one copy and for those that send one per subscription. - Stop the network thread without paho's loop_stop(): its join has no timeout, and it raises if the thread clears its reference while ending by itself. Set _thread_terminate and join the thread kept at loop_start() for at most the loop timeout. - A failed Subscribe puts the messages it claimed back into the unclaimed queue, so a retry still gets persistent-session messages. - An ack that arrives after its wait gave up is dropped, so it cannot answer a later operation that reuses the mid. - When paho closes the connection before CONNACK (MQTT 3.1.1 codes 1 and 2 with reconnect_on_failure=False), say that the broker closed it without accepting it, and why that usually happens. - Filter queues are bounded like the unclaimed queue, with a warning on the next read when messages were dropped. - Listen decodes before taking messages and drops only an undecodable one; Subscribe And Validate skips it. - Listen and Subscribe And Validate stop waiting when the connection drops and report why. - _connected is derived from the CONNACK and disconnect state. * Shorten the unit tests' loop timeout to keep the layer well under 3 s * Count restored-message overflow, expire abandoned acks, document the paho internals used - A failed Subscribe that returns messages to a full unclaimed queue now drops the oldest and counts them, like any other overflow. - A late ack is discarded only for ABANDONED_ACK_SECONDS after its wait gave up, so a lost ack cannot block a later operation that reuses the mid after the counter wraps. - Comment on the paho 2.1 attribute _stop() sets, and note that a thread that outlives the join keeps its socket. - CHANGELOG: the single-copy broker limit of the overlap handling.
This pull request was closed.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
This code (self._messages[topic] = []) in "subscribe" and "listen" will abandon some messages that would not return to scripts when multiple messages received and multiple “subscribe” or “listen” called in very short time.
This modification suites to repair the deficiency above.