From a3635bc45ff40c1d511ebb8971686ef795439f0d Mon Sep 17 00:00:00 2001 From: "roberto@precise64" Date: Thu, 24 May 2012 14:30:27 +0200 Subject: [PATCH] another master refactoring --- cluster.c | 44 +++++++++ master.c | 290 ++++++++++++++++++++++++------------------------------ uwsgi.h | 4 +- 3 files changed, 178 insertions(+), 160 deletions(-) diff --git a/cluster.c b/cluster.c index 63127021..2c2cdb86 100644 --- a/cluster.c +++ b/cluster.c @@ -359,5 +359,49 @@ void manage_cluster_announce(char *key, uint16_t keylen, char *val, uint16_t val } } +void manage_cluster_message(char *cluster_opt_buf, int cluster_opt_size) { + + struct uwsgi_cluster_node nucn; + + switch (uwsgi.wsgi_requests[0]->uh.modifier1) { + case 95: + memset(&nucn, 0, sizeof(struct uwsgi_cluster_node)); + +#ifdef __BIG_ENDIAN__ + uwsgi.wsgi_requests[0]->uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize); +#endif + uwsgi_hooked_parse(uwsgi.wsgi_requests[0]->buffer, uwsgi.wsgi_requests[0]->uh.pktsize, manage_cluster_announce, &nucn); + if (nucn.name[0] != 0) { + uwsgi_cluster_add_node(&nucn, CLUSTER_NODE_DYNAMIC); + } + break; + case 96: +#ifdef __BIG_ENDIAN__ + uwsgi.wsgi_requests[0]->uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize); +#endif + uwsgi_log_verbose("%.*s\n", uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); + break; + case 98: + if (kill(getpid(), SIGHUP)) { + uwsgi_error("kill()"); + } + break; + case 99: + if (uwsgi.cluster_nodes) + break; + if (uwsgi.wsgi_requests[0]->uh.modifier2 == 0) { + uwsgi_log("requested configuration data, sending %d bytes\n", cluster_opt_size); + sendto(uwsgi.cluster_fd, cluster_opt_buf, cluster_opt_size, 0, (struct sockaddr *) &uwsgi.mc_cluster_addr, sizeof(uwsgi.mc_cluster_addr)); + } + break; + case 73: +#ifdef __BIG_ENDIAN__ + uwsgi.wsgi_requests[0]->uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize); +#endif + uwsgi_log_verbose("[uWSGI cluster %s] new node available: %.*s\n", uwsgi.cluster, uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); + break; + } +} + #endif diff --git a/master.c b/master.c index 4c65d045..1984b7e5 100644 --- a/master.c +++ b/master.c @@ -29,12 +29,12 @@ void suspend_resume_them_all(int signum) { } // subscribe/unsubscribe if needed - struct uwsgi_string_list *subscriptions = uwsgi.subscriptions; - while(subscriptions) { - uwsgi_log("%s %s\n", suspend ? "unsubscribing from" : "subscribing to", subscriptions->value); - uwsgi_subscribe(subscriptions->value, suspend); - subscriptions = subscriptions->next; - } + struct uwsgi_string_list *subscriptions = uwsgi.subscriptions; + while (subscriptions) { + uwsgi_log("%s %s\n", suspend ? "unsubscribing from" : "subscribing to", subscriptions->value); + uwsgi_subscribe(subscriptions->value, suspend); + subscriptions = subscriptions->next; + } for (i = 1; i <= uwsgi.numproc; i++) { @@ -99,13 +99,13 @@ void *logger_thread_loop(void *noarg) { // block all signals sigset_t smask; - sigfillset(&smask); - pthread_sigmask(SIG_BLOCK, &smask, NULL); + sigfillset(&smask); + pthread_sigmask(SIG_BLOCK, &smask, NULL); logpoll.events = POLLIN; logpoll.fd = uwsgi.shared->worker_log_pipe[0]; - for(;;) { + for (;;) { int ret = poll(&logpoll, 1, -1); if (ret > 0 && logpoll.revents & POLLIN) { pthread_mutex_lock(&uwsgi.threaded_logger_lock); @@ -121,27 +121,28 @@ void *cache_sweeper_loop(void *noarg) { int i; // block all signals - sigset_t smask; - sigfillset(&smask); - pthread_sigmask(SIG_BLOCK, &smask, NULL); + sigset_t smask; + sigfillset(&smask); + pthread_sigmask(SIG_BLOCK, &smask, NULL); - if (!uwsgi.cache_expire_freq) uwsgi.cache_expire_freq = 3; + if (!uwsgi.cache_expire_freq) + uwsgi.cache_expire_freq = 3; // remove expired cache items TODO use rb_tree timeouts - for(;;) { + for (;;) { sleep(uwsgi.cache_expire_freq); uint64_t freed_items = 0; // skip the first slot - for (i = 1; i < (int) uwsgi.cache_max_items; i++) { - uwsgi_wlock(uwsgi.cache_lock); - if (uwsgi.cache_items[i].expires) { - if (uwsgi.cache_items[i].expires < (uint64_t) uwsgi.current_time) { - uwsgi_cache_del(NULL, 0, i); + for (i = 1; i < (int) uwsgi.cache_max_items; i++) { + uwsgi_wlock(uwsgi.cache_lock); + if (uwsgi.cache_items[i].expires) { + if (uwsgi.cache_items[i].expires < (uint64_t) uwsgi.current_time) { + uwsgi_cache_del(NULL, 0, i); freed_items++; - } - } - uwsgi_rwunlock(uwsgi.cache_lock); - } + } + } + uwsgi_rwunlock(uwsgi.cache_lock); + } if (uwsgi.cache_report_freed_items && freed_items > 0) { uwsgi_log("freed %llu cache items\n", (unsigned long long) freed_items); } @@ -236,7 +237,7 @@ void uwsgi_subscribe(char *subscription, uint8_t cmd) { modifier1[-1] = ','; } - clear: +clear: free(udp_address); @@ -331,12 +332,11 @@ int master_loop(char **argv, char **environ) { #ifdef UWSGI_MULTICAST char *cluster_opt_buf = NULL; - int cluster_opt_size = 4; + int cluster_opt_size = 4; char *cptrbuf; uint16_t ustrlen; struct uwsgi_header *uh; - struct uwsgi_cluster_node nucn; #endif #endif @@ -419,9 +419,9 @@ int master_loop(char **argv, char **environ) { if (uwsgi.cache_max_items > 0 && !uwsgi.cache_no_expire) { if (pthread_create(&cache_sweeper, NULL, cache_sweeper_loop, NULL)) { - uwsgi_error("pthread_create()"); + uwsgi_error("pthread_create()"); uwsgi_log("unable to run the cache sweeper !!!\n"); - } + } else { uwsgi_log("cache sweeper thread enabled\n"); } @@ -630,7 +630,7 @@ int master_loop(char **argv, char **environ) { if (uwsgi.requested_cheaper_algo) { uwsgi.cheaper_algo = NULL; struct uwsgi_cheaper_algo *uca = uwsgi.cheaper_algos; - while(uca) { + while (uca) { if (!strcmp(uca->name, uwsgi.requested_cheaper_algo)) { uwsgi.cheaper_algo = uca->func; break; @@ -650,6 +650,7 @@ int master_loop(char **argv, char **environ) { for (;;) { //uwsgi_log("ready_to_reload %d %d\n", ready_to_reload, uwsgi.numproc); + // run master_cycle hook for every plugin for (i = 0; i < uwsgi.gp_cnt; i++) { if (uwsgi.gp[i]->master_cycle) { uwsgi.gp[i]->master_cycle(); @@ -737,7 +738,8 @@ int master_loop(char **argv, char **environ) { // cheaper management if (uwsgi.cheaper && !uwsgi.cheap && !uwsgi.to_heaven && !uwsgi.to_hell && !uwsgi.workers[0].suspended) { - if (!uwsgi_calc_cheaper()) return 0; + if (!uwsgi_calc_cheaper()) + return 0; } @@ -831,6 +833,8 @@ int master_loop(char **argv, char **environ) { } } } + + // wait for event rlen = event_queue_wait(uwsgi.master_queue, check_interval, &interesting_fd); if (rlen == 0) { @@ -1012,44 +1016,8 @@ int master_loop(char **argv, char **environ) { goto health_cycle; } - switch (uwsgi.wsgi_requests[0]->uh.modifier1) { - case 95: - memset(&nucn, 0, sizeof(struct uwsgi_cluster_node)); + manage_cluster_message(cluster_opt_buf, cluster_opt_size); -#ifdef __BIG_ENDIAN__ - uwsgi.wsgi_requests[0]->uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize); -#endif - uwsgi_hooked_parse(uwsgi.wsgi_requests[0]->buffer, uwsgi.wsgi_requests[0]->uh.pktsize, manage_cluster_announce, &nucn); - if (nucn.name[0] != 0) { - uwsgi_cluster_add_node(&nucn, CLUSTER_NODE_DYNAMIC); - } - break; - case 96: -#ifdef __BIG_ENDIAN__ - uwsgi.wsgi_requests[0]->uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize); -#endif - uwsgi_log_verbose("%.*s\n", uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); - break; - case 98: - if (kill(getpid(), SIGHUP)) { - uwsgi_error("kill()"); - } - break; - case 99: - if (uwsgi.cluster_nodes) - break; - if (uwsgi.wsgi_requests[0]->uh.modifier2 == 0) { - uwsgi_log("requested configuration data, sending %d bytes\n", cluster_opt_size); - sendto(uwsgi.cluster_fd, cluster_opt_buf, cluster_opt_size, 0, (struct sockaddr *) &uwsgi.mc_cluster_addr, sizeof(uwsgi.mc_cluster_addr)); - } - break; - case 73: -#ifdef __BIG_ENDIAN__ - uwsgi.wsgi_requests[0]->uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize); -#endif - uwsgi_log_verbose("[uWSGI cluster %s] new node available: %.*s\n", uwsgi.cluster, uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); - break; - } goto health_cycle; } #endif @@ -1163,7 +1131,7 @@ int master_loop(char **argv, char **environ) { } - health_cycle: +health_cycle: now = time(NULL); if (now - uwsgi.current_time < 1) { continue; @@ -1345,7 +1313,7 @@ int master_loop(char **argv, char **environ) { } #ifdef UWSGI_SPOOLER struct uwsgi_spooler *uspool = uwsgi.spoolers; - while(uspool) { + while (uspool) { if (uspool->harakiri > 0 && uspool->harakiri < (time_t) uwsgi.current_time) { uwsgi_log("*** HARAKIRI ON THE SPOOLER (pid: %d) ***\n", uspool->pid); kill(uspool->pid, SIGKILL); @@ -1373,7 +1341,7 @@ int master_loop(char **argv, char **environ) { } // resubscribe every 10 cycles by default - if ((uwsgi.subscriptions && ((uwsgi.master_cycles % uwsgi.subscribe_freq) == 0 || uwsgi.master_cycles == 1)) && !uwsgi.to_heaven && !uwsgi.to_hell && !uwsgi.workers[0].suspended) { + if ((uwsgi.subscriptions && ((uwsgi.master_cycles % uwsgi.subscribe_freq) == 0 || uwsgi.master_cycles == 1)) && !uwsgi.to_heaven && !uwsgi.to_hell && !uwsgi.workers[0].suspended) { struct uwsgi_string_list *subscriptions = uwsgi.subscriptions; while (subscriptions) { uwsgi_subscribe(subscriptions->value, 0); @@ -1408,31 +1376,36 @@ int master_loop(char **argv, char **environ) { } - if (diedpid > 0) { - // check for deadlocks first - struct uwsgi_lock_item *uli = uwsgi.registered_locks; - while(uli) { - if (!uli->can_deadlock) goto nextlock; - pid_t locked_pid = 0; + // no one died + if (diedpid <= 0) + continue; + + // check for deadlocks first + struct uwsgi_lock_item *uli = uwsgi.registered_locks; + while (uli) { + if (!uli->can_deadlock) + goto nextlock; + pid_t locked_pid = 0; + if (uli->rw) { + locked_pid = uwsgi_rwlock_check(uli); + } + else { + locked_pid = uwsgi_lock_check(uli); + } + if (locked_pid == diedpid) { + uwsgi_log("[deadlock-detector] pid %d was holding lock %s (%p)\n", (int) diedpid, uli->id, uli->lock_ptr); if (uli->rw) { - locked_pid = uwsgi_rwlock_check(uli); + uwsgi_rwunlock(uli); } else { - locked_pid = uwsgi_lock_check(uli); + uwsgi_unlock(uli); } - if (locked_pid == diedpid) { - uwsgi_log("[deadlock-detector] pid %d was holding lock %s (%p)\n", (int) diedpid, uli->id, uli->lock_ptr); - if (uli->rw) { - uwsgi_rwunlock(uli); - } - else { - uwsgi_unlock(uli); - } - } -nextlock: - uli = uli->next; } +nextlock: + uli = uli->next; } + + // reload gateways and daemons only on normal workflow if (!uwsgi.to_heaven && !uwsgi.to_hell) { @@ -1440,7 +1413,7 @@ nextlock: /* reload the spooler */ struct uwsgi_spooler *uspool = uwsgi.spoolers; pid_found = 0; - while(uspool) { + while (uspool) { if (uspool->pid > 0 && diedpid == uspool->pid) { uwsgi_log("OOOPS the spooler is no more...trying respawn...\n"); uspool->respawned++; @@ -1500,7 +1473,6 @@ nextlock: } - if (diedpid <= 0) continue; /* What happens here ? @@ -1517,7 +1489,7 @@ nextlock: // check spooler, mules, gateways and daemons #ifdef UWSGI_SPOOLER struct uwsgi_spooler *uspool = uwsgi.spoolers; - while(uspool) { + while (uspool) { if (uspool->pid > 0 && diedpid == uspool->pid) { uwsgi_log("spooler (pid: %d) annihilated\n", (int) diedpid); goto next; @@ -1558,83 +1530,83 @@ nextlock: else if (WIFSTOPPED(waitpid_status)) { uwsgi_log("subprocess %d stopped\n", (int) diedpid); } - next: +next: continue; } - else { - if (uwsgi.to_heaven) { - ready_to_reload++; - uwsgi.workers[uwsgi.mywid].pid = 0; - // only to be safe :P - uwsgi.workers[uwsgi.mywid].harakiri = 0; - continue; - } - else if (uwsgi.to_hell) { - ready_to_die++; - uwsgi.workers[uwsgi.mywid].pid = 0; - // only to be safe :P - uwsgi.workers[uwsgi.mywid].harakiri = 0; - continue; - } - else if (uwsgi.to_outworld) { - uwsgi.lazy_respawned++; - uwsgi.workers[uwsgi.mywid].destroy = 0; - uwsgi.workers[uwsgi.mywid].pid = 0; - // only to be safe :P - uwsgi.workers[uwsgi.mywid].harakiri = 0; - } + // ok a worker died... + if (uwsgi.to_heaven) { + ready_to_reload++; + uwsgi.workers[uwsgi.mywid].pid = 0; + // only to be safe :P + uwsgi.workers[uwsgi.mywid].harakiri = 0; + continue; + } + else if (uwsgi.to_hell) { + ready_to_die++; + uwsgi.workers[uwsgi.mywid].pid = 0; + // only to be safe :P + uwsgi.workers[uwsgi.mywid].harakiri = 0; + continue; + } + else if (uwsgi.to_outworld) { + uwsgi.lazy_respawned++; + uwsgi.workers[uwsgi.mywid].destroy = 0; + uwsgi.workers[uwsgi.mywid].pid = 0; + // only to be safe :P + uwsgi.workers[uwsgi.mywid].harakiri = 0; + } - if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_FAILED_APP_CODE) { - uwsgi_log("OOPS ! failed loading app in worker %d (pid %d) :( trying again...\n", uwsgi.mywid, (int) diedpid); - } - else if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_DE_HIJACKED_CODE) { - uwsgi_log("...restoring worker %d (pid: %d)...\n", uwsgi.mywid, (int) diedpid); - } - else if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_EXCEPTION_CODE) { - uwsgi_log("... monitored exception detected, respawning worker %d (pid: %d)...\n", uwsgi.mywid, (int) diedpid); - } - else if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_QUIET_CODE) { - // noop - } - else if (uwsgi.workers[uwsgi.mywid].manage_next_request) { - if (WIFSIGNALED(waitpid_status)) { - uwsgi_log("DAMN ! worker %d (pid: %d) died, killed by signal %d :( trying respawn ...\n", uwsgi.mywid, (int) diedpid, (int) WTERMSIG(waitpid_status)); - } - else { - uwsgi_log("DAMN ! worker %d (pid: %d) died :( trying respawn ...\n", uwsgi.mywid, (int) diedpid); - } - } - - if (uwsgi.workers[uwsgi.mywid].cheaped == 1) { - uwsgi.workers[uwsgi.mywid].pid = 0; - uwsgi_log("uWSGI worker %d cheaped.\n", uwsgi.mywid); - uwsgi.workers[uwsgi.mywid].harakiri = 0; - continue; - } - gettimeofday(&last_respawn, NULL); - if (last_respawn.tv_sec <= uwsgi.respawn_delta + check_interval) { - last_respawn_rate++; - if (last_respawn_rate > uwsgi.numproc) { - if (uwsgi.forkbomb_delay > 0) { - uwsgi_log("worker respawning too fast !!! i have to sleep a bit (%d seconds)...\n", uwsgi.forkbomb_delay); - /* use --forkbomb-delay 0 to disable sleeping */ - sleep(uwsgi.forkbomb_delay); - } - last_respawn_rate = 0; - } + if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_FAILED_APP_CODE) { + uwsgi_log("OOPS ! failed loading app in worker %d (pid %d) :( trying again...\n", uwsgi.mywid, (int) diedpid); + } + else if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_DE_HIJACKED_CODE) { + uwsgi_log("...restoring worker %d (pid: %d)...\n", uwsgi.mywid, (int) diedpid); + } + else if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_EXCEPTION_CODE) { + uwsgi_log("... monitored exception detected, respawning worker %d (pid: %d)...\n", uwsgi.mywid, (int) diedpid); + } + else if (WIFEXITED(waitpid_status) && WEXITSTATUS(waitpid_status) == UWSGI_QUIET_CODE) { + // noop + } + else if (uwsgi.workers[uwsgi.mywid].manage_next_request) { + if (WIFSIGNALED(waitpid_status)) { + uwsgi_log("DAMN ! worker %d (pid: %d) died, killed by signal %d :( trying respawn ...\n", uwsgi.mywid, (int) diedpid, (int) WTERMSIG(waitpid_status)); } else { + uwsgi_log("DAMN ! worker %d (pid: %d) died :( trying respawn ...\n", uwsgi.mywid, (int) diedpid); + } + } + + if (uwsgi.workers[uwsgi.mywid].cheaped == 1) { + uwsgi.workers[uwsgi.mywid].pid = 0; + uwsgi_log("uWSGI worker %d cheaped.\n", uwsgi.mywid); + uwsgi.workers[uwsgi.mywid].harakiri = 0; + continue; + } + gettimeofday(&last_respawn, NULL); + if (last_respawn.tv_sec <= uwsgi.respawn_delta + check_interval) { + last_respawn_rate++; + if (last_respawn_rate > uwsgi.numproc) { + if (uwsgi.forkbomb_delay > 0) { + uwsgi_log("worker respawning too fast !!! i have to sleep a bit (%d seconds)...\n", uwsgi.forkbomb_delay); + /* use --forkbomb-delay 0 to disable sleeping */ + sleep(uwsgi.forkbomb_delay); + } last_respawn_rate = 0; } - gettimeofday(&last_respawn, NULL); - uwsgi.respawn_delta = last_respawn.tv_sec; - - if (uwsgi_respawn_worker(uwsgi.mywid)) - return 0; - } + else { + last_respawn_rate = 0; + } + gettimeofday(&last_respawn, NULL); + uwsgi.respawn_delta = last_respawn.tv_sec; + + if (uwsgi_respawn_worker(uwsgi.mywid)) + return 0; + + // end of the loop } // never here diff --git a/uwsgi.h b/uwsgi.h index 5c43761d..31b1c086 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -1,6 +1,6 @@ /* uWSGI */ -/* indent -i8 -br -brs -brf -l0 -npsl -nip -npcs -npsl -di1 */ +/* indent -i8 -br -brs -brf -l0 -npsl -nip -npcs -npsl -di1 -il0 */ #ifdef __cplusplus extern "C" { @@ -2853,6 +2853,8 @@ int uwsgi_stats_str(struct uwsgi_stats *, char *); char *uwsgi_substitute(char *, char *, char *); +void manage_cluster_message(char *, int); + #ifdef UWSGI_SSL #include "openssl/conf.h" #include "openssl/ssl.h"