|
46 | 46 |
|
47 | 47 | # Callback when connection is accidentally lost.
|
48 | 48 | def on_connection_interrupted(connection, error_code):
|
49 |
| - print("Connection interrupted. error_code:{}".format(error_code)) |
| 49 | + print("Connection interrupted. error_code: {}".format(error_code)) |
| 50 | + |
50 | 51 |
|
51 | 52 | # Callback when an interrupted connection is re-established.
|
52 | 53 | def on_connection_resumed(connection, error_code, session_present):
|
53 |
| - print("Connection resumed. error_code:{} session_present:{}".format(error_code, session_present)) |
| 54 | + print("Connection resumed. error_code: {} session_present: {}".format(error_code, session_present)) |
54 | 55 |
|
55 |
| - if not session_present: |
56 |
| - print("Server did not save state. Resubscribing to existing topics...") |
| 56 | + if not error_code and not session_present: |
| 57 | + print("Session did not persist. Resubscribing to existing topics...") |
57 | 58 | resubscribe_future, _ = connection.resubscribe_existing_topics()
|
58 | 59 |
|
59 | 60 | # Cannot synchronously wait for resubscribe result because we're on the connection's event-loop thread,
|
60 | 61 | # evaluate result with a callback instead.
|
61 |
| - def on_resubscribed(resubscribe_future): |
62 |
| - try: |
63 |
| - resubscribe_results = resubscribe_future.result() |
64 |
| - print("Resubscribe complete. results:{}".format(resubscribe_results)) |
65 |
| - for topic, qos in resubscribe_results['topics']: |
66 |
| - assert(qos) |
67 |
| - except Exception as e: |
68 |
| - print("Resubscribe failed. Exiting", e) |
69 |
| - exit(-1) |
| 62 | + resubscribe_future.add_done_callback(on_resubscribe_complete) |
| 63 | + |
| 64 | + |
| 65 | +def on_resubscribe_complete(resubscribe_future): |
| 66 | + resubscribe_results = resubscribe_future.result() |
| 67 | + print("Resubscribe results: {}".format(resubscribe_results)) |
70 | 68 |
|
71 |
| - resubscribe_future.add_done_callback(on_resubscribed) |
| 69 | + for topic, qos in resubscribe_results['topics']: |
| 70 | + if qos is None: |
| 71 | + exit("Server rejected resubscribe to topic: {}".format(topic)) |
72 | 72 |
|
73 | 73 |
|
74 | 74 | # Callback when the subscribed topic receives a message
|
@@ -121,7 +121,9 @@ def on_message_received(topic, message):
|
121 | 121 | qos=mqtt.QoS.AT_LEAST_ONCE,
|
122 | 122 | callback=on_message_received)
|
123 | 123 |
|
124 |
| - subscribe_future.result() |
| 124 | + subscribe_result = subscribe_future.result() |
| 125 | + if subscribe_result['qos'] is None: |
| 126 | + raise RuntimeError("Server rejected subscription") |
125 | 127 | print("Subscribed!")
|
126 | 128 |
|
127 | 129 | # Publish message to server desired number of times.
|
|
0 commit comments