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.
This commit is contained in:
Roger A. Light
2022-10-09 22:17:47 +01:00
parent 428d39b22b
commit 16feb14a57
30 changed files with 310 additions and 206 deletions
+1 -2
View File
@@ -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 */
+1
View File
@@ -18,6 +18,7 @@ Contributors:
#include <string.h>
#include <sqlite3.h>
#include <stdio.h>
#include <stdlib.h>
#include <cjson/cJSON.h>
+6 -10
View File
@@ -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++;
+1
View File
@@ -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;
}
+1 -1
View File
@@ -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);
-3
View File
@@ -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);
}
-11
View File
@@ -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;
-1
View File
@@ -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);
-4
View File
@@ -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:
-22
View File
@@ -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;
+13 -7
View File
@@ -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)
+10 -6
View File
@@ -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)
+111
View File
@@ -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)
@@ -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)
@@ -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)
@@ -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)
@@ -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)
@@ -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)
+9 -10
View File
@@ -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)
+11 -13
View File
@@ -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)
@@ -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)
+9 -11
View File
@@ -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)
+7 -6
View File
@@ -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)
+7 -6
View File
@@ -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)
+21 -10
View File
@@ -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)
+9 -10
View File
@@ -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)
+1
View File
@@ -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
+18 -9
View File
@@ -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")
+16 -15
View File
@@ -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'),
+11
View File
@@ -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)