From 90b6512d524c7ee5077952843e01ea127cfd51ec Mon Sep 17 00:00:00 2001 From: Unbit Date: Fri, 11 Jul 2014 11:46:23 +0200 Subject: [PATCH] pluggable subscription algos --- core/subscription.c | 302 ++++++++++++++++++-------------- plugins/corerouter/corerouter.c | 3 +- plugins/corerouter/cr_map.c | 4 +- uwsgi.h | 28 ++- 4 files changed, 193 insertions(+), 144 deletions(-) diff --git a/core/subscription.c b/core/subscription.c index 0dfecf0f..fc3bcb32 100644 --- a/core/subscription.c +++ b/core/subscription.c @@ -128,133 +128,7 @@ struct uwsgi_subscribe_slot *uwsgi_get_subscribe_slot(struct uwsgi_subscribe_slo return NULL; } -// least reference count -static struct uwsgi_subscribe_node *uwsgi_subscription_algo_lrc(struct uwsgi_subscribe_slot *current_slot, struct uwsgi_subscribe_node *node) { - // if node is NULL we are in the second step (in lrc mode we do not use the first step) - if (node) - return NULL; - - struct uwsgi_subscribe_node *choosen_node = NULL; - node = current_slot->nodes; - uint64_t min_rc = 0; - while (node) { - if (!node->death_mark) { - if (min_rc == 0 || node->reference < min_rc) { - min_rc = node->reference; - choosen_node = node; - if (min_rc == 0 && !(node->next && node->next->reference <= node->reference && node->next->last_requests <= node->last_requests)) - break; - } - } - node = node->next; - } - - if (choosen_node) { - choosen_node->reference++; - } - - return choosen_node; -} - -// weighted least reference count -static struct uwsgi_subscribe_node *uwsgi_subscription_algo_wlrc(struct uwsgi_subscribe_slot *current_slot, struct uwsgi_subscribe_node *node) { - // if node is NULL we are in the second step (in wlrc mode we do not use the first step) - if (node) - return NULL; - - struct uwsgi_subscribe_node *choosen_node = NULL; - node = current_slot->nodes; - double min_rc = 0; - while (node) { - if (!node->death_mark) { - // node->weight is always >= 1, we can safely use it as divider - double ref = (double) node->reference / (double) node->weight; - double next_node_ref = 0; - if (node->next) - next_node_ref = (double) node->next->reference / (double) node->next->weight; - - if (min_rc == 0 || ref < min_rc) { - min_rc = ref; - choosen_node = node; - if (min_rc == 0 && !(node->next && next_node_ref <= ref && node->next->last_requests <= node->last_requests)) - break; - } - } - node = node->next; - } - - if (choosen_node) { - choosen_node->reference++; - } - - return choosen_node; -} - -// weighted round robin algo -static struct uwsgi_subscribe_node *uwsgi_subscription_algo_wrr(struct uwsgi_subscribe_slot *current_slot, struct uwsgi_subscribe_node *node) { - // if node is NULL we are in the second step - if (node) { - if (node->death_mark == 0 && node->wrr > 0) { - node->wrr--; - node->reference++; - return node; - } - return NULL; - } - - // no wrr > 0 node found, reset them - node = current_slot->nodes; - uint64_t min_weight = 0; - while (node) { - if (!node->death_mark) { - if (min_weight == 0 || node->weight < min_weight) - min_weight = node->weight; - } - node = node->next; - } - - // now set wrr - node = current_slot->nodes; - struct uwsgi_subscribe_node *choosen_node = NULL; - while (node) { - if (!node->death_mark) { - node->wrr = node->weight / min_weight; - choosen_node = node; - } - node = node->next; - } - if (choosen_node) { - choosen_node->wrr--; - choosen_node->reference++; - } - return choosen_node; -} - -void uwsgi_subscription_set_algo(char *algo) { - - if (!algo) - goto wrr; - - if (!strcmp(algo, "wrr")) { - uwsgi.subscription_algo = uwsgi_subscription_algo_wrr; - return; - } - - if (!strcmp(algo, "lrc")) { - uwsgi.subscription_algo = uwsgi_subscription_algo_lrc; - return; - } - - if (!strcmp(algo, "wlrc")) { - uwsgi.subscription_algo = uwsgi_subscription_algo_wlrc; - return; - } - -wrr: - uwsgi.subscription_algo = uwsgi_subscription_algo_wrr; -} - -struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slot **slot, char *key, uint16_t keylen) { +struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slot **slot, char *key, uint16_t keylen, struct uwsgi_subscription_client *client) { if (keylen > 0xff) return NULL; @@ -287,14 +161,14 @@ struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slo continue; } - struct uwsgi_subscribe_node *choosen_node = uwsgi.subscription_algo(current_slot, node); + struct uwsgi_subscribe_node *choosen_node = current_slot->algo(current_slot, node, client); if (choosen_node) return choosen_node; node = node->next; } - return uwsgi.subscription_algo(current_slot, node); + return current_slot->algo(current_slot, node, client); } struct uwsgi_subscribe_node *uwsgi_get_subscribe_node_by_name(struct uwsgi_subscribe_slot **slot, char *key, uint16_t keylen, char *val, uint16_t vallen) { @@ -436,6 +310,10 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo node->cores = usr->cores; node->load = usr->load; node->weight = usr->weight; + node->backup_level = usr->backup_level; + if (usr->proto_len > 0) { + node->proto = usr->proto[0]; + } if (!node->weight) node->weight = 1; node->last_requests = 0; @@ -468,6 +346,10 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo node->cores = usr->cores; node->load = usr->load; node->weight = usr->weight; + node->backup_level = usr->backup_level; + if (usr->proto_len > 0) { + node->proto = usr->proto[0]; + } node->unix_check = usr->unix_check; if (!node->weight) node->weight = 1; @@ -539,6 +421,10 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo current_slot->nodes->cores = usr->cores; current_slot->nodes->load = usr->load; current_slot->nodes->weight = usr->weight; + current_slot->nodes->backup_level = usr->backup_level; + if (usr->proto_len > 0) { + current_slot->nodes->proto = usr->proto[0]; + } current_slot->nodes->unix_check = usr->unix_check; if (!current_slot->nodes->weight) current_slot->nodes->weight = 1; @@ -570,6 +456,9 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo current_slot->prev = old_slot; current_slot->next = NULL; + current_slot->algo = usr->algo; + if (!current_slot->algo) current_slot->algo = uwsgi.subscription_algo; + if (!slot[hash_key] || current_slot->prev == NULL) { slot[hash_key] = current_slot; @@ -861,13 +750,6 @@ int uwsgi_no_subscriptions(struct uwsgi_subscribe_slot **slot) { return 1; } -struct uwsgi_subscribe_slot **uwsgi_subscription_init_ht() { - if (!uwsgi.subscription_algo) { - uwsgi_subscription_set_algo(NULL); - } - return uwsgi_calloc(sizeof(struct uwsgi_subscription_slot *) * UMAX16); -} - void uwsgi_subscribe(char *subscription, uint8_t cmd) { size_t subfile_size; @@ -1105,3 +987,151 @@ void uwsgi_subscribe_all(uint8_t cmd, int verbose) { } } + + +// least reference count +static struct uwsgi_subscribe_node *uwsgi_subscription_algo_lrc(struct uwsgi_subscribe_slot *current_slot, struct uwsgi_subscribe_node *node, struct uwsgi_subscription_client *client) { + // if node is NULL we are in the second step (in lrc mode we do not use the first step) + if (node) + return NULL; + + struct uwsgi_subscribe_node *choosen_node = NULL; + node = current_slot->nodes; + uint64_t min_rc = 0; + while (node) { + if (!node->death_mark) { + if (min_rc == 0 || node->reference < min_rc) { + min_rc = node->reference; + choosen_node = node; + if (min_rc == 0 && !(node->next && node->next->reference <= node->reference && node->next->last_requests <= node->last_requests)) + break; + } + } + node = node->next; + } + + if (choosen_node) { + choosen_node->reference++; + } + + return choosen_node; +} + +// weighted least reference count +static struct uwsgi_subscribe_node *uwsgi_subscription_algo_wlrc(struct uwsgi_subscribe_slot *current_slot, struct uwsgi_subscribe_node *node, struct uwsgi_subscription_client *client) { + // if node is NULL we are in the second step (in wlrc mode we do not use the first step) + if (node) + return NULL; + + struct uwsgi_subscribe_node *choosen_node = NULL; + node = current_slot->nodes; + double min_rc = 0; + while (node) { + if (!node->death_mark) { + // node->weight is always >= 1, we can safely use it as divider + double ref = (double) node->reference / (double) node->weight; + double next_node_ref = 0; + if (node->next) + next_node_ref = (double) node->next->reference / (double) node->next->weight; + + if (min_rc == 0 || ref < min_rc) { + min_rc = ref; + choosen_node = node; + if (min_rc == 0 && !(node->next && next_node_ref <= ref && node->next->last_requests <= node->last_requests)) + break; + } + } + node = node->next; + } + + if (choosen_node) { + choosen_node->reference++; + } + + return choosen_node; +} + +// weighted round robin algo +static struct uwsgi_subscribe_node *uwsgi_subscription_algo_wrr(struct uwsgi_subscribe_slot *current_slot, struct uwsgi_subscribe_node *node, struct uwsgi_subscription_client *client) { + // if node is NULL we are in the second step + if (node) { + if (node->death_mark == 0 && node->wrr > 0) { + node->wrr--; + node->reference++; + return node; + } + return NULL; + } + + // no wrr > 0 node found, reset them + node = current_slot->nodes; + uint64_t min_weight = 0; + while (node) { + if (!node->death_mark) { + if (min_weight == 0 || node->weight < min_weight) + min_weight = node->weight; + } + node = node->next; + } + + // now set wrr + node = current_slot->nodes; + struct uwsgi_subscribe_node *choosen_node = NULL; + while (node) { + if (!node->death_mark) { + node->wrr = node->weight / min_weight; + choosen_node = node; + } + node = node->next; + } + if (choosen_node) { + choosen_node->wrr--; + choosen_node->reference++; + } + return choosen_node; +} + +void uwsgi_subscription_init_algos() { + + uwsgi_register_subscription_algo("wrr", uwsgi_subscription_algo_wrr); + uwsgi_register_subscription_algo("lrc", uwsgi_subscription_algo_lrc); + uwsgi_register_subscription_algo("wlrc", uwsgi_subscription_algo_wlrc); +} + +void uwsgi_subscription_set_algo(char *algo) { + if (!uwsgi.subscription_algos) { + uwsgi_register_subscription_algo("wrr", uwsgi_subscription_algo_wrr); + uwsgi_register_subscription_algo("lrc", uwsgi_subscription_algo_lrc); + uwsgi_register_subscription_algo("wlrc", uwsgi_subscription_algo_wlrc); + } + if (!algo) + goto wrr; + uwsgi.subscription_algo = uwsgi_subscription_algo_get(algo, strlen(algo)); + if (uwsgi.subscription_algo) return ; + +wrr: + uwsgi.subscription_algo = uwsgi_subscription_algo_wrr; +} + +// we are lazy for subscription algos, we initialize them only if needed +struct uwsgi_subscribe_slot **uwsgi_subscription_init_ht() { + if (!uwsgi.subscription_algo) { + uwsgi_subscription_set_algo(NULL); + } + return uwsgi_calloc(sizeof(struct uwsgi_subscription_slot *) * UMAX16); +} + +struct uwsgi_subscribe_node *(*uwsgi_subscription_algo_get(char *name , size_t len))(struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *) { + struct uwsgi_string_list *usl = NULL; + uwsgi_foreach(usl, uwsgi.subscription_algos) { + if (!uwsgi_strncmp(usl->value, usl->len, name, len)) { + return (struct uwsgi_subscribe_node *(*)(struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *)) usl->custom_ptr; + } + } + return NULL; +} + +void uwsgi_register_subscription_algo(char *name, struct uwsgi_subscribe_node *(*func)(struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *)) { + struct uwsgi_string_list *usl = uwsgi_string_new_list(&uwsgi.subscription_algos, name); + usl->custom_ptr = func; +} diff --git a/plugins/corerouter/corerouter.c b/plugins/corerouter/corerouter.c index 604a6984..46253e09 100644 --- a/plugins/corerouter/corerouter.c +++ b/plugins/corerouter/corerouter.c @@ -292,8 +292,7 @@ void corerouter_manage_subscription(char *key, uint16_t keylen, char *val, uint1 usr->notify_len = vallen; } else if (!uwsgi_strncmp("algo", 4, key, keylen)) { - usr->algo = val; - usr->algo_len = vallen; + usr->algo = uwsgi_subscription_algo_get(val, vallen); } else if (!uwsgi_strncmp("backup", 6, key, keylen)) { usr->backup_level = uwsgi_str_num(val, vallen); diff --git a/plugins/corerouter/cr_map.c b/plugins/corerouter/cr_map.c index 066031e8..06e1e177 100644 --- a/plugins/corerouter/cr_map.c +++ b/plugins/corerouter/cr_map.c @@ -48,7 +48,7 @@ int uwsgi_cr_map_use_pattern(struct uwsgi_corerouter *ucr, struct corerouter_pee int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) { - peer->un = uwsgi_get_subscribe_node(ucr->subscriptions, peer->key, peer->key_len); + peer->un = uwsgi_get_subscribe_node(ucr->subscriptions, peer->key, peer->key_len, NULL); if (peer->un && peer->un->len) { peer->instance_address = peer->un->name; peer->instance_address_len = peer->un->len; @@ -73,7 +73,7 @@ split: #ifdef UWSGI_DEBUG uwsgi_log("trying with %.*s\n", name_len, name); #endif - peer->un = uwsgi_get_subscribe_node(ucr->subscriptions, name, name_len); + peer->un = uwsgi_get_subscribe_node(ucr->subscriptions, name, name_len, NULL); if (!peer->un) { char *next = memchr(name+1, '.', name_len-1); if (next) { diff --git a/uwsgi.h b/uwsgi.h index 51a9b9bb..0675990b 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -1689,6 +1689,8 @@ struct uwsgi_fsmon { struct uwsgi_fsmon *next; }; +struct uwsgi_subscription_client; + struct uwsgi_server { // store the machine hostname @@ -2638,7 +2640,7 @@ struct uwsgi_server { struct uwsgi_string_list *subscriptions; struct uwsgi_string_list *subscriptions2; - struct uwsgi_subscribe_node *(*subscription_algo) (struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *); + struct uwsgi_subscribe_node *(*subscription_algo) (struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *); int subscription_dotsplit; int never_swap; @@ -2737,6 +2739,7 @@ struct uwsgi_server { struct uwsgi_string_list *hook_as_on_demand_vassal; uint64_t max_requests_delta; char *emperor_chdir_attr; + struct uwsgi_string_list *subscription_algos; }; struct uwsgi_rpc { @@ -3336,8 +3339,7 @@ struct uwsgi_subscribe_req { char *proto; uint16_t proto_len; - char *algo; - uint16_t algo_len; + struct uwsgi_subscribe_node *(*algo) (struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *); }; void uwsgi_nuclear_blast(); @@ -3567,6 +3569,11 @@ struct uwsgi_mule_farm *uwsgi_mule_farm_new(struct uwsgi_mule_farm **, struct uw int uwsgi_farm_has_mule(struct uwsgi_farm *, int); struct uwsgi_farm *get_farm_by_name(char *); +struct uwsgi_subscription_client { + int fd; + union uwsgi_sockaddr sockaddr; + char *cookie; +}; struct uwsgi_subscribe_node { @@ -3606,6 +3613,11 @@ struct uwsgi_subscribe_node { struct uwsgi_subscribe_slot *slot; struct uwsgi_subscribe_node *next; + + // uWSGI 2.1 + uint64_t backup_level; + //here the solution is a bit hacky, we take the first letter of the proto ('u','\0' -> uwsgi, 'h' -> http, 'f' -> fastcgi, 's' -> scgi) + char proto; }; struct uwsgi_subscribe_slot { @@ -3628,6 +3640,9 @@ struct uwsgi_subscribe_slot { uint8_t sni_enabled; #endif + // uWSGI 2.1 (algo is required) + struct uwsgi_subscribe_node *(*algo) (struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *); + }; void mule_send_msg(int, char *, size_t); @@ -3637,7 +3652,7 @@ void create_signal_pipe(int *); void create_msg_pipe(int *, int); struct uwsgi_subscribe_slot *uwsgi_get_subscribe_slot(struct uwsgi_subscribe_slot **, char *, uint16_t); struct uwsgi_subscribe_node *uwsgi_get_subscribe_node_by_name(struct uwsgi_subscribe_slot **, char *, uint16_t, char *, uint16_t); -struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slot **, char *, uint16_t); +struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slot **, char *, uint16_t, struct uwsgi_subscription_client *); int uwsgi_remove_subscribe_node(struct uwsgi_subscribe_slot **, struct uwsgi_subscribe_node *); struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slot **, struct uwsgi_subscribe_req *); @@ -4840,6 +4855,11 @@ int uwsgi_webdav_multistatus_propstat_new(struct uwsgi_buffer *); int uwsgi_webdav_multistatus_propstat_close(struct uwsgi_buffer *); int uwsgi_webdav_multistatus_prop_new(struct uwsgi_buffer *); int uwsgi_webdav_multistatus_prop_close(struct uwsgi_buffer *); + +struct uwsgi_subscribe_node *(*uwsgi_subscription_algo_get(char * , size_t))(struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *); + +void uwsgi_subscription_init_algos(void); +void uwsgi_register_subscription_algo(char *, struct uwsgi_subscribe_node *(*) (struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *)); #ifdef __cplusplus } #endif