From 4b09ead07d0738172aee5dc96515506f526bb64d Mon Sep 17 00:00:00 2001 From: Unbit Date: Sun, 12 Jan 2014 09:04:57 +0100 Subject: [PATCH] refactored recured subscription system --- core/subscription.c | 723 +++++++++++++++++++++++--------------------- 1 file changed, 376 insertions(+), 347 deletions(-) diff --git a/core/subscription.c b/core/subscription.c index a9f40315..67648a01 100644 --- a/core/subscription.c +++ b/core/subscription.c @@ -23,46 +23,52 @@ extern struct uwsgi_server uwsgi; #ifdef UWSGI_SSL static void uwsgi_subscription_sni_check(struct uwsgi_subscribe_slot *current_slot, struct uwsgi_subscribe_req *usr) { if (usr->sni_key_len > 0 && usr->sni_crt_len > 0) { - if (!current_slot->sni_enabled) { - char *sni_key = uwsgi_concat2n(usr->sni_key, usr->sni_key_len, "", 0); - char *sni_crt = uwsgi_concat2n(usr->sni_crt, usr->sni_crt_len, "", 0); - char *sni_ca = NULL; - if (usr->sni_ca_len > 0) { - sni_ca = uwsgi_concat2n(usr->sni_ca, usr->sni_ca_len, "", 0); - } - char *servername = NULL; - char *colon = memchr(current_slot->key, ':', current_slot->keylen); - if (colon) { - servername = uwsgi_concat2n(current_slot->key, colon-current_slot->key, "", 0); - } - else { - servername = uwsgi_concat2n(current_slot->key, current_slot->keylen, "", 0); - } - if (uwsgi_ssl_add_sni_item(servername, sni_crt, sni_key, uwsgi.sni_dir_ciphers , sni_ca)) { - current_slot->sni_enabled = 1; - } - if (sni_key) free(sni_key); - if (sni_crt) free(sni_crt); - if (sni_ca) free(sni_ca); - } - } + if (!current_slot->sni_enabled) { + char *sni_key = uwsgi_concat2n(usr->sni_key, usr->sni_key_len, "", 0); + char *sni_crt = uwsgi_concat2n(usr->sni_crt, usr->sni_crt_len, "", 0); + char *sni_ca = NULL; + if (usr->sni_ca_len > 0) { + sni_ca = uwsgi_concat2n(usr->sni_ca, usr->sni_ca_len, "", 0); + } + char *servername = NULL; + char *colon = memchr(current_slot->key, ':', current_slot->keylen); + if (colon) { + servername = uwsgi_concat2n(current_slot->key, colon - current_slot->key, "", 0); + } + else { + servername = uwsgi_concat2n(current_slot->key, current_slot->keylen, "", 0); + } + if (uwsgi_ssl_add_sni_item(servername, sni_crt, sni_key, uwsgi.sni_dir_ciphers, sni_ca)) { + current_slot->sni_enabled = 1; + } + if (sni_key) + free(sni_key); + if (sni_crt) + free(sni_crt); + if (sni_ca) + free(sni_ca); + } + } } #endif int uwsgi_subscription_credentials_check(struct uwsgi_subscribe_slot *slot, struct uwsgi_subscribe_req *usr) { - struct uwsgi_string_list *usl = NULL; - uwsgi_foreach(usl, uwsgi.subscriptions_credentials_check_dir) { - char *filename = uwsgi_concat2n(usl->value, usl->len, slot->key, slot->keylen); - struct stat st; - int ret = stat(filename, &st); - free(filename); - if (ret != 0) continue; - if (st.st_uid != usr->uid) continue; - if (st.st_gid != usr->gid) continue; - // accepted... - return 1; - } - return 0; + struct uwsgi_string_list *usl = NULL; + uwsgi_foreach(usl, uwsgi.subscriptions_credentials_check_dir) { + char *filename = uwsgi_concat2n(usl->value, usl->len, slot->key, slot->keylen); + struct stat st; + int ret = stat(filename, &st); + free(filename); + if (ret != 0) + continue; + if (st.st_uid != usr->uid) + continue; + if (st.st_gid != usr->gid) + continue; + // accepted... + return 1; + } + return 0; } struct uwsgi_subscribe_slot *uwsgi_get_subscribe_slot(struct uwsgi_subscribe_slot **slot, char *key, uint16_t keylen) { @@ -346,7 +352,7 @@ int uwsgi_remove_subscribe_node(struct uwsgi_subscribe_slot **slot, struct uwsgi // first check if i am the only node if ((!prev_slot && !next_slot) || next_slot == node_slot) { #ifdef UWSGI_SSL - if (uwsgi.subscriptions_sign_check_dir) { + if (node_slot->sign_ctx) { EVP_PKEY_free(node_slot->sign_public_key); EVP_MD_CTX_destroy(node_slot->sign_ctx); } @@ -375,7 +381,7 @@ int uwsgi_remove_subscribe_node(struct uwsgi_subscribe_slot **slot, struct uwsgi } #ifdef UWSGI_SSL - if (uwsgi.subscriptions_sign_check_dir) { + if (node_slot->sign_ctx) { EVP_PKEY_free(node_slot->sign_public_key); EVP_MD_CTX_destroy(node_slot->sign_ctx); } @@ -388,6 +394,8 @@ end: return ret; } +static int subscription_new_sign_ctx(struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_req *); + struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slot **slot, struct uwsgi_subscribe_req *usr) { struct uwsgi_subscribe_slot *current_slot = uwsgi_get_subscribe_slot(slot, usr->key, usr->keylen), *old_slot = NULL, *a_slot; @@ -412,7 +420,7 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo if (!uwsgi_strncmp(node->name, node->len, usr->address, usr->address_len)) { #ifdef UWSGI_SSL // this should avoid sending sniffed packets... - if (uwsgi.subscriptions_sign_check_dir && usr->unix_check <= node->unix_check) { + if (current_slot->sign_ctx && usr->unix_check <= node->unix_check) { uwsgi_log("[uwsgi-subscription for pid %d] invalid (sniffed ?) packet sent for slot: %.*s node: %.*s unix_check: %lu\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, (unsigned long) usr->unix_check); return NULL; } @@ -433,12 +441,8 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo } #ifdef UWSGI_SSL - if (uwsgi.subscriptions_sign_check_dir && usr->unix_check < (uwsgi_now() - (time_t) uwsgi.subscriptions_sign_check_tolerance)) { - uwsgi_log("[uwsgi-subscription for pid %d] invalid (sniffed ?) packet sent for slot: %.*s node: %.*s unix_check: %lu\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, (unsigned long) usr->unix_check); - return NULL; - } // check here as we are sure the node will be added - uwsgi_subscription_sni_check(current_slot, usr); + uwsgi_subscription_sni_check(current_slot, usr); #endif node = uwsgi_malloc(sizeof(struct uwsgi_subscribe_node)); @@ -478,56 +482,22 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo if (node->notify[0]) { char buf[1024]; int ret = snprintf(buf, 1024, "[subscription ack] %.*s => new node: %.*s", usr->keylen, usr->key, usr->address_len, usr->address); - if (ret > 0 && ret < 1024) uwsgi_notify_msg(node->notify, buf, ret); + if (ret > 0 && ret < 1024) + uwsgi_notify_msg(node->notify, buf, ret); } return node; } else { + current_slot = uwsgi_malloc(sizeof(struct uwsgi_subscribe_slot)); #ifdef UWSGI_SSL - FILE *kf = NULL; - if (uwsgi.subscriptions_sign_check_dir) { - if (usr->sign_len == 0 || usr->base_len == 0) return NULL; - if (usr->unix_check < (uwsgi_now() - (time_t) uwsgi.subscriptions_sign_check_tolerance)) { - uwsgi_log("[uwsgi-subscription for pid %d] invalid (sniffed ?) packet sent for slot: %.*s node: %.*s unix_check: %lu\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, (unsigned long) usr->unix_check); - return NULL; - } - char *keyfile = uwsgi_sanitize_cert_filename(uwsgi.subscriptions_sign_check_dir, usr->key, usr->keylen); - kf = fopen(keyfile, "r"); - free(keyfile); - if (!kf) - return NULL; - + if (uwsgi.subscriptions_sign_check_dir && !subscription_new_sign_ctx(current_slot, usr)) { + free(current_slot); + return NULL; } #endif - current_slot = uwsgi_malloc(sizeof(struct uwsgi_subscribe_slot)); uint32_t hash = djb33x_hash(usr->key, usr->keylen); int hash_key = hash % 0xffff; current_slot->hash = hash_key; -#ifdef UWSGI_SSL - if (uwsgi.subscriptions_sign_check_dir) { - current_slot->sign_public_key = PEM_read_PUBKEY(kf, NULL, NULL, NULL); - fclose(kf); - if (!current_slot->sign_public_key) { - uwsgi_log("unable to load public key for %.*s\n", usr->keylen, usr->key); - free(current_slot); - return NULL; - } - current_slot->sign_ctx = EVP_MD_CTX_create(); - if (!current_slot->sign_ctx) { - uwsgi_log("unable to initialize EVP context for %.*s\n", usr->keylen, usr->key); - EVP_PKEY_free(current_slot->sign_public_key); - free(current_slot); - return NULL; - } - - if (!uwsgi_subscription_sign_check(current_slot, usr)) { - EVP_PKEY_free(current_slot->sign_public_key); - EVP_MD_CTX_destroy(current_slot->sign_ctx); - free(current_slot); - return NULL; - } - } -#endif current_slot->keylen = usr->keylen; memcpy(current_slot->key, usr->key, usr->keylen); if (uwsgi.subscriptions_credentials_check_dir) { @@ -598,9 +568,10 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo if (current_slot->nodes->notify[0]) { char buf[1024]; - int ret = snprintf(buf, 1024, "[subscription ack] %.*s => new node: %.*s", usr->keylen, usr->key, usr->address_len, usr->address); - if (ret > 0 && ret < 1024) uwsgi_notify_msg(current_slot->nodes->notify, buf, ret); - } + int ret = snprintf(buf, 1024, "[subscription ack] %.*s => new node: %.*s", usr->keylen, usr->key, usr->address_len, usr->address); + if (ret > 0 && ret < 1024) + uwsgi_notify_msg(current_slot->nodes->notify, buf, ret); + } return current_slot->nodes; } @@ -608,24 +579,24 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo static void send_subscription(int sfd, char *host, char *message, uint16_t message_size) { - int fd = sfd; - struct sockaddr_in udp_addr; - struct sockaddr_un un_addr; - ssize_t ret; + int fd = sfd; + struct sockaddr_in udp_addr; + struct sockaddr_un un_addr; + ssize_t ret; - char *udp_port = strchr(host, ':'); + char *udp_port = strchr(host, ':'); if (fd == -1) { if (udp_port) { - fd = socket(AF_INET, SOCK_DGRAM, 0); + fd = socket(AF_INET, SOCK_DGRAM, 0); } else { - fd = socket(AF_UNIX, SOCK_DGRAM, 0); + fd = socket(AF_UNIX, SOCK_DGRAM, 0); + } + if (fd < 0) { + uwsgi_error("send_subscription()/socket()"); + return; } - if (fd < 0) { - uwsgi_error("send_subscription()/socket()"); - return; - } uwsgi_socket_nb(fd); } else if (fd == -2) { @@ -634,115 +605,130 @@ static void send_subscription(int sfd, char *host, char *message, uint16_t messa if (udp_port) { if (inet_fd == -1) { inet_fd = socket(AF_INET, SOCK_DGRAM, 0); - if (inet_fd < 0) { - uwsgi_error("send_subscription()/socket()"); - return; - } + if (inet_fd < 0) { + uwsgi_error("send_subscription()/socket()"); + return; + } uwsgi_socket_nb(inet_fd); } fd = inet_fd; } else { if (unix_fd == -1) { - unix_fd = socket(AF_UNIX, SOCK_DGRAM, 0); - if (unix_fd < 0) { - uwsgi_error("send_subscription()/socket()"); - return; - } - uwsgi_socket_nb(unix_fd); - } - fd = unix_fd; + unix_fd = socket(AF_UNIX, SOCK_DGRAM, 0); + if (unix_fd < 0) { + uwsgi_error("send_subscription()/socket()"); + return; + } + uwsgi_socket_nb(unix_fd); + } + fd = unix_fd; } } - if (udp_port) { - udp_port[0] = 0; - memset(&udp_addr, 0, sizeof(struct sockaddr_in)); - udp_addr.sin_family = AF_INET; - udp_addr.sin_port = htons(atoi(udp_port + 1)); - udp_addr.sin_addr.s_addr = inet_addr(host); - ret = sendto(fd, message, message_size, 0, (struct sockaddr *) &udp_addr, sizeof(udp_addr)); - udp_port[0] = ':'; - } - else { - memset(&un_addr, 0, sizeof(struct sockaddr_un)); - un_addr.sun_family = AF_UNIX; - // use 102 as the magic number - strncat(un_addr.sun_path, host, 102); + if (udp_port) { + udp_port[0] = 0; + memset(&udp_addr, 0, sizeof(struct sockaddr_in)); + udp_addr.sin_family = AF_INET; + udp_addr.sin_port = htons(atoi(udp_port + 1)); + udp_addr.sin_addr.s_addr = inet_addr(host); + ret = sendto(fd, message, message_size, 0, (struct sockaddr *) &udp_addr, sizeof(udp_addr)); + udp_port[0] = ':'; + } + else { + memset(&un_addr, 0, sizeof(struct sockaddr_un)); + un_addr.sun_family = AF_UNIX; + // use 102 as the magic number + strncat(un_addr.sun_path, host, 102); if (uwsgi.subscriptions_use_credentials) { // could be useless as internally the socket could add them automagically - ret = uwsgi_pass_cred2(fd, message, message_size, (struct sockaddr *) &un_addr, sizeof(un_addr)); + ret = uwsgi_pass_cred2(fd, message, message_size, (struct sockaddr *) &un_addr, sizeof(un_addr)); } else { - ret = sendto(fd, message, message_size, 0, (struct sockaddr *) &un_addr, sizeof(un_addr)); + ret = sendto(fd, message, message_size, 0, (struct sockaddr *) &un_addr, sizeof(un_addr)); } - } + } - if (ret < 0) { - uwsgi_error("send_subscription()/sendto()"); - } + if (ret < 0) { + uwsgi_error("send_subscription()/sendto()"); + } if (sfd == -1) - close(fd); + close(fd); } static struct uwsgi_buffer *uwsgi_subscription_ub(char *key, size_t keysize, uint8_t modifier1, uint8_t modifier2, uint8_t cmd, char *socket_name, char *sign, char *sni_key, char *sni_crt, char *sni_ca) { - struct uwsgi_buffer *ub = uwsgi_buffer_new(4096); + struct uwsgi_buffer *ub = uwsgi_buffer_new(4096); - // make space for uwsgi header - ub->pos = 4; + // make space for uwsgi header + ub->pos = 4; - if (uwsgi_buffer_append_keyval(ub, "key", 3, key, keysize)) goto end; - if (uwsgi_buffer_append_keyval(ub, "address", 7, socket_name, strlen(socket_name))) goto end; - if (uwsgi_buffer_append_keynum(ub, "modifier1", 9, modifier1)) goto end; - if (uwsgi_buffer_append_keynum(ub, "modifier2", 9, modifier2)) goto end; - if (uwsgi_buffer_append_keynum(ub, "cores", 5, uwsgi.numproc * uwsgi.cores)) goto end; - if (uwsgi_buffer_append_keynum(ub, "load", 4, uwsgi.shared->load)) goto end; - if (uwsgi.auto_weight) { - if (uwsgi_buffer_append_keynum(ub, "weight", 6, uwsgi.numproc * uwsgi.cores )) goto end; - } - else { - if (uwsgi_buffer_append_keynum(ub, "weight", 6, uwsgi.weight )) goto end; - } + if (uwsgi_buffer_append_keyval(ub, "key", 3, key, keysize)) + goto end; + if (uwsgi_buffer_append_keyval(ub, "address", 7, socket_name, strlen(socket_name))) + goto end; + if (uwsgi_buffer_append_keynum(ub, "modifier1", 9, modifier1)) + goto end; + if (uwsgi_buffer_append_keynum(ub, "modifier2", 9, modifier2)) + goto end; + if (uwsgi_buffer_append_keynum(ub, "cores", 5, uwsgi.numproc * uwsgi.cores)) + goto end; + if (uwsgi_buffer_append_keynum(ub, "load", 4, uwsgi.shared->load)) + goto end; + if (uwsgi.auto_weight) { + if (uwsgi_buffer_append_keynum(ub, "weight", 6, uwsgi.numproc * uwsgi.cores)) + goto end; + } + else { + if (uwsgi_buffer_append_keynum(ub, "weight", 6, uwsgi.weight)) + goto end; + } - if (sni_key) { - if (uwsgi_buffer_append_keyval(ub, "sni_key", 7, sni_key, strlen(sni_key))) goto end; - } + if (sni_key) { + if (uwsgi_buffer_append_keyval(ub, "sni_key", 7, sni_key, strlen(sni_key))) + goto end; + } - if (sni_crt) { - if (uwsgi_buffer_append_keyval(ub, "sni_crt", 7, sni_crt, strlen(sni_crt))) goto end; - } + if (sni_crt) { + if (uwsgi_buffer_append_keyval(ub, "sni_crt", 7, sni_crt, strlen(sni_crt))) + goto end; + } - if (sni_ca) { - if (uwsgi_buffer_append_keyval(ub, "sni_ca", 6, sni_ca, strlen(sni_ca))) goto end; - } + if (sni_ca) { + if (uwsgi_buffer_append_keyval(ub, "sni_ca", 6, sni_ca, strlen(sni_ca))) + goto end; + } if (uwsgi.subscription_notify_socket) { - if (uwsgi_buffer_append_keyval(ub, "notify", 6, uwsgi.subscription_notify_socket, strlen(uwsgi.subscription_notify_socket))) goto end; + if (uwsgi_buffer_append_keyval(ub, "notify", 6, uwsgi.subscription_notify_socket, strlen(uwsgi.subscription_notify_socket))) + goto end; } else if (uwsgi.notify_socket_fd > -1 && uwsgi.notify_socket) { - if (uwsgi_buffer_append_keyval(ub, "notify", 6, uwsgi.notify_socket, strlen(uwsgi.notify_socket))) goto end; + if (uwsgi_buffer_append_keyval(ub, "notify", 6, uwsgi.notify_socket, strlen(uwsgi.notify_socket))) + goto end; } #ifdef UWSGI_SSL - if (sign) { - if (uwsgi_buffer_append_keynum(ub, "unix", 4, (uwsgi_now() + (time_t) cmd) )) goto end; + if (sign) { + if (uwsgi_buffer_append_keynum(ub, "unix", 4, (uwsgi_now() + (time_t) cmd))) + goto end; - unsigned int signature_len = 0; - char *signature = uwsgi_rsa_sign(sign, ub->buf + 4, ub->pos - 4, &signature_len); - if (signature && signature_len > 0) { - if (uwsgi_buffer_append_keyval(ub, "sign", 4, signature, signature_len)) { - free(signature); - goto end; - } - free(signature); - } - } + unsigned int signature_len = 0; + char *signature = uwsgi_rsa_sign(sign, ub->buf + 4, ub->pos - 4, &signature_len); + if (signature && signature_len > 0) { + if (uwsgi_buffer_append_keyval(ub, "sign", 4, signature, signature_len)) { + free(signature); + goto end; + } + free(signature); + } + } #endif - // add uwsgi header - if (uwsgi_buffer_set_uh(ub, 224, cmd)) goto end; - + // add uwsgi header + if (uwsgi_buffer_set_uh(ub, 224, cmd)) + goto end; + return ub; end: @@ -752,19 +738,20 @@ end: void uwsgi_send_subscription_from_fd(int fd, char *udp_address, char *key, size_t keysize, uint8_t modifier1, uint8_t modifier2, uint8_t cmd, char *socket_name, char *sign, char *sni_key, char *sni_crt, char *sni_ca) { - if (socket_name == NULL && !uwsgi.sockets) - return; + if (socket_name == NULL && !uwsgi.sockets) + return; - if (!socket_name) { - socket_name = uwsgi.sockets->name; - } + if (!socket_name) { + socket_name = uwsgi.sockets->name; + } - struct uwsgi_buffer *ub = uwsgi_subscription_ub(key, keysize, modifier1, modifier2, cmd, socket_name, sign, sni_key, sni_crt, sni_ca); + struct uwsgi_buffer *ub = uwsgi_subscription_ub(key, keysize, modifier1, modifier2, cmd, socket_name, sign, sni_key, sni_crt, sni_ca); - if (!ub) return; + if (!ub) + return; - send_subscription(fd, udp_address, ub->buf, ub->pos); - uwsgi_buffer_destroy(ub); + send_subscription(fd, udp_address, ub->buf, ub->pos); + uwsgi_buffer_destroy(ub); } @@ -773,19 +760,63 @@ void uwsgi_send_subscription(char *udp_address, char *key, size_t keysize, uint8 } #ifdef UWSGI_SSL -int uwsgi_subscription_sign_check(struct uwsgi_subscribe_slot *slot, struct uwsgi_subscribe_req *usr) { - +static int subscription_is_safe(struct uwsgi_subscribe_req *usr) { struct uwsgi_string_list *usl = NULL; - uwsgi_foreach(usl, uwsgi.subscriptions_sign_skip_uid) { - if (usl->custom == 0) { - usl->custom = atoi(usl->value); - } - if (usr->uid > 0 && usr->uid == (uid_t) usl->custom) { - return 1; - } + uwsgi_foreach(usl, uwsgi.subscriptions_sign_skip_uid) { + if (usl->custom == 0) { + usl->custom = atoi(usl->value); + } + if (usr->uid > 0 && usr->uid == (uid_t) usl->custom) { + return 1; + } + } + return 0; +} +static int subscription_new_sign_ctx(struct uwsgi_subscribe_slot *slot, struct uwsgi_subscribe_req *usr) { + if (subscription_is_safe(usr)) return 1; + + if (usr->sign_len == 0 || usr->base_len == 0) + return 0; + + if (usr->unix_check < (uwsgi_now() - (time_t) uwsgi.subscriptions_sign_check_tolerance)) { + uwsgi_log("[uwsgi-subscription for pid %d] invalid (sniffed ?) packet sent for slot: %.*s node: %.*s unix_check: %lu\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, (unsigned long) usr->unix_check); + return 0; + } + + char *keyfile = uwsgi_sanitize_cert_filename(uwsgi.subscriptions_sign_check_dir, usr->key, usr->keylen); + FILE *kf = fopen(keyfile, "r"); + free(keyfile); + if (!kf) return 0; + slot->sign_public_key = PEM_read_PUBKEY(kf, NULL, NULL, NULL); + fclose(kf); + if (!slot->sign_public_key) { + uwsgi_log("unable to load public key for %.*s\n", usr->keylen, usr->key); + return 0; + } + slot->sign_ctx = EVP_MD_CTX_create(); + if (!slot->sign_ctx) { + uwsgi_log("unable to initialize EVP context for %.*s\n", usr->keylen, usr->key); + EVP_PKEY_free(slot->sign_public_key); + return 0; } - if (usr->sign_len == 0 || usr->base_len == 0) return 0; + if (!uwsgi_subscription_sign_check(slot, usr)) { + EVP_PKEY_free(slot->sign_public_key); + EVP_MD_CTX_destroy(slot->sign_ctx); + return 0; + } + + return 1; +} +int uwsgi_subscription_sign_check(struct uwsgi_subscribe_slot *slot, struct uwsgi_subscribe_req *usr) { + if (subscription_is_safe(usr)) return 1; + + if (usr->sign_len == 0 || usr->base_len == 0) + return 0; + + if (!slot->sign_ctx) { + if (!subscription_new_sign_ctx(slot, usr)) return 0; + } if (EVP_VerifyInit_ex(slot->sign_ctx, uwsgi.subscriptions_sign_check_md, NULL) == 0) { ERR_print_errors_fp(stderr); @@ -827,133 +858,133 @@ struct uwsgi_subscribe_slot **uwsgi_subscription_init_ht() { void uwsgi_subscribe(char *subscription, uint8_t cmd) { - size_t subfile_size; - size_t i; - char *key = NULL; - int keysize = 0; - char *modifier1 = NULL; - int modifier1_len = 0; - char *socket_name = NULL; - char *udp_address = subscription; - char *udp_port = NULL; - char *subscription_key = NULL; - char *sign = NULL; + size_t subfile_size; + size_t i; + char *key = NULL; + int keysize = 0; + char *modifier1 = NULL; + int modifier1_len = 0; + char *socket_name = NULL; + char *udp_address = subscription; + char *udp_port = NULL; + char *subscription_key = NULL; + char *sign = NULL; - // check for explicit socket_name - char *equal = strchr(subscription, '='); - if (equal) { - socket_name = subscription; - if (socket_name[0] == '=') { - equal = strchr(socket_name + 1, '='); - if (!equal) - return; - *equal = '\0'; - struct uwsgi_socket *us = uwsgi_get_shared_socket_by_num(atoi(socket_name + 1)); - if (!us) - return; - socket_name = us->name; - } - *equal = '\0'; - udp_address = equal + 1; - } + // check for explicit socket_name + char *equal = strchr(subscription, '='); + if (equal) { + socket_name = subscription; + if (socket_name[0] == '=') { + equal = strchr(socket_name + 1, '='); + if (!equal) + return; + *equal = '\0'; + struct uwsgi_socket *us = uwsgi_get_shared_socket_by_num(atoi(socket_name + 1)); + if (!us) + return; + socket_name = us->name; + } + *equal = '\0'; + udp_address = equal + 1; + } - // check for unix socket - if (udp_address[0] != '/') { - udp_port = strchr(udp_address, ':'); - if (!udp_port) { - if (equal) - *equal = '='; - return; - } - subscription_key = strchr(udp_port + 1, ':'); - } - else { - subscription_key = strchr(udp_address + 1, ':'); - } + // check for unix socket + if (udp_address[0] != '/') { + udp_port = strchr(udp_address, ':'); + if (!udp_port) { + if (equal) + *equal = '='; + return; + } + subscription_key = strchr(udp_port + 1, ':'); + } + else { + subscription_key = strchr(udp_address + 1, ':'); + } - if (!subscription_key) { - if (equal) - *equal = '='; - return; - } + if (!subscription_key) { + if (equal) + *equal = '='; + return; + } - udp_address = uwsgi_concat2n(udp_address, subscription_key - udp_address, "", 0); + udp_address = uwsgi_concat2n(udp_address, subscription_key - udp_address, "", 0); if (subscription_key[1] == '@') { - if (!uwsgi_file_exists(subscription_key + 2)) - goto clear; - char *lines = uwsgi_open_and_read(subscription_key + 2, &subfile_size, 1, NULL); - if (subfile_size > 0) { - key = lines; - for (i = 0; i < subfile_size; i++) { - if (lines[i] == 0) { - if (keysize > 0) { - if (key[0] != '#' && key[0] != '\n') { - modifier1 = strchr(key, ','); - if (modifier1) { - modifier1[0] = 0; - modifier1++; - modifier1_len = strlen(modifier1); - keysize = strlen(key); - } - uwsgi_send_subscription(udp_address, key, keysize, uwsgi_str_num(modifier1, modifier1_len), 0, cmd, socket_name, sign, NULL, NULL, NULL); - modifier1 = NULL; - modifier1_len = 0; - } - } - break; - } - else if (lines[i] == '\n') { - if (keysize > 0) { - if (key[0] != '#' && key[0] != '\n') { - lines[i] = 0; - modifier1 = strchr(key, ','); - if (modifier1) { - modifier1[0] = 0; - modifier1++; - modifier1_len = strlen(modifier1); - keysize = strlen(key); - } - uwsgi_send_subscription(udp_address, key, keysize, uwsgi_str_num(modifier1, modifier1_len), 0, cmd, socket_name, sign, NULL, NULL, NULL); - modifier1 = NULL; - modifier1_len = 0; - lines[i] = '\n'; - } - } - key = lines + i + 1; - keysize = 0; - continue; - } - keysize++; - } - } - free(lines); - } - else { - modifier1 = strchr(subscription_key + 1, ','); - if (modifier1) { - modifier1[0] = 0; - modifier1++; + if (!uwsgi_file_exists(subscription_key + 2)) + goto clear; + char *lines = uwsgi_open_and_read(subscription_key + 2, &subfile_size, 1, NULL); + if (subfile_size > 0) { + key = lines; + for (i = 0; i < subfile_size; i++) { + if (lines[i] == 0) { + if (keysize > 0) { + if (key[0] != '#' && key[0] != '\n') { + modifier1 = strchr(key, ','); + if (modifier1) { + modifier1[0] = 0; + modifier1++; + modifier1_len = strlen(modifier1); + keysize = strlen(key); + } + uwsgi_send_subscription(udp_address, key, keysize, uwsgi_str_num(modifier1, modifier1_len), 0, cmd, socket_name, sign, NULL, NULL, NULL); + modifier1 = NULL; + modifier1_len = 0; + } + } + break; + } + else if (lines[i] == '\n') { + if (keysize > 0) { + if (key[0] != '#' && key[0] != '\n') { + lines[i] = 0; + modifier1 = strchr(key, ','); + if (modifier1) { + modifier1[0] = 0; + modifier1++; + modifier1_len = strlen(modifier1); + keysize = strlen(key); + } + uwsgi_send_subscription(udp_address, key, keysize, uwsgi_str_num(modifier1, modifier1_len), 0, cmd, socket_name, sign, NULL, NULL, NULL); + modifier1 = NULL; + modifier1_len = 0; + lines[i] = '\n'; + } + } + key = lines + i + 1; + keysize = 0; + continue; + } + keysize++; + } + } + free(lines); + } + else { + modifier1 = strchr(subscription_key + 1, ','); + if (modifier1) { + modifier1[0] = 0; + modifier1++; - sign = strchr(modifier1 + 1, ','); - if (sign) { - *sign = 0; - sign++; - } - modifier1_len = strlen(modifier1); - } + sign = strchr(modifier1 + 1, ','); + if (sign) { + *sign = 0; + sign++; + } + modifier1_len = strlen(modifier1); + } - uwsgi_send_subscription(udp_address, subscription_key + 1, strlen(subscription_key + 1), uwsgi_str_num(modifier1, modifier1_len), 0, cmd, socket_name, sign, NULL, NULL, NULL); - if (modifier1) - modifier1[-1] = ','; - if (sign) - sign[-1] = ','; - } + uwsgi_send_subscription(udp_address, subscription_key + 1, strlen(subscription_key + 1), uwsgi_str_num(modifier1, modifier1_len), 0, cmd, socket_name, sign, NULL, NULL, NULL); + if (modifier1) + modifier1[-1] = ','; + if (sign) + sign[-1] = ','; + } clear: - if (equal) - *equal = '='; - free(udp_address); + if (equal) + *equal = '='; + free(udp_address); } @@ -972,27 +1003,16 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) { char *s2_sni_crt = NULL; char *s2_sni_ca = NULL; - if (uwsgi_kvlist_parse(arg, strlen(arg), ',', '=', - "server", &s2_server, - "key", &s2_key, - "socket", &s2_socket, - "addr", &s2_addr, - "weight", &s2_weight, - "modifier1", &s2_modifier1, - "modifier2", &s2_modifier2, - "sign", &s2_sign, - "check", &s2_check, - "sni_key", &s2_sni_key, - "sni_crt", &s2_sni_crt, - "sni_ca", &s2_sni_ca, - NULL)) { + if (uwsgi_kvlist_parse(arg, strlen(arg), ',', '=', "server", &s2_server, "key", &s2_key, "socket", &s2_socket, "addr", &s2_addr, "weight", &s2_weight, "modifier1", &s2_modifier1, "modifier2", &s2_modifier2, "sign", &s2_sign, "check", &s2_check, "sni_key", &s2_sni_key, "sni_crt", &s2_sni_crt, "sni_ca", &s2_sni_ca, NULL)) { return; } - if (!s2_server || !s2_key) goto end; + if (!s2_server || !s2_key) + goto end; if (s2_check) { - if (uwsgi_file_exists(s2_check)) goto end; + if (uwsgi_file_exists(s2_check)) + goto end; } if (s2_weight) { @@ -1022,39 +1042,48 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) { uwsgi_send_subscription(s2_server, s2_key, strlen(s2_key), modifier1, modifier2, cmd, s2_addr, s2_sign, s2_sni_key, s2_sni_crt, s2_sni_ca); end: - if (s2_server) free(s2_server); - if (s2_key) free(s2_key); - if (s2_socket) free(s2_socket); - if (s2_addr) free(s2_addr); - if (s2_weight) free(s2_weight); - if (s2_modifier1) free(s2_modifier1); - if (s2_modifier2) free(s2_modifier2); - if (s2_sign) free(s2_sign); - if (s2_check) free(s2_check); + if (s2_server) + free(s2_server); + if (s2_key) + free(s2_key); + if (s2_socket) + free(s2_socket); + if (s2_addr) + free(s2_addr); + if (s2_weight) + free(s2_weight); + if (s2_modifier1) + free(s2_modifier1); + if (s2_modifier2) + free(s2_modifier2); + if (s2_sign) + free(s2_sign); + if (s2_check) + free(s2_check); } void uwsgi_subscribe_all(uint8_t cmd, int verbose) { - if (uwsgi.subscriptions_blocked) return; + if (uwsgi.subscriptions_blocked) + return; // -- subscribe struct uwsgi_string_list *subscriptions = uwsgi.subscriptions; - while (subscriptions) { + while (subscriptions) { if (verbose) { - uwsgi_log("%s %s\n", cmd ? "unsubscribing from" : "subscribing to", subscriptions->value); + uwsgi_log("%s %s\n", cmd ? "unsubscribing from" : "subscribing to", subscriptions->value); } - uwsgi_subscribe(subscriptions->value, cmd); - subscriptions = subscriptions->next; - } + uwsgi_subscribe(subscriptions->value, cmd); + subscriptions = subscriptions->next; + } // --subscribe2 subscriptions = uwsgi.subscriptions2; - while (subscriptions) { - if (verbose) { - uwsgi_log("%s %s\n", cmd ? "unsubscribing from" : "subscribing to", subscriptions->value); - } - uwsgi_subscribe2(subscriptions->value, cmd); - subscriptions = subscriptions->next; - } + while (subscriptions) { + if (verbose) { + uwsgi_log("%s %s\n", cmd ? "unsubscribing from" : "subscribing to", subscriptions->value); + } + uwsgi_subscribe2(subscriptions->value, cmd); + subscriptions = subscriptions->next; + } } -