completed regexp support in subscription system

This commit is contained in:
roberto@debian32
2011-11-01 08:33:52 +01:00
parent 5594503a56
commit 713370f343
4 changed files with 76 additions and 33 deletions
+4 -2
View File
@@ -57,6 +57,7 @@ struct uwsgi_fastrouter {
char *subscription_server;
struct uwsgi_subscribe_slot *subscriptions;
int subscription_regexp;
int socket_timeout;
@@ -109,6 +110,7 @@ struct option fastrouter_options[] = {
{"fastrouter-cheap", no_argument, &ufr.cheap, 1},
{"fastrouter-subscription-server", required_argument, 0, LONG_ARGS_FASTROUTER_SUBSCRIPTION_SERVER},
{"fastrouter-subscription-slot", required_argument, 0, LONG_ARGS_FASTROUTER_SUBSCRIPTION_SLOT},
{"fastrouter-subscription-use-regexp", no_argument, &ufr.subscription_regexp, 1},
{"fastrouter-timeout", required_argument, 0, LONG_ARGS_FASTROUTER_TIMEOUT},
{0, 0, 0, 0},
};
@@ -435,7 +437,7 @@ void fastrouter_loop() {
if (len > 0) {
memset(&usr, 0, sizeof(struct uwsgi_subscribe_req));
uwsgi_hooked_parse(bbuf+4, len-4, fastrouter_manage_subscription, &usr);
if (uwsgi_add_subscribe_node(&ufr.subscriptions, &usr, 0) && ufr.i_am_cheap) {
if (uwsgi_add_subscribe_node(&ufr.subscriptions, &usr, ufr.subscription_regexp) && ufr.i_am_cheap) {
struct uwsgi_fastrouter_socket *ufr_sock = ufr.sockets;
while(ufr_sock) {
event_queue_add_fd_read(ufr.queue, ufr_sock->fd);
@@ -516,7 +518,7 @@ void fastrouter_loop() {
fr_session->instance_address = tmp_socket_name;
}
else if (ufr.subscription_server) {
fr_session->un = uwsgi_get_subscribe_node(&ufr.subscriptions, fr_session->hostname, fr_session->hostname_len, 0);
fr_session->un = uwsgi_get_subscribe_node(&ufr.subscriptions, fr_session->hostname, fr_session->hostname_len, ufr.subscription_regexp);
if (fr_session->un && fr_session->un->len) {
fr_session->instance_address = fr_session->un->name;
fr_session->instance_address_len = fr_session->un->len;
+4 -2
View File
@@ -39,6 +39,7 @@ struct uwsgi_http {
int server;
char *subscription_server;
int subscription_regexp;
char *pattern;
int pattern_len;
@@ -76,6 +77,7 @@ struct option http_options[] = {
{"http-use-cluster", no_argument, &uhttp.use_cluster, 1},
{"http-events", required_argument, 0, LONG_ARGS_HTTP_EVENTS},
{"http-subscription-server", required_argument, 0, LONG_ARGS_HTTP_SUBSCRIPTION_SERVER},
{"http-subscription-use-regexp", no_argument, &uhttp.subscription_regexp, 1},
{"http-timeout", required_argument, 0, LONG_ARGS_HTTP_TIMEOUT},
{0, 0, 0, 0},
};
@@ -557,7 +559,7 @@ void http_loop() {
if (len > 0) {
memset(&usr, 0, sizeof(struct uwsgi_subscribe_req));
uwsgi_hooked_parse(bbuf+4, len-4, http_manage_subscription, &usr);
uwsgi_add_subscribe_node(&uhttp.subscriptions, &usr, 0);
uwsgi_add_subscribe_node(&uhttp.subscriptions, &usr, uhttp.subscription_regexp);
}
}
else {
@@ -637,7 +639,7 @@ void http_loop() {
uhttp_session->instance_address_len = uhttp.to_len;
}
else if (uhttp.subscription_server) {
uhttp_session->un = uwsgi_get_subscribe_node(&uhttp.subscriptions, uhttp_session->hostname, uhttp_session->hostname_len, 0);
uhttp_session->un = uwsgi_get_subscribe_node(&uhttp.subscriptions, uhttp_session->hostname, uhttp_session->hostname_len, uhttp.subscription_regexp);
if (uhttp_session->un && uhttp_session->un->len) {
uhttp_session->instance_address = uhttp_session->un->name;
uhttp_session->instance_address_len = uhttp_session->un->len;
+67 -29
View File
@@ -19,32 +19,21 @@
struct uwsgi_subscribe_slot *uwsgi_get_subscribe_slot(struct uwsgi_subscribe_slot **slot, char *key, uint16_t keylen, int regexp) {
struct uwsgi_subscribe_slot *current_slot = *slot;
#ifdef UWSGI_PCRE
int match;
#endif
if (keylen > 0xff) return NULL;
while(current_slot) {
#ifdef UWSGI_PCRE
match = 0;
if (regexp) {
if (uwsgi_regexp_match(current_slot->pattern, current_slot->pattern_extra, key, keylen)) {
match = 1;
if (uwsgi_regexp_match(current_slot->pattern, current_slot->pattern_extra, key, keylen) >= 0) {
return current_slot;
}
}
else {
#endif
if (!uwsgi_strncmp(key, keylen, current_slot->key, current_slot->keylen)) {
#ifdef UWSGI_PCRE
match = 1;
}
}
if (match) {
#endif
// auto optimization
if (current_slot->prev) {
// auto optimization
if (current_slot->prev) {
if (current_slot->hits > current_slot->prev->hits) {
struct uwsgi_subscribe_slot *slot_parent = current_slot->prev->prev, *slot_prev = current_slot->prev;
if (slot_parent) {
@@ -60,10 +49,13 @@ struct uwsgi_subscribe_slot *uwsgi_get_subscribe_slot(struct uwsgi_subscribe_slo
current_slot->next = slot_prev;
current_slot->prev = slot_parent;
}
}
}
return current_slot;
return current_slot;
}
#ifdef UWSGI_PCRE
}
#endif
current_slot = current_slot->next;
}
@@ -150,7 +142,7 @@ void uwsgi_remove_subscribe_node(struct uwsgi_subscribe_slot **slot, struct uwsg
struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slot **slot, struct uwsgi_subscribe_req *usr, int regexp) {
struct uwsgi_subscribe_slot *current_slot = uwsgi_get_subscribe_slot(slot, usr->key, usr->keylen, regexp), *old_slot = NULL, *a_slot;
struct uwsgi_subscribe_slot *current_slot = uwsgi_get_subscribe_slot(slot, usr->key, usr->keylen, 0), *old_slot = NULL, *a_slot;
struct uwsgi_subscribe_node *node, *old_node = NULL;
if (usr->address_len > 0xff) return NULL;
@@ -182,12 +174,6 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
}
else {
a_slot = *slot;
while(a_slot) {
old_slot = a_slot;
a_slot = a_slot->next;
}
current_slot = uwsgi_malloc(sizeof(struct uwsgi_subscribe_slot));
current_slot->keylen = usr->keylen;
memcpy(current_slot->key, usr->key, usr->keylen);
@@ -216,16 +202,68 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
current_slot->nodes->next = NULL;
if (old_slot) {
old_slot->next = current_slot;
#ifdef UWSGI_PCRE
// if key is a regexp, order it by keylen
if (regexp) {
old_slot = NULL;
a_slot = *slot;
while(a_slot) {
if (a_slot->keylen > current_slot->keylen) {
old_slot = a_slot;
break;
}
a_slot = a_slot->next;
}
if (old_slot) {
current_slot->prev = old_slot->prev;
old_slot->prev = current_slot;
if (current_slot->prev) {
old_slot->prev->next = current_slot;
}
current_slot->next = old_slot;
}
else {
a_slot = *slot;
while(a_slot) {
old_slot = a_slot;
a_slot = a_slot->next;
}
if (old_slot) {
old_slot->next = current_slot;
}
current_slot->prev = old_slot;
current_slot->next = NULL;
}
}
else {
#endif
a_slot = *slot;
while(a_slot) {
old_slot = a_slot;
a_slot = a_slot->next;
}
current_slot->prev = old_slot;
current_slot->next = NULL;
if (!*slot) {
if (old_slot) {
old_slot->next = current_slot;
}
current_slot->prev = old_slot;
current_slot->next = NULL;
#ifdef UWSGI_PCRE
}
#endif
if (!*slot || current_slot->prev == NULL) {
*slot = current_slot;
}
uwsgi_log("[uwsgi-subscription] new pool: %.*s\n", usr->keylen, usr->key);
uwsgi_log("[uwsgi-subscription] %.*s => new node: %.*s\n", usr->keylen, usr->key, usr->address_len, usr->address);
return current_slot->nodes;
+1
View File
@@ -3385,6 +3385,7 @@ static int manage_base_opt(int i, char *optarg) {
}
return 1;
case LONG_ARGS_SUBSCRIBE_TO:
uwsgi.master_process = 1;
uwsgi_string_new_list(&uwsgi.subscriptions, optarg);
return 1;
#ifdef __linux__