From 16feb14a57fd0f3351a6791c795c086f69f6978f Mon Sep 17 00:00:00 2001 From: "Roger A. Light" Date: Sun, 9 Oct 2022 22:17:47 +0100 Subject: [PATCH] Persistence interface updates and sqlite plugin fixes The MOSQ_EVT_PERSIST_CLIENT_MSG_CLEAR event has been removed, due to never being called. It is the responsibility of the plugin to remove client subscriptions and client messages when the client is removed. Lots of persist test improvements and additions - mostly checking item counts. --- include/mosquitto_broker.h | 3 +- plugins/persist-sqlite/base_msgs.c | 1 + plugins/persist-sqlite/client_msgs.c | 16 +-- plugins/persist-sqlite/clients.c | 1 + plugins/persist-sqlite/persist_sqlite.h | 2 +- plugins/persist-sqlite/plugin.c | 3 - src/database.c | 11 -- src/mosquitto_broker_internal.h | 1 - src/plugin_callbacks.c | 4 - src/plugin_persist.c | 22 ---- .../broker/15-persist-client-msg-in-v3-1-1.py | 20 ++-- test/broker/15-persist-client-msg-in-v5-0.py | 16 ++- .../15-persist-client-msg-out-clear-v3-1-1.py | 111 ++++++++++++++++++ .../15-persist-client-msg-out-dup-v3-1-1.py | 21 ++-- .../15-persist-client-msg-out-queue-v3-1-1.py | 16 ++- .../15-persist-client-msg-out-v3-1-1-db.py | 15 +-- .../15-persist-client-msg-out-v3-1-1.py | 13 +- test/broker/15-persist-client-msg-out-v5-0.py | 15 +-- test/broker/15-persist-client-v3-1-1.py | 19 ++- test/broker/15-persist-client-v5-0.py | 24 ++-- .../15-persist-publish-properties-v5-0.py | 16 +-- test/broker/15-persist-retain-clear.py | 20 ++-- test/broker/15-persist-retain-v3-1-1.py | 13 +- test/broker/15-persist-retain-v5-0.py | 13 +- test/broker/15-persist-subscription-v3-1-1.py | 31 +++-- test/broker/15-persist-subscription-v5-0.py | 19 ++- test/broker/Makefile | 1 + test/broker/persist_sqlite.py | 27 +++-- test/broker/test.py | 31 ++--- test/mosq_test.py | 11 ++ 30 files changed, 310 insertions(+), 206 deletions(-) create mode 100755 test/broker/15-persist-client-msg-out-clear-v3-1-1.py diff --git a/include/mosquitto_broker.h b/include/mosquitto_broker.h index 1ce1cac6..dac728c1 100644 --- a/include/mosquitto_broker.h +++ b/include/mosquitto_broker.h @@ -94,8 +94,7 @@ enum mosquitto_plugin_event { MOSQ_EVT_PERSIST_CLIENT_MSG_ADD = 26, MOSQ_EVT_PERSIST_CLIENT_MSG_DELETE = 27, MOSQ_EVT_PERSIST_CLIENT_MSG_UPDATE = 28, - MOSQ_EVT_PERSIST_CLIENT_MSG_CLEAR = 29, - MOSQ_EVT_PERSIST_CLIENT_MSG_LOAD = 30, + MOSQ_EVT_PERSIST_CLIENT_MSG_LOAD = 29, }; /* Data for the MOSQ_EVT_RELOAD event */ diff --git a/plugins/persist-sqlite/base_msgs.c b/plugins/persist-sqlite/base_msgs.c index bf4abdce..c3147406 100644 --- a/plugins/persist-sqlite/base_msgs.c +++ b/plugins/persist-sqlite/base_msgs.c @@ -18,6 +18,7 @@ Contributors: #include #include +#include #include #include diff --git a/plugins/persist-sqlite/client_msgs.c b/plugins/persist-sqlite/client_msgs.c index 22736be9..a6b2f51d 100644 --- a/plugins/persist-sqlite/client_msgs.c +++ b/plugins/persist-sqlite/client_msgs.c @@ -111,16 +111,12 @@ int persist_sqlite__client_msg_update_cb(int event, void *event_data, void *user } -int persist_sqlite__client_msg_clear_cb(int event, void *event_data, void *userdata) +int persist_sqlite__client_msg_clear(struct mosquitto_sqlite *ms, const char *client_id, int direction) { - struct mosquitto_evt_persist_client_msg *ed = event_data; - struct mosquitto_sqlite *ms = userdata; - int rc = 1; + int rc = MOSQ_ERR_UNKNOWN; - UNUSED(event); - - if(ed->direction == mosq_bmd_all){ - if(sqlite3_bind_text(ms->client_msg_clear_all_stmt, 1, ed->client_id, (int)strlen(ed->client_id), SQLITE_STATIC) == SQLITE_OK){ + if(direction == mosq_bmd_all){ + if(sqlite3_bind_text(ms->client_msg_clear_all_stmt, 1, client_id, (int)strlen(client_id), SQLITE_STATIC) == SQLITE_OK){ ms->event_count++; rc = sqlite3_step(ms->client_msg_clear_all_stmt); if(rc == SQLITE_DONE){ @@ -131,8 +127,8 @@ int persist_sqlite__client_msg_clear_cb(int event, void *event_data, void *userd } sqlite3_reset(ms->client_msg_clear_all_stmt); }else{ - if(sqlite3_bind_text(ms->client_msg_clear_stmt, 1, ed->client_id, (int)strlen(ed->client_id), SQLITE_STATIC) == SQLITE_OK - && sqlite3_bind_int64(ms->client_msg_clear_stmt, 2, ed->direction) == SQLITE_OK + if(sqlite3_bind_text(ms->client_msg_clear_stmt, 1, client_id, (int)strlen(client_id), SQLITE_STATIC) == SQLITE_OK + && sqlite3_bind_int64(ms->client_msg_clear_stmt, 2, direction) == SQLITE_OK ){ ms->event_count++; diff --git a/plugins/persist-sqlite/clients.c b/plugins/persist-sqlite/clients.c index 8c8efa33..b8809831 100644 --- a/plugins/persist-sqlite/clients.c +++ b/plugins/persist-sqlite/clients.c @@ -101,6 +101,7 @@ int persist_sqlite__client_remove_cb(int event, void *event_data, void *userdata rc = MOSQ_ERR_UNKNOWN; } } + persist_sqlite__client_msg_clear(ms, ed->client_id, mosq_bmd_all); return rc; } diff --git a/plugins/persist-sqlite/persist_sqlite.h b/plugins/persist-sqlite/persist_sqlite.h index 6568a34d..76215d97 100644 --- a/plugins/persist-sqlite/persist_sqlite.h +++ b/plugins/persist-sqlite/persist_sqlite.h @@ -61,7 +61,7 @@ int persist_sqlite__client_add_cb(int event, void *event_data, void *userdata); int persist_sqlite__client_update_cb(int event, void *event_data, void *userdata); int persist_sqlite__client_remove_cb(int event, void *event_data, void *userdata); int persist_sqlite__client_msg_add_cb(int event, void *event_data, void *userdata); -int persist_sqlite__client_msg_clear_cb(int event, void *event_data, void *userdata); +int persist_sqlite__client_msg_clear(struct mosquitto_sqlite *ms, const char *client_id, int direction); int persist_sqlite__client_msg_remove_cb(int event, void *event_data, void *userdata); int persist_sqlite__client_msg_update_cb(int event, void *event_data, void *userdata); int persist_sqlite__base_msg_add_cb(int event, void *event_data, void *userdata); diff --git a/plugins/persist-sqlite/plugin.c b/plugins/persist-sqlite/plugin.c index 44d166b1..55898df8 100644 --- a/plugins/persist-sqlite/plugin.c +++ b/plugins/persist-sqlite/plugin.c @@ -167,8 +167,6 @@ int mosquitto_plugin_init(mosquitto_plugin_id_t *identifier, void **user_data, s if(rc) goto fail; rc = mosquitto_callback_register(plg_id, MOSQ_EVT_PERSIST_CLIENT_MSG_UPDATE, persist_sqlite__client_msg_update_cb, NULL, &plg_data); if(rc) goto fail; - rc = mosquitto_callback_register(plg_id, MOSQ_EVT_PERSIST_CLIENT_MSG_CLEAR, persist_sqlite__client_msg_clear_cb, NULL, &plg_data); - if(rc) goto fail; rc = mosquitto_callback_register(plg_id, MOSQ_EVT_TICK, persist_sqlite__tick_cb, NULL, &plg_data); if(rc) goto fail; @@ -205,7 +203,6 @@ int mosquitto_plugin_cleanup(void *user_data, struct mosquitto_opt *options, int mosquitto_callback_unregister(plg_id, MOSQ_EVT_PERSIST_CLIENT_MSG_ADD, persist_sqlite__client_msg_add_cb, NULL); mosquitto_callback_unregister(plg_id, MOSQ_EVT_PERSIST_CLIENT_MSG_DELETE, persist_sqlite__client_msg_remove_cb, NULL); mosquitto_callback_unregister(plg_id, MOSQ_EVT_PERSIST_CLIENT_MSG_UPDATE, persist_sqlite__client_msg_update_cb, NULL); - mosquitto_callback_unregister(plg_id, MOSQ_EVT_PERSIST_CLIENT_MSG_CLEAR, persist_sqlite__client_msg_clear_cb, NULL); mosquitto_callback_unregister(plg_id, MOSQ_EVT_TICK, persist_sqlite__tick_cb, NULL); } diff --git a/src/database.c b/src/database.c index 24e2fabd..610c8903 100644 --- a/src/database.c +++ b/src/database.c @@ -780,27 +780,16 @@ int db__messages_delete_outgoing(struct mosquitto *context) int db__messages_delete(struct mosquitto *context, bool force_free) { - bool clear_incoming = false, clear_outgoing = false; if(!context) return MOSQ_ERR_INVAL; if(force_free || context->clean_start || (context->bridge && context->bridge->clean_start)){ db__messages_delete_incoming(context); - clear_incoming = true; } if(force_free || (context->bridge && context->bridge->clean_start_local) || (context->bridge == NULL && context->clean_start)){ db__messages_delete_outgoing(context); - clear_outgoing = true; - } - - if(clear_incoming && clear_outgoing){ - plugin_persist__handle_client_msg_clear(context, mosq_bmd_all); - }else if(clear_incoming){ - plugin_persist__handle_client_msg_clear(context, mosq_bmd_in); - }else if(clear_outgoing){ - plugin_persist__handle_client_msg_clear(context, mosq_bmd_out); } return MOSQ_ERR_SUCCESS; diff --git a/src/mosquitto_broker_internal.h b/src/mosquitto_broker_internal.h index b40a4132..729d0d88 100644 --- a/src/mosquitto_broker_internal.h +++ b/src/mosquitto_broker_internal.h @@ -875,7 +875,6 @@ void plugin_persist__handle_subscription_delete(struct mosquitto *context, const void plugin_persist__handle_client_msg_add(struct mosquitto *context, const struct mosquitto_client_msg *cmsg); void plugin_persist__handle_client_msg_delete(struct mosquitto *context, const struct mosquitto_client_msg *cmsg); void plugin_persist__handle_client_msg_update(struct mosquitto *context, const struct mosquitto_client_msg *cmsg); -void plugin_persist__handle_client_msg_clear(struct mosquitto *context, uint8_t direction); void plugin_persist__handle_base_msg_add(struct mosquitto_base_msg *base_msg); void plugin_persist__handle_base_msg_delete(struct mosquitto_base_msg *base_msg); void plugin_persist__handle_retain_msg_set(struct mosquitto_base_msg *base_msg); diff --git a/src/plugin_callbacks.c b/src/plugin_callbacks.c index 210bb0fa..dbecee1b 100644 --- a/src/plugin_callbacks.c +++ b/src/plugin_callbacks.c @@ -79,8 +79,6 @@ static const char *get_event_name(int event) return "persist-client-msg-delete"; case MOSQ_EVT_PERSIST_CLIENT_MSG_UPDATE: return "persist-client-msg-update"; - case MOSQ_EVT_PERSIST_CLIENT_MSG_CLEAR: - return "persist-client-msg-clear"; case MOSQ_EVT_PERSIST_CLIENT_MSG_LOAD: return "persist-client-msg-load"; default: @@ -147,8 +145,6 @@ static struct mosquitto__callback **plugin__get_callback_base(struct mosquitto__ return &security_options->plugin_callbacks.persist_client_msg_delete; case MOSQ_EVT_PERSIST_CLIENT_MSG_UPDATE: return &security_options->plugin_callbacks.persist_client_msg_update; - case MOSQ_EVT_PERSIST_CLIENT_MSG_CLEAR: - return &security_options->plugin_callbacks.persist_client_msg_clear; case MOSQ_EVT_PERSIST_BASE_MSG_ADD: return &security_options->plugin_callbacks.persist_base_msg_add; case MOSQ_EVT_PERSIST_BASE_MSG_DELETE: diff --git a/src/plugin_persist.c b/src/plugin_persist.c index 2b8e0d31..0966464f 100644 --- a/src/plugin_persist.c +++ b/src/plugin_persist.c @@ -88,8 +88,6 @@ void plugin_persist__handle_client_update(struct mosquitto *context) struct mosquitto__security_options *opts; struct mosquitto_message_v5 will; - if(db.shutdown) return; - UNUSED(will); /* FIXME */ if(db.shutdown) return; @@ -276,26 +274,6 @@ void plugin_persist__handle_client_msg_update(struct mosquitto *context, const s } -void plugin_persist__handle_client_msg_clear(struct mosquitto *context, uint8_t direction) -{ - struct mosquitto_evt_persist_client_msg event_data; - struct mosquitto__callback *cb_base; - struct mosquitto__security_options *opts; - - if(db.shutdown || context->is_persisted == false) return; - - opts = &db.config->security_options; - memset(&event_data, 0, sizeof(event_data)); - - event_data.client_id = context->id; - event_data.direction = direction; - - DL_FOREACH(opts->plugin_callbacks.persist_client_msg_clear, cb_base){ - cb_base->cb(MOSQ_EVT_PERSIST_CLIENT_MSG_CLEAR, &event_data, cb_base->userdata); - } -} - - void plugin_persist__handle_base_msg_add(struct mosquitto_base_msg *msg) { struct mosquitto_evt_persist_base_msg event_data; diff --git a/test/broker/15-persist-client-msg-in-v3-1-1.py b/test/broker/15-persist-client-msg-in-v3-1-1.py index f18ebce7..a6e9fbf4 100755 --- a/test/broker/15-persist-client-msg-in-v3-1-1.py +++ b/test/broker/15-persist-client-msg-in-v3-1-1.py @@ -7,12 +7,12 @@ from mosq_test_helper import * persist_help = persist_module() port = mosq_test.get_port() +persist_help.init(port) conf_file = os.path.basename(__file__).replace('.py', '.conf') persist_help.write_config(conf_file, port) rc = 1 -persist_help.init(port) client_id = "persist-client-msg-in-v3-1-1" proto_ver = 4 @@ -20,6 +20,7 @@ proto_ver = 4 helper_id = "persist-client-msg-in-v3-1-1-helper" topic = "client-msg-in/2" qos = 2 +stde = b"" connect_packet = mosq_test.gen_connect(client_id, proto_ver=proto_ver, clean_session=False) connack_packet1 = mosq_test.gen_connack(rc=0, proto_ver=proto_ver) @@ -52,11 +53,10 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port, clients=1, client_msgs_in=2, base_msgs=2) # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -78,6 +78,12 @@ try: helper.send(pubrec2_packet) mosq_test.do_receive_send(helper, pubrel2_packet, pubcomp2_packet, "pubcomp2 receive") + # Kill broker + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port, clients=1) + rc = broker_terminate_rc finally: if broker is not None: @@ -85,7 +91,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-client-msg-in-v5-0.py b/test/broker/15-persist-client-msg-in-v5-0.py index 2118a9b2..39a4765d 100755 --- a/test/broker/15-persist-client-msg-in-v5-0.py +++ b/test/broker/15-persist-client-msg-in-v5-0.py @@ -53,11 +53,9 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + broker_terminate_rc = mosq_test.terminate_broker(broker) + + persist_help.check_counts(port, clients=1, client_msgs_in=2, base_msgs=2) # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -79,6 +77,12 @@ try: helper.send(pubrec2_packet) mosq_test.do_receive_send(helper, pubrel2_packet, pubcomp2_packet, "pubcomp2 receive") + # Kill broker + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port, clients=1) + rc = broker_terminate_rc finally: if broker is not None: @@ -86,7 +90,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-client-msg-out-clear-v3-1-1.py b/test/broker/15-persist-client-msg-out-clear-v3-1-1.py new file mode 100755 index 00000000..01332ac5 --- /dev/null +++ b/test/broker/15-persist-client-msg-out-clear-v3-1-1.py @@ -0,0 +1,111 @@ +#!/usr/bin/env python3 + +# Connect a client, add a subscription, disconnect, send a message with a +# different client, restore, reconnect with clear start, check it is not received. + +from mosq_test_helper import * +persist_help = persist_module() + +port = mosq_test.get_port() +conf_file = os.path.basename(__file__).replace('.py', '.conf') +persist_help.write_config(conf_file, port) + +rc = 1 + +persist_help.init(port) + +keepalive = 10 +client_id = "persist-client-msg-v3-1-1" +proto_ver = 4 + +helper_id = "persist-client-msg-v3-1-1-helper" +topic0 = "client-msg/0" +topic1 = "client-msg/1" +topic2 = "client-msg/2" + +connect_packet = mosq_test.gen_connect(client_id, keepalive=keepalive, proto_ver=proto_ver, clean_session=False) +connect_packet_clear = mosq_test.gen_connect(client_id, keepalive=keepalive, proto_ver=proto_ver, clean_session=True) +connack_packet1 = mosq_test.gen_connack(rc=0, proto_ver=proto_ver) +connack_packet2 = mosq_test.gen_connack(rc=0, flags=1, proto_ver=proto_ver) +mid = 1 +subscribe_packet0 = mosq_test.gen_subscribe(mid, topic0, qos=0, proto_ver=proto_ver) +suback_packet0 = mosq_test.gen_suback(mid=mid, qos=0, proto_ver=proto_ver) +subscribe_packet1 = mosq_test.gen_subscribe(mid, topic1, qos=1, proto_ver=proto_ver) +suback_packet1 = mosq_test.gen_suback(mid=mid, qos=1, proto_ver=proto_ver) +subscribe_packet2 = mosq_test.gen_subscribe(mid, topic2, qos=2, proto_ver=proto_ver) +suback_packet2 = mosq_test.gen_suback(mid=mid, qos=2, proto_ver=proto_ver) + +connect_packet_helper = mosq_test.gen_connect(helper_id, keepalive=keepalive, proto_ver=proto_ver, clean_session=True) +publish_packet0 = mosq_test.gen_publish(topic=topic0, qos=0, payload="message", proto_ver=proto_ver) +mid = 1 +publish_packet1 = mosq_test.gen_publish(topic=topic1, qos=1, payload="message", mid=mid, proto_ver=proto_ver) +puback_packet = mosq_test.gen_puback(mid=mid, proto_ver=proto_ver) +mid = 2 +publish_packet2 = mosq_test.gen_publish(topic=topic2, qos=2, payload="message", mid=mid, proto_ver=proto_ver) +pubrec_packet = mosq_test.gen_pubrec(mid=mid, proto_ver=proto_ver) +pubrel_packet = mosq_test.gen_pubrel(mid=mid, proto_ver=proto_ver) +pubcomp_packet = mosq_test.gen_pubcomp(mid=mid, proto_ver=proto_ver) + +broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) + +con = None +try: + # Connect client, subscribe, disconnect + sock = mosq_test.do_client_connect(connect_packet, connack_packet1, timeout=5, port=port) + mosq_test.do_send_receive(sock, subscribe_packet0, suback_packet0, "suback 0") + mosq_test.do_send_receive(sock, subscribe_packet1, suback_packet1, "suback 1") + mosq_test.do_send_receive(sock, subscribe_packet2, suback_packet2, "suback 2") + sock.close() + + # Connect helper and publish + helper = mosq_test.do_client_connect(connect_packet_helper, connack_packet1, timeout=5, port=port) + helper.send(publish_packet0) + mosq_test.do_send_receive(helper, publish_packet1, puback_packet, "puback helper") + mosq_test.do_send_receive(helper, publish_packet2, pubrec_packet, "pubrec helper") + mosq_test.do_send_receive(helper, pubrel_packet, pubcomp_packet, "pubcomp helper") + helper.close() + + # Kill broker + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port, clients=1, client_msgs_out=2, base_msgs=2, subscriptions=3) + + # Restart broker + broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) + + # Connect client again, it should have a session + sock = mosq_test.do_client_connect(connect_packet, connack_packet2, timeout=5, port=port) + + # Does the client get the messages - don't complete the flows + mosq_test.expect_packet(sock, "publish 1", publish_packet1) + mosq_test.expect_packet(sock, "publish 2", publish_packet2) + sock.close() + + # Connect client again and clear the session + sock = mosq_test.do_client_connect(connect_packet_clear, connack_packet1, timeout=5, port=port) + # If there are messages, the ping will fail + mosq_test.do_ping(sock) + + # Kill broker + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port) + + rc = broker_terminate_rc +finally: + if broker is not None: + broker.terminate() + if mosq_test.wait_for_subprocess(broker): + print("broker not terminated") + if rc == 0: rc=1 + (_, stde) = broker.communicate() + os.remove(conf_file) + rc += persist_help.cleanup(port) + + if rc: + print(stde.decode('utf-8')) + + +exit(rc) diff --git a/test/broker/15-persist-client-msg-out-dup-v3-1-1.py b/test/broker/15-persist-client-msg-out-dup-v3-1-1.py index 81adf1ae..6fced226 100755 --- a/test/broker/15-persist-client-msg-out-dup-v3-1-1.py +++ b/test/broker/15-persist-client-msg-out-dup-v3-1-1.py @@ -51,7 +51,7 @@ try: mosq_test.do_send_receive(sock, subscribe_packet, suback_packet, "suback") sock.close() - #persist_help.check_counts(port, clients=1, client_msgs=0, base_msgs=0, retains=0, subscriptions=1) + #persist_help.check_counts(port, clients=1, subscriptions=1) # Helper - send message then disconnect sock = mosq_test.do_client_connect(connect2_packet, connack2_packet, timeout=5, port=port) @@ -59,13 +59,13 @@ try: mosq_test.do_send_receive(sock, pubrel_packet, pubcomp_packet, "pubcomp") sock.close() - #persist_help.check_counts(port, clients=1, client_msgs=1, base_msgs=1, retains=0, subscriptions=1) + #persist_help.check_counts(port, clients=1, client_msgs=1, base_msgs=1, subscriptions=1) # Reconnect, receive publish, disconnect sock = mosq_test.do_client_connect(connect1_packet, connack1_packet2, timeout=5, port=port) mosq_test.expect_packet(sock, "publish 1", publish_packet_r1) - #persist_help.check_counts(port, clients=1, client_msgs=1, base_msgs=1, retains=0, subscriptions=1) + #persist_help.check_counts(port, clients=1, client_msgs=1, base_msgs=1, subscriptions=1) # Reconnect, receive publish, disconnect - dup should now be set sock = mosq_test.do_client_connect(connect1_packet, connack1_packet2, timeout=5, port=port) @@ -74,24 +74,19 @@ try: #con.close() #con = None - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 - (stdo, stde) = broker.communicate() + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) broker = None - persist_help.check_counts(port, clients=1, client_msgs=1, base_msgs=1, retains=0, subscriptions=1) + persist_help.check_counts(port, clients=1, client_msgs_out=1, base_msgs=1, subscriptions=1) # Check client - persist_help.check_client(port, client_id, None, 0, 0, port, 0, 2, 1, -1, 0) + persist_help.check_client(port, client_id, None, 0, 0, port, 0, 2, 1, 4294967295, 0) # Check subscription persist_help.check_subscription(port, client_id, topic, qos, 0) # Check stored message - store_id = persist_help.check_store_msg(port, 0, topic, payload_b, source_id, None, len(payload_b), source_mid, port, qos, 0) + store_id = persist_help.check_base_msg(port, 0, topic, payload_b, source_id, None, len(payload_b), source_mid, port, qos, 0) # Check client msg persist_help.check_client_msg(port, client_id, store_id, 1, persist_help.dir_out, 1, qos, 0, persist_help.ms_wait_for_pubrec) @@ -103,7 +98,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() if con is not None: con.close() os.remove(conf_file) diff --git a/test/broker/15-persist-client-msg-out-queue-v3-1-1.py b/test/broker/15-persist-client-msg-out-queue-v3-1-1.py index 8e51807f..5a979e40 100755 --- a/test/broker/15-persist-client-msg-out-queue-v3-1-1.py +++ b/test/broker/15-persist-client-msg-out-queue-v3-1-1.py @@ -66,11 +66,10 @@ try: helper.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port, clients=1, client_msgs_out=2, base_msgs=2, subscriptions=3) # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -89,6 +88,11 @@ try: # If there are messages, the ping will fail mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port, clients=1, subscriptions=3) + rc = broker_terminate_rc finally: if broker is not None: @@ -96,7 +100,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-client-msg-out-v3-1-1-db.py b/test/broker/15-persist-client-msg-out-v3-1-1-db.py index 61b8a4a4..e6f6df29 100755 --- a/test/broker/15-persist-client-msg-out-v3-1-1-db.py +++ b/test/broker/15-persist-client-msg-out-v3-1-1-db.py @@ -48,24 +48,19 @@ try: mosq_test.do_send_receive(sock, publish_packet, puback_packet, "puback") sock.close() - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 - (stdo, stde) = broker.communicate() + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) broker = None - persist_help.check_counts(port, clients=1, client_msgs=1, base_msgs=1, retains=0, subscriptions=1) + persist_help.check_counts(port, clients=1, client_msgs_out=1, base_msgs=1, subscriptions=1) # Check client - persist_help.check_client(port, client_id, None, 0, 0, port, 0, 2, 1, -1, 0) + persist_help.check_client(port, client_id, None, 0, 0, port, 0, 2, 1, 4294967295, 0) # Check subscription persist_help.check_subscription(port, client_id, topic, qos, 0) # Check stored message - store_id = persist_help.check_store_msg(port, 0, topic, payload_b, source_id, None, len(payload_b), mid, port, qos, 0) + store_id = persist_help.check_base_msg(port, 0, topic, payload_b, source_id, None, len(payload_b), mid, port, qos, 0) # Check client msg persist_help.check_client_msg(port, client_id, store_id, 0, persist_help.dir_out, 1, qos, 0, persist_help.ms_queued) @@ -77,7 +72,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-client-msg-out-v3-1-1.py b/test/broker/15-persist-client-msg-out-v3-1-1.py index d8a1a393..266d4c21 100755 --- a/test/broker/15-persist-client-msg-out-v3-1-1.py +++ b/test/broker/15-persist-client-msg-out-v3-1-1.py @@ -65,11 +65,8 @@ try: helper.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -88,6 +85,10 @@ try: # If there are messages, the ping will fail mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port, clients=1, subscriptions=3) + rc = broker_terminate_rc finally: if broker is not None: @@ -95,7 +96,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-client-msg-out-v5-0.py b/test/broker/15-persist-client-msg-out-v5-0.py index 9dbda14c..b86c16cf 100755 --- a/test/broker/15-persist-client-msg-out-v5-0.py +++ b/test/broker/15-persist-client-msg-out-v5-0.py @@ -45,7 +45,6 @@ pubrec_packet = mosq_test.gen_pubrec(mid=mid, proto_ver=proto_ver) pubrel_packet = mosq_test.gen_pubrel(mid=mid, proto_ver=proto_ver) pubcomp_packet = mosq_test.gen_pubcomp(mid=mid, proto_ver=proto_ver) - broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) con = None @@ -66,11 +65,8 @@ try: helper.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -89,6 +85,11 @@ try: # If there are messages, the ping will fail mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + + persist_help.check_counts(port, clients=1, subscriptions=3) + rc = broker_terminate_rc finally: if broker is not None: @@ -96,7 +97,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-client-v3-1-1.py b/test/broker/15-persist-client-v3-1-1.py index 41098bd5..d529773f 100755 --- a/test/broker/15-persist-client-v3-1-1.py +++ b/test/broker/15-persist-client-v3-1-1.py @@ -37,11 +37,8 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -62,10 +59,8 @@ try: sock.close() # Kill broker - broker.terminate() - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated (2)") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -75,6 +70,10 @@ try: mosq_test.do_ping(sock) sock.close() + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port) + rc = broker_terminate_rc finally: if broker is not None: @@ -82,7 +81,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (3)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-client-v5-0.py b/test/broker/15-persist-client-v5-0.py index 9b67999e..605761c3 100755 --- a/test/broker/15-persist-client-v5-0.py +++ b/test/broker/15-persist-client-v5-0.py @@ -35,14 +35,11 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + broker_terminate_rc = mosq_test.terminate_broker(broker) - persist_help.check_counts(port, clients=1, client_msgs=0, base_msgs=0, retains=0, subscriptions=0) - persist_help.check_client(port, "persist-client-v5-0", None, 0, 1, port, 10000, 2, 1, 60, 0) + persist_help.check_counts(port, clients=1) + # FIXME - port persist_help.check_client(port, "persist-client-v5-0", None, 0, 1, port, 10000, 2, 1, 60, 0) + persist_help.check_client(port, "persist-client-v5-0", None, 0, 1, None, 10000, 2, 1, 60, 0) # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -63,12 +60,10 @@ try: sock.close() # Kill broker - broker.terminate() - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None - persist_help.check_counts(port, clients=0, client_msgs=0, base_msgs=0, retains=0, subscriptions=0) + persist_help.check_counts(port) # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -78,6 +73,9 @@ try: mosq_test.do_ping(sock) sock.close() + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port) rc = broker_terminate_rc finally: @@ -86,7 +84,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-publish-properties-v5-0.py b/test/broker/15-persist-publish-properties-v5-0.py index 35b3f707..ffa00100 100755 --- a/test/broker/15-persist-publish-properties-v5-0.py +++ b/test/broker/15-persist-publish-properties-v5-0.py @@ -23,11 +23,12 @@ connack_packet = mosq_test.gen_connack(rc=0, proto_ver=proto_ver) props = mqtt5_props.gen_byte_prop(mqtt5_props.PROP_PAYLOAD_FORMAT_INDICATOR, 1) props += mqtt5_props.gen_string_prop(mqtt5_props.PROP_CONTENT_TYPE, "plain/text") props += mqtt5_props.gen_string_prop(mqtt5_props.PROP_RESPONSE_TOPIC, "/dev/null") -#props += mqtt5_props.gen_string_prop(mqtt5_props.PROP_CORRELATION_DATA, "2357289375902345") +props += mqtt5_props.gen_string_prop(mqtt5_props.PROP_CORRELATION_DATA, "2357289375902345") props += mqtt5_props.gen_string_pair_prop(mqtt5_props.PROP_USER_PROPERTY, "name", "value4") props += mqtt5_props.gen_string_pair_prop(mqtt5_props.PROP_USER_PROPERTY, "name", "value3") props += mqtt5_props.gen_string_pair_prop(mqtt5_props.PROP_USER_PROPERTY, "name", "value2") props += mqtt5_props.gen_string_pair_prop(mqtt5_props.PROP_USER_PROPERTY, "name", "value1") +#props += mqtt5_props.gen_uint32_prop(mqtt5_props.PROP_MESSAGE_EXPIRY_INTERVAL, 60) publish_packet = mosq_test.gen_publish(topic, qos=qos, payload="retained message 1", retain=True, proto_ver=proto_ver, properties=props) mid = 1 @@ -56,11 +57,8 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -73,6 +71,10 @@ try: mosq_test.expect_packet(sock, "publish", publish_packet) mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port, base_msgs=1, retain_msgs=1) + rc = broker_terminate_rc finally: if broker is not None: @@ -80,7 +82,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-retain-clear.py b/test/broker/15-persist-retain-clear.py index ce20a663..6eb93790 100755 --- a/test/broker/15-persist-retain-clear.py +++ b/test/broker/15-persist-retain-clear.py @@ -57,11 +57,8 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -78,11 +75,8 @@ try: mosq_test.do_send_receive(sock, publish2_clear_packet, publish2_clear_echo, "clear retain flag") # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -96,6 +90,10 @@ try: # If the other retained message still exists, this will cause an error mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port, base_msgs=1, retain_msgs=1) + rc = broker_terminate_rc finally: if broker is not None: @@ -103,7 +101,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-retain-v3-1-1.py b/test/broker/15-persist-retain-v3-1-1.py index 191b763e..c3088a37 100755 --- a/test/broker/15-persist-retain-v3-1-1.py +++ b/test/broker/15-persist-retain-v3-1-1.py @@ -58,11 +58,8 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -75,6 +72,10 @@ try: mosq_test.receive_unordered(sock, publish1_packet, publish3_packet, "publish 1 / 3") mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port, base_msgs=2, retain_msgs=2) + rc = broker_terminate_rc finally: if broker is not None: @@ -82,7 +83,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-retain-v5-0.py b/test/broker/15-persist-retain-v5-0.py index f520a205..e1e82fcb 100755 --- a/test/broker/15-persist-retain-v5-0.py +++ b/test/broker/15-persist-retain-v5-0.py @@ -58,11 +58,8 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -75,6 +72,10 @@ try: mosq_test.receive_unordered(sock, publish1_packet, publish3_packet, "publish 1 / 3") mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port, base_msgs=2, retain_msgs=2) + rc = broker_terminate_rc finally: if broker is not None: @@ -82,7 +83,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (2)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-subscription-v3-1-1.py b/test/broker/15-persist-subscription-v3-1-1.py index 899fd21e..10961262 100755 --- a/test/broker/15-persist-subscription-v3-1-1.py +++ b/test/broker/15-persist-subscription-v3-1-1.py @@ -35,6 +35,7 @@ topic2 = "subscription/2" packets = {} packets["connect"] = mosq_test.gen_connect(client_id, proto_ver=proto_ver, clean_session=False) +packets["connect_clear"] = mosq_test.gen_connect(client_id, proto_ver=proto_ver, clean_session=True) packets["connack1"] = mosq_test.gen_connack(rc=0, proto_ver=proto_ver) packets["connack2"] = mosq_test.gen_connack(rc=0, flags=1, proto_ver=proto_ver) mid = 1 @@ -71,11 +72,8 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -97,10 +95,8 @@ try: sock.close() # Kill broker - broker.terminate() - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated (2)") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -117,6 +113,21 @@ try: mosq_test.do_receive_send(sock, packets["publish1"], packets["puback1"], "publish 1") mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port, clients=1, subscriptions=2) + + # Restart broker + broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) + + # Connect client again, but clear the session + sock = mosq_test.do_client_connect(packets["connect_clear"], packets["connack1"], timeout=5, port=port) + mosq_test.do_ping(sock) + + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port) + rc = broker_terminate_rc finally: if broker is not None: @@ -124,7 +135,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (3)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/15-persist-subscription-v5-0.py b/test/broker/15-persist-subscription-v5-0.py index 53c6840d..b55e3ddc 100755 --- a/test/broker/15-persist-subscription-v5-0.py +++ b/test/broker/15-persist-subscription-v5-0.py @@ -82,11 +82,8 @@ try: sock.close() # Kill broker - broker.terminate() - broker_terminate_rc = 0 - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -108,10 +105,8 @@ try: sock.close() # Kill broker - broker.terminate() - if mosq_test.wait_for_subprocess(broker): - print("broker not terminated (2)") - broker_terminate_rc = 1 + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None # Restart broker broker = mosq_test.start_broker(filename=os.path.basename(__file__), use_conf=True, port=port) @@ -128,6 +123,10 @@ try: mosq_test.do_receive_send(sock, packets["publish1"], packets["puback1"], "publish 1") mosq_test.do_ping(sock) + (broker_terminate_rc, stde) = mosq_test.terminate_broker(broker) + broker = None + persist_help.check_counts(port, clients=1, subscriptions=2) + rc = broker_terminate_rc finally: if broker is not None: @@ -135,7 +134,7 @@ finally: if mosq_test.wait_for_subprocess(broker): print("broker not terminated (3)") if rc == 0: rc=1 - (stdo, stde) = broker.communicate() + (_, stde) = broker.communicate() os.remove(conf_file) rc += persist_help.cleanup(port) diff --git a/test/broker/Makefile b/test/broker/Makefile index 5ed3977b..d5cc5893 100644 --- a/test/broker/Makefile +++ b/test/broker/Makefile @@ -261,6 +261,7 @@ endif 15 : ./15-persist-client-msg-in-v3-1-1.py persist_sqlite ./15-persist-client-msg-in-v5-0.py persist_sqlite + ./15-persist-client-msg-out-clear-v3-1-1.py persist_sqlite ./15-persist-client-msg-out-dup-v3-1-1.py persist_sqlite ./15-persist-client-msg-out-queue-v3-1-1.py persist_sqlite ./15-persist-client-msg-out-v3-1-1-db.py persist_sqlite diff --git a/test/broker/persist_sqlite.py b/test/broker/persist_sqlite.py index b3b3b607..2b4d5777 100755 --- a/test/broker/persist_sqlite.py +++ b/test/broker/persist_sqlite.py @@ -60,7 +60,7 @@ def cleanup(port): return rc -def check_counts(port, clients, client_msgs, base_msgs, retains, subscriptions): +def check_counts(port, clients=0, client_msgs_in=0, client_msgs_out=0, base_msgs=0, retain_msgs=0, subscriptions=0): con = sqlite3.connect(f"{port}/mosquitto.sqlite3") cur = con.cursor() cur.execute('SELECT COUNT(*) FROM clients') @@ -68,10 +68,15 @@ def check_counts(port, clients, client_msgs, base_msgs, retains, subscriptions): if row[0] != clients: raise ValueError("Found %d clients, expected %d" % (row[0], clients)) - cur.execute('SELECT COUNT(*) FROM client_msgs') + cur.execute('SELECT COUNT(*) FROM client_msgs WHERE direction=0') row = cur.fetchone() - if row[0] != client_msgs: - raise ValueError("Found %d client_msgs, expected %d" % (row[0], client_msgs)) + if row[0] != client_msgs_in: + raise ValueError("Found %d client_msgs_in, expected %d" % (row[0], client_msgs_in)) + + cur.execute('SELECT COUNT(*) FROM client_msgs WHERE direction=1') + row = cur.fetchone() + if row[0] != client_msgs_out: + raise ValueError("Found %d client_msgs_out, expected %d" % (row[0], client_msgs_out)) cur.execute('SELECT COUNT(*) FROM subscriptions') row = cur.fetchone() @@ -85,8 +90,8 @@ def check_counts(port, clients, client_msgs, base_msgs, retains, subscriptions): cur.execute('SELECT COUNT(*) FROM retains') row = cur.fetchone() - if row[0] != retains: - raise ValueError("Found %d retains, expected %d" % (row[0], retains)) + if row[0] != retain_msgs: + raise ValueError("Found %d retain_msgs, expected %d" % (row[0], retain_msgs)) con.close() @@ -94,6 +99,10 @@ def check_client(port, client_id, username, will_delay_time, session_expiry_time listener_port, max_packet_size, max_qos, retain_available, session_expiry_interval, will_delay_interval): + # "Fix" the infinite session expiry interval as mangled by an int32 conversion. + if session_expiry_interval == 4294967295: + session_expiry_interval = -1 + con = sqlite3.connect(f"{port}/mosquitto.sqlite3") cur = con.cursor() cur.execute('SELECT client_id, username, will_delay_time, session_expiry_time, ' + @@ -105,7 +114,7 @@ def check_client(port, client_id, username, will_delay_time, session_expiry_time if row[0] != client_id: raise ValueError("Invalid client_id %s / %s" % (row[0], client_id)) - if row[1] != username: + if username is not None and row[1] != username: raise ValueError("Invalid username %s / %s" % (row[1], username)) if (will_delay_time == 0 and row[2] != 0) or (will_delay_time != 0 and row[2] == 0): @@ -114,7 +123,7 @@ def check_client(port, client_id, username, will_delay_time, session_expiry_time if (session_expiry_time == 0 and row[3] != 0) or (session_expiry_time != 0 and row[3] == 0): raise ValueError("Invalid session_expiry_time %d / %d" % (row[3], session_expiry_time)) - if row[4] != listener_port: + if listener_port is not None and row[4] != listener_port: raise ValueError("Invalid listener_port %d / %d" % (row[4], listener_port)) if row[5] != max_packet_size: @@ -188,7 +197,7 @@ def check_client_msg(port, client_id, store_id, dup, direction, mid, qos, retain con.close() -def check_store_msg(port, expiry_time, topic, payload, source_id, source_username, +def check_base_msg(port, expiry_time, topic, payload, source_id, source_username, payloadlen, source_mid, source_port, qos, retain, idx=0): con = sqlite3.connect(f"{port}/mosquitto.sqlite3") diff --git a/test/broker/test.py b/test/broker/test.py index 89ba2b67..a028e30f 100755 --- a/test/broker/test.py +++ b/test/broker/test.py @@ -218,21 +218,22 @@ tests = [ (1, './14-dynsec-role-invalid.py'), (1, './14-dynsec-role.py'), - (1, './15-persist-client-msg-in-v3-1-1.py', 'persist_sqlite'), - (1, './15-persist-client-msg-in-v5-0.py', 'persist_sqlite'), - (1, './15-persist-client-msg-out-dup-v3-1-1.py', 'persist_sqlite'), - (1, './15-persist-client-msg-out-queue-v3-1-1.py', 'persist_sqlite'), - (1, './15-persist-client-msg-out-v3-1-1-db.py', 'persist_sqlite'), - (1, './15-persist-client-msg-out-v3-1-1.py', 'persist_sqlite'), - (1, './15-persist-client-msg-out-v5-0.py', 'persist_sqlite'), - (1, './15-persist-client-v3-1-1.py', 'persist_sqlite'), - (1, './15-persist-client-v5-0.py', 'persist_sqlite'), - (1, './15-persist-publish-properties-v5-0.py', 'persist_sqlite'), - (1, './15-persist-retain-clear.py', 'persist_sqlite'), - (1, './15-persist-retain-v3-1-1.py', 'persist_sqlite'), - (1, './15-persist-retain-v5-0.py', 'persist_sqlite'), - (1, './15-persist-subscription-v3-1-1.py', 'persist_sqlite'), - (1, './15-persist-subscription-v5-0.py', 'persist_sqlite'), + (1, './15-persist-client-msg-in-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-client-msg-in-v5-0.py', 'persist_sqlite'), + (1, './15-persist-client-msg-out-clear-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-client-msg-out-dup-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-client-msg-out-queue-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-client-msg-out-v3-1-1-db.py', 'persist_sqlite'), + (1, './15-persist-client-msg-out-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-client-msg-out-v5-0.py', 'persist_sqlite'), + (1, './15-persist-client-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-client-v5-0.py', 'persist_sqlite'), + (1, './15-persist-publish-properties-v5-0.py', 'persist_sqlite'), + (1, './15-persist-retain-clear.py', 'persist_sqlite'), + (1, './15-persist-retain-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-retain-v5-0.py', 'persist_sqlite'), + (1, './15-persist-subscription-v3-1-1.py', 'persist_sqlite'), + (1, './15-persist-subscription-v5-0.py', 'persist_sqlite'), (1, './16-cmd-args.py'), (1, './16-config-includedir.py'), diff --git a/test/mosq_test.py b/test/mosq_test.py index 053873fb..2c797bde 100644 --- a/test/mosq_test.py +++ b/test/mosq_test.py @@ -115,6 +115,17 @@ def wait_for_subprocess(client,timeout=10,terminate_timeout=2): pass return rc + +def terminate_broker(broker): + broker.terminate() + (_, stde) = broker.communicate() + if wait_for_subprocess(broker): + print("broker not terminated") + return (1, stde) + else: + return (0, stde) + + def pub_helper(port, proto_ver=4): connect_packet = gen_connect("pub-helper", proto_ver=proto_ver) connack_packet = gen_connack(rc=0, proto_ver=proto_ver)