Improve performance with lots of queued message

Split message queue in two queues: in-flight and queued to avoid the
need to iterate over all messages.

Signed-off-by: Pierre Fersing <pierre.fersing@bleemeo.com>
This commit is contained in:
Pierre Fersing
2016-04-18 16:24:13 +02:00
parent 62402f7b60
commit 44f23252a0
6 changed files with 308 additions and 225 deletions
+4 -2
View File
@@ -195,8 +195,10 @@ struct mosquitto {
bool is_dropping;
bool is_bridge;
struct mosquitto__bridge *bridge;
struct mosquitto_client_msg *msgs;
struct mosquitto_client_msg *last_msg;
struct mosquitto_client_msg *inflight_msgs;
struct mosquitto_client_msg *last_inflight_msg;
struct mosquitto_client_msg *queued_msgs;
struct mosquitto_client_msg *last_queued_msg;
int msg_count;
int msg_count12;
struct mosquitto__acl_user *acl_list;
+16 -5
View File
@@ -70,8 +70,10 @@ struct mosquitto *context__init(struct mosquitto_db *db, mosq_sock_t sock)
}
}
context->bridge = NULL;
context->msgs = NULL;
context->last_msg = NULL;
context->inflight_msgs = NULL;
context->last_inflight_msg = NULL;
context->queued_msgs = NULL;
context->last_queued_msg = NULL;
context->msg_count = 0;
context->msg_count12 = 0;
#ifdef WITH_TLS
@@ -166,15 +168,24 @@ void context__cleanup(struct mosquitto_db *db, struct mosquitto *context, bool d
context->will = NULL;
}
if(do_free || context->clean_session){
msg = context->msgs;
msg = context->inflight_msgs;
while(msg){
next = msg->next;
db__msg_store_deref(db, &msg->store);
mosquitto__free(msg);
msg = next;
}
context->msgs = NULL;
context->last_msg = NULL;
context->inflight_msgs = NULL;
context->last_inflight_msg = NULL;
msg = context->queued_msgs;
while(msg){
next = msg->next;
db__msg_store_deref(db, &msg->store);
mosquitto__free(msg);
msg = next;
}
context->queued_msgs = NULL;
context->last_queued_msg = NULL;
}
if(do_free){
mosquitto__free(context);
+250 -207
View File
File diff suppressed because it is too large Load Diff
+18 -6
View File
@@ -85,9 +85,11 @@ int handle__connect(struct mosquitto_db *db, struct mosquitto *context)
uint8_t will, will_retain, will_qos, clean_session;
uint8_t username_flag, password_flag;
char *username = NULL, *password = NULL;
bool process_queued_msgs = false;
int rc;
struct mosquitto__acl_user *acl_tail;
struct mosquitto_client_msg *msg_tail, *msg_prev;
struct mosquitto_client_msg **msgs;
struct mosquitto *found_context;
int slen;
struct mosquitto__subleaf *leaf;
@@ -442,9 +444,11 @@ int handle__connect(struct mosquitto_db *db, struct mosquitto *context)
context->clean_session = clean_session;
if(context->clean_session == false && found_context->clean_session == false){
if(found_context->msgs){
context->msgs = found_context->msgs;
found_context->msgs = NULL;
if(found_context->inflight_msgs || found_context->queued_msgs){
context->inflight_msgs = found_context->inflight_msgs;
context->queued_msgs = found_context->queued_msgs;
found_context->inflight_msgs = NULL;
found_context->queued_msgs = NULL;
db__message_reconnect_reset(db, context);
}
context->subs = found_context->subs;
@@ -532,7 +536,8 @@ int handle__connect(struct mosquitto_db *db, struct mosquitto *context)
/* Remove any queued messages that are no longer allowed through ACL,
* assuming a possible change of username. */
msg_tail = context->msgs;
msgs = &(context->inflight_msgs);
msg_tail = *msgs;
msg_prev = NULL;
while(msg_tail){
if(msg_tail->direction == mosq_md_out){
@@ -543,10 +548,11 @@ int handle__connect(struct mosquitto_db *db, struct mosquitto *context)
mosquitto__free(msg_tail);
msg_tail = msg_prev->next;
}else{
context->msgs = context->msgs->next;
*msgs = (*msgs)->next;
mosquitto__free(msg_tail);
msg_tail = context->msgs;
msg_tail = (*msgs);
}
// XXX: why it does not update last_msg if msg_tail was the last message ?
}else{
msg_prev = msg_tail;
msg_tail = msg_tail->next;
@@ -555,6 +561,12 @@ int handle__connect(struct mosquitto_db *db, struct mosquitto *context)
msg_prev = msg_tail;
msg_tail = msg_tail->next;
}
if (msg_tail == NULL && !process_queued_msgs){
msgs = &(context->queued_msgs);
msg_tail = *msgs;
msg_prev = NULL;
process_queued_msgs = true;
}
}
HASH_ADD_KEYPTR(hh_id, db->contexts_by_id, context->id, strlen(context->id), context);
+1
View File
@@ -499,6 +499,7 @@ int db__message_insert(struct mosquitto_db *db, struct mosquitto *context, uint1
int db__message_release(struct mosquitto_db *db, struct mosquitto *context, uint16_t mid, enum mosquitto_msg_direction dir);
int db__message_update(struct mosquitto *context, uint16_t mid, enum mosquitto_msg_direction dir, enum mosquitto_msg_state state);
int db__message_write(struct mosquitto_db *db, struct mosquitto *context);
void db__message_dequeue_first(struct mosquitto *context);
int db__messages_delete(struct mosquitto_db *db, struct mosquitto *context);
int db__messages_easy_queue(struct mosquitto_db *db, struct mosquitto *context, const char *topic, int qos, uint32_t payloadlen, const void *payload, int retain);
int db__message_store(struct mosquitto_db *db, const char *source, uint16_t source_mid, char *topic, int qos, uint32_t payloadlen, mosquitto__payload_uhpa *payload, int retain, struct mosquitto_msg_store **stored, dbid_t store_id);
+19 -5
View File
@@ -71,13 +71,14 @@ static int persist__client_messages_write(struct mosquitto_db *db, FILE *db_fptr
dbid_t i64temp;
uint16_t i16temp, slen;
uint8_t i8temp;
bool process_queued_message = false;
struct mosquitto_client_msg *cmsg;
assert(db);
assert(db_fptr);
assert(context);
cmsg = context->msgs;
cmsg = context->inflight_msgs;
while(cmsg){
slen = strlen(context->id);
@@ -115,6 +116,10 @@ static int persist__client_messages_write(struct mosquitto_db *db, FILE *db_fptr
write_e(db_fptr, &i8temp, sizeof(uint8_t));
cmsg = cmsg->next;
if (cmsg == NULL && !process_queued_message){
cmsg = context->queued_msgs;
process_queued_message = true;
}
}
return MOSQ_ERR_SUCCESS;
@@ -408,6 +413,7 @@ error:
static int persist__client_msg_restore(struct mosquitto_db *db, const char *client_id, uint16_t mid, uint8_t qos, uint8_t retain, uint8_t direction, uint8_t state, uint8_t dup, uint64_t store_id)
{
struct mosquitto_client_msg *cmsg;
struct mosquitto_client_msg **msgs, **last_msg;
struct mosquitto_msg_store_load *load;
struct mosquitto *context;
@@ -442,12 +448,20 @@ static int persist__client_msg_restore(struct mosquitto_db *db, const char *clie
log__printf(NULL, MOSQ_LOG_ERR, "Error restoring persistent database, message store corrupt.");
return 1;
}
if(context->msgs){
context->last_msg->next = cmsg;
if (state == mosq_ms_queued){
msgs = &(context->queued_msgs);
last_msg = &(context->last_queued_msg);
}else{
context->msgs = cmsg;
msgs = &(context->inflight_msgs);
last_msg = &(context->last_inflight_msg);
}
context->last_msg = cmsg;
if(*msgs){
(*last_msg)->next = cmsg;
}else{
*msgs = cmsg;
}
*last_msg = cmsg;
return MOSQ_ERR_SUCCESS;
}