Compare commits

...
26 Commits
Author SHA1 Message Date
roberto@quantal64 f698b558e1 uWSGI 1.3-rc3 2012-09-13 18:46:56 +02:00
roberto@quantal64 ca832900a0 ensure master_cleanup is run by the master 2012-09-13 16:58:05 +02:00
roberto@quantal64 6a6227a8e7 added a warning for --map-socket invalid syntax 2012-09-13 14:09:55 +02:00
roberto@quantal64 8539cdb5a7 ported build system to python 3.3 2012-09-13 14:05:47 +02:00
roberto@quantal64 5cea1de1ba added master_cleanup hook and carbon flush on stop/reload 2012-09-12 19:03:22 +02:00
roberto@quantal64 e735b9e31f applied busyness_backlog.diff by Łukasz Mierzwa 2012-09-12 18:12:55 +02:00
roberto@quantal64 9ff8f8899f added --log-drain 2012-09-12 08:08:32 +02:00
roberto@quantal64 d908c4ef35 allows building without udp 2012-09-07 19:45:58 +02:00
roberto@quantal64 207f50726d oops, use classic perl syntax for push 2012-09-07 19:36:21 +02:00
roberto@quantal64 212559fcce put sleep() in each psgi cleanup hook to show its power :P 2012-09-07 19:13:34 +02:00
roberto@quantal64 cf480e499c implemented psgi cleanup handlers, and fixed a leak in rpc call 2012-09-07 19:09:22 +02:00
roberto@quantal64 f2e27f5154 improved process limit detection 2012-09-07 14:54:57 +02:00
roberto@quantal64 d932a4feae fixed typcasting 2012-09-07 14:40:06 +02:00
roberto@quantal64 da6e05072f explicit link with librt for alternatives clock 2012-09-07 06:58:26 +02:00
roberto@quantal64 13277d9d71 improved cheaper_busyness 2012-09-06 17:10:03 +02:00
roberto@quantal64 f17f8b337f fixed idle mode 2012-09-06 16:21:12 +02:00
roberto@quantal64 ffb454e414 headers must be sent before close() 2012-09-06 15:22:17 +02:00
roberto@quantal64 424f5a09e7 call close() always, even if client disconnect 2012-09-06 15:02:12 +02:00
roberto@quantal64 c29a6f6446 removed STOP/TSTP signal from gateways 2012-09-06 11:20:21 +02:00
roberto@quantal64 6f7edf0941 fixed uwsgi_buffer realloc len 2012-09-04 19:18:38 +02:00
roberto@quantal64 a4c4b7e8d2 fixed uwsgi_buffer realloc 2012-09-04 18:13:16 +02:00
roberto@quantal64 837f6f3a54 initialize spooler locks before app loading 2012-09-04 08:21:40 +02:00
roberto@quantal64 13041da6bf fixed piping 2012-09-04 08:13:02 +02:00
roberto@quantal64 046824dae7 fixed usage of mmap() 2012-09-03 14:11:42 +02:00
roberto@quantal64 269a588d9d uWSGI 1.3-rc2 2012-09-01 11:42:42 +02:00
roberto@quantal64 50ae12518a Added tag 1.3-rc2 for changeset 14524da00a8b 2012-09-01 11:42:31 +02:00
28 changed files with 380 additions and 74 deletions
+1
View File
@@ -52,3 +52,4 @@ e1568fd16b7b586cc72deb4dfccbbe64ae0b84df 1.1
24c8fe5a26e557db41cf1cee45b8964d5885e484 1.2-rc1
29de0fb320bc0a1ce84972f9f249360e315d3c10 1.2-rc2
c3cdecbf2bac591336baddd14f9ac22e18d6e200 1.2
14524da00a8b382dffb1d16a0e70cdf4a946f369 1.3-rc2
+4 -1
View File
@@ -18,11 +18,14 @@ int uwsgi_buffer_append(struct uwsgi_buffer *ub, char *buf, size_t len) {
size_t remains = ub->len - ub->pos;
if (len > remains) {
char *new_buf = realloc(ub->buf, ub->len + UMAX(len, (size_t) uwsgi.page_size));
size_t chunk_size = UMAX(len, (size_t) uwsgi.page_size);
char *new_buf = realloc(ub->buf, ub->len + chunk_size);
if (!new_buf) {
uwsgi_error("realloc()");
return -1;
}
ub->buf = new_buf;
ub->len += chunk_size;
}
memcpy(ub->buf + ub->pos, buf, len);
+2
View File
@@ -71,6 +71,8 @@ void gateway_respawn(int id) {
signal(SIGUSR1, SIG_IGN);
signal(SIGUSR2, SIG_IGN);
signal(SIGPIPE, SIG_IGN);
signal(SIGSTOP, SIG_IGN);
signal(SIGTSTP, SIG_IGN);
ug->loop(id, ug->data);
// never here !!! (i hope)
+24 -5
View File
@@ -214,16 +214,26 @@ void uwsgi_setup_workers() {
// allocate memory for cores
uwsgi.workers[i].cores = (struct uwsgi_core *) uwsgi_calloc_shared(sizeof(struct uwsgi_core) * uwsgi.cores);
// this is a trick for avoiding too much memory areas
void *ts = uwsgi_calloc_shared(sizeof(void *) * uwsgi.max_apps * uwsgi.cores);
void *buffers = uwsgi_malloc_shared(uwsgi.buffer_size * uwsgi.cores);
void *hvec = uwsgi_malloc_shared(sizeof(struct iovec) * uwsgi.vec_size * uwsgi.cores);
void *post_buf = NULL;
if (uwsgi.post_buffering > 0)
post_buf = uwsgi_malloc_shared(uwsgi.post_buffering_bufsize * uwsgi.cores);
for (j = 0; j < uwsgi.cores; j++) {
// allocate shared memory for thread states (required for some language, like python)
uwsgi.workers[i].cores[j].ts = uwsgi_calloc_shared(sizeof(void *) * uwsgi.max_apps);
uwsgi.workers[i].cores[j].ts = ts + ((sizeof(void *) * uwsgi.max_apps) * j);
// raw per-request buffer
uwsgi.workers[i].cores[j].buffer = uwsgi_malloc_shared(uwsgi.buffer_size);
uwsgi.workers[i].cores[j].buffer = buffers + (uwsgi.buffer_size * j);
// iovec for uwsgi vars
uwsgi.workers[i].cores[j].hvec = uwsgi_malloc_shared(sizeof(struct iovec) * uwsgi.vec_size);
if (uwsgi.post_buffering > 0)
uwsgi.workers[i].cores[j].post_buf = uwsgi_malloc_shared(uwsgi.post_buffering_bufsize);
uwsgi.workers[i].cores[j].hvec = hvec + ((sizeof(struct iovec) * uwsgi.vec_size) * j);
if (post_buf)
uwsgi.workers[i].cores[j].post_buf = post_buf + (uwsgi.post_buffering_bufsize *j);
}
// master does not need to following steps...
if (i == 0) continue;
uwsgi.workers[i].signal_pipe[0] = -1;
@@ -232,4 +242,13 @@ void uwsgi_setup_workers() {
snprintf(uwsgi.workers[i].snapshot_name, 0xff, "uWSGI snapshot %d", i);
}
uint64_t total_memory = (sizeof(struct uwsgi_app) * uwsgi.max_apps) + (sizeof(struct uwsgi_core) * uwsgi.cores) + (sizeof(void *) * uwsgi.max_apps * uwsgi.cores) +
(uwsgi.buffer_size * uwsgi.cores) + (sizeof(struct iovec) * uwsgi.vec_size * uwsgi.cores);
if (uwsgi.post_buffering > 0) {
total_memory += (uwsgi.post_buffering_bufsize * uwsgi.cores);
}
total_memory *= (uwsgi.numproc + uwsgi.master_process);
uwsgi_log("mapped %llu bytes (%llu KB) for %d cores\n", total_memory, total_memory / 1024, uwsgi.cores*uwsgi.numproc);
}
+23 -2
View File
@@ -265,6 +265,15 @@ int uwsgi_master_log(void) {
ssize_t rlen = read(uwsgi.shared->worker_log_pipe[0], uwsgi.log_master_buf, uwsgi.log_master_bufsize);
if (rlen > 0) {
#ifdef UWSGI_PCRE
struct uwsgi_regexp_list *url = uwsgi.log_drain_rules;
while(url) {
if (uwsgi_regexp_match(url->pattern, url->pattern_extra, uwsgi.log_master_buf, rlen) >= 0) {
return 0;
}
url = url->next;
}
#endif
if (uwsgi.choosen_logger) {
struct uwsgi_logger *ul = uwsgi.choosen_logger;
while(ul) {
@@ -275,7 +284,6 @@ int uwsgi_master_log(void) {
else {
rlen = write(uwsgi.original_log_fd, uwsgi.log_master_buf, rlen);
}
// TODO allow uwsgi.logger = func
return 0;
}
@@ -600,6 +608,8 @@ int master_loop(char **argv, char **environ) {
uwsgi_unix_signal(SIGURG, uwsgi_restore_auto_snapshot);
}
atexit(uwsgi_master_cleanup_hooks);
uwsgi.master_queue = event_queue_init();
/* route signals to workers... */
@@ -1179,12 +1189,23 @@ health_cycle:
uwsgi.current_time = uwsgi_now();
if (!last_request_timecheck)
last_request_timecheck = uwsgi.current_time;
int busy_workers = 0;
for (i = 1; i <= uwsgi.numproc; i++) {
if (uwsgi.workers[i].cheaped == 0 && uwsgi.workers[i].pid > 0) {
if (uwsgi.workers[i].busy == 1) {
busy_workers = 1;
break;
}
}
}
if (last_request_count != uwsgi.workers[0].requests) {
last_request_timecheck = uwsgi.current_time;
last_request_count = uwsgi.workers[0].requests;
}
// a bit of over-engeneering to avoid clock skews
else if (last_request_timecheck < uwsgi.current_time && (uwsgi.current_time - last_request_timecheck > uwsgi.idle)) {
else if (last_request_timecheck < uwsgi.current_time && (uwsgi.current_time - last_request_timecheck > uwsgi.idle) && !busy_workers) {
uwsgi_log("workers have been inactive for more than %d seconds (%llu-%llu)\n", uwsgi.idle, (unsigned long long) uwsgi.current_time, (unsigned long long) last_request_timecheck);
uwsgi.cheap = 1;
if (uwsgi.die_on_idle) {
+27
View File
@@ -5,6 +5,30 @@ extern struct uwsgi_server uwsgi;
void worker_wakeup() {
}
void uwsgi_master_cleanup_hooks(void) {
int j;
// could be an inherited atexit hook
if (uwsgi.mywid > 0) return ;
uwsgi.cleaning = 1;
for (j = 0; j < uwsgi.gp_cnt; j++) {
if (uwsgi.gp[j]->master_cleanup) {
uwsgi.gp[j]->master_cleanup();
}
}
for (j = 0; j < 256; j++) {
if (uwsgi.p[j]->master_cleanup) {
uwsgi.p[j]->master_cleanup();
}
}
}
int uwsgi_calc_cheaper(void) {
int i;
@@ -240,6 +264,9 @@ void uwsgi_reload(char **argv) {
waitpid(WAIT_ANY, &waitpid_status, WNOHANG);
}
// call master cleanup hooks
uwsgi_master_cleanup_hooks();
// call atexit user exec
uwsgi_exec_atexit();
+1 -1
View File
@@ -475,7 +475,7 @@ ssize_t uwsgi_send_message(int fd, uint8_t modifier1, uint8_t modifier2, char *m
// transfer data from one socket to another
if (pfd >= 0 && plen > 0) {
ret = uwsgi_pipe_sized(pfd, fd, timeout, plen);
ret = uwsgi_pipe_sized(pfd, fd, plen, timeout);
if (ret < 0) return -1;
}
+4
View File
@@ -1526,6 +1526,10 @@ void uwsgi_map_sockets() {
while (usl) {
char *colon = strchr(usl->value, ':');
if (!colon) {
uwsgi_log("invalid socket mapping, must be socket:worker[,worker...]\n");
exit(1);
}
if ((int) uwsgi_str_num(usl->value, colon - usl->value) == uwsgi_get_socket_num(uwsgi_sock)) {
enabled = 0;
char *p = strtok(colon + 1, ",");
+34 -2
View File
@@ -1112,6 +1112,7 @@ void sanitize_args() {
uwsgi.ignore_write_errors = 1;
}
if (uwsgi.cheaper_count > 0 && uwsgi.cheaper_count >= uwsgi.numproc) {
uwsgi_log("invalid cheaper value: must be lower than processes\n");
exit(1);
@@ -3133,7 +3134,8 @@ void *uwsgi_malloc_shared(size_t size) {
void *addr = mmap(NULL, size, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANON, -1, 0);
if (addr == NULL) {
if (addr == MAP_FAILED) {
uwsgi_log("unable to allocate %llu bytes (%lluMB)\n", (unsigned long long )size, (unsigned long long) (size/(1024*1024)));
uwsgi_error("mmap()");
exit(1);
}
@@ -3177,6 +3179,36 @@ struct uwsgi_string_list *uwsgi_string_new_list(struct uwsgi_string_list **list,
return uwsgi_string;
}
#ifdef UWSGI_PCRE
struct uwsgi_regexp_list *uwsgi_regexp_new_list(struct uwsgi_regexp_list **list, char *value) {
struct uwsgi_regexp_list *url = *list, *old_url;
if (!url) {
*list = uwsgi_malloc(sizeof(struct uwsgi_regexp_list));
url = *list;
}
else {
while (url) {
old_url = url;
url = url->next;
}
url = uwsgi_malloc(sizeof(struct uwsgi_regexp_list));
old_url->next = url;
}
if (uwsgi_regexp_build(value, &url->pattern, &url->pattern_extra)) {
exit(1);
}
url->next = NULL;
url->custom = 0;
return url;
}
#endif
char *uwsgi_string_get_list(struct uwsgi_string_list **list, int pos, size_t * len) {
struct uwsgi_string_list *uwsgi_string = *list;
@@ -4602,7 +4634,7 @@ timeout:
ssize_t uwsgi_pipe_sized(int src, int dst, size_t required, int timeout) {
char buf[8192];
size_t written = -1;
size_t written = 0;
ssize_t len;
while(written < required) {
+39 -11
View File
@@ -349,6 +349,9 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"logger-list", no_argument, 0, "list enabled loggers", uwsgi_opt_true, &uwsgi.loggers_list, 0},
{"loggers-list", no_argument, 0, "list enabled loggers", uwsgi_opt_true, &uwsgi.loggers_list, 0},
{"threaded-logger", no_argument, 0, "offload log writing to a thread", uwsgi_opt_true, &uwsgi.threaded_logger, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
#ifdef UWSGI_PCRE
{"log-drain", required_argument, 0, "drain (do not show) log lines matching the specified regexp", uwsgi_opt_add_regexp_list, &uwsgi.log_drain_rules, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
#endif
#ifdef UWSGI_ZEROMQ
{"log-zeromq", required_argument, 0, "send logs to a zeromq server", uwsgi_opt_set_logger, "zeromq", UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
#endif
@@ -1938,6 +1941,16 @@ int uwsgi_start(void *v_argv) {
uwsgi_error("setrlimit()");
}
}
if (!getrlimit(RLIMIT_NPROC, &uwsgi.rl_nproc)) {
if (uwsgi.rl_nproc.rlim_cur != RLIM_INFINITY) {
uwsgi_log("your processes number limit is %d\n", (int) uwsgi.rl_nproc.rlim_cur);
if ((int)uwsgi.rl_nproc.rlim_cur < uwsgi.numproc+uwsgi.master_process) {
uwsgi.numproc = uwsgi.rl_nproc.rlim_cur - 1;
uwsgi_log("!!! number of workers adjusted to %d due to system limits !!!\n", uwsgi.numproc);
}
}
}
#endif
#ifndef __OpenBSD__
@@ -1970,7 +1983,6 @@ int uwsgi_start(void *v_argv) {
}
#endif
uwsgi_log_initial("your memory page size is %d bytes\n", uwsgi.page_size);
if (uwsgi.buffer_size > 65536) {
@@ -2075,10 +2087,6 @@ int uwsgi_start(void *v_argv) {
uwsgi_log("cores allocated...\n");
#endif
if (uwsgi.cores > 1) {
uwsgi_log("allocated %llu bytes (%llu KB) for %d cores per worker.\n", (uint64_t) (sizeof(struct uwsgi_core) * uwsgi.cores), (uint64_t) ((sizeof(struct uwsgi_core) * uwsgi.cores) / 1024), uwsgi.cores);
}
if (uwsgi.vhost) {
uwsgi_log("VirtualHosting mode enabled.\n");
}
@@ -2328,6 +2336,23 @@ int uwsgi_start(void *v_argv) {
}
}
#ifdef UWSGI_SPOOLER
// initialize locks and socket as soon as possibile, as the master could enqueue tasks
if (uwsgi.spoolers != NULL && uwsgi.sockets) {
create_signal_pipe(uwsgi.shared->spooler_signal_pipe);
struct uwsgi_spooler *uspool = uwsgi.spoolers;
while (uspool) {
// lock is required even in EXTERNAL mode
uspool->lock = uwsgi_lock_init(uwsgi_concat2("spooler on ", uspool->dir));
if (uspool->mode == UWSGI_SPOOLER_EXTERNAL) goto next;
create_signal_pipe(uspool->signal_pipe);
next:
uspool = uspool->next;
}
}
#endif
// preinit apps (create the language environment)
for (i = 0; i < 256; i++) {
if (uwsgi.p[i]->preinit_apps) {
@@ -2460,15 +2485,11 @@ int uwsgi_start(void *v_argv) {
#ifdef UWSGI_SPOOLER
if (uwsgi.spoolers != NULL && uwsgi.sockets) {
create_signal_pipe(uwsgi.shared->spooler_signal_pipe);
struct uwsgi_spooler *uspool = uwsgi.spoolers;
while (uspool) {
// lock is required even in EXTERNAL mode
uspool->lock = uwsgi_lock_init(uwsgi_concat2("spooler on ", uspool->dir));
if (uspool->mode == UWSGI_SPOOLER_EXTERNAL) goto next;
create_signal_pipe(uspool->signal_pipe);
if (uspool->mode == UWSGI_SPOOLER_EXTERNAL) goto next2;
uspool->pid = spooler_start(uspool);
next:
next2:
uspool = uspool->next;
}
}
@@ -3155,6 +3176,13 @@ void uwsgi_opt_add_string_list(char *opt, char *value, void *list) {
uwsgi_string_new_list(ptr, value);
}
#ifdef UWSGI_PCRE
void uwsgi_opt_add_regexp_list(char *opt, char *value, void *list) {
struct uwsgi_regexp_list **ptr = (struct uwsgi_regexp_list **) list;
uwsgi_regexp_new_list(ptr, value);
}
#endif
void uwsgi_opt_add_shared_socket(char *opt, char *value, void *protocol) {
uwsgi_new_shared_socket(generate_socket_name(value));
}
+4 -1
View File
@@ -68,7 +68,7 @@ void carbon_master_cycle() {
if (last_update == 0) last_update = uwsgi_now();
// update
if (uwsgi.current_time - last_update >= u_carbon.freq) {
if (uwsgi.current_time - last_update >= u_carbon.freq || uwsgi.cleaning) {
for (i = 0; i < uwsgi.numproc; i++) {
u_carbon.current_busyness_values[i] = uwsgi.workers[i+1].running_time - u_carbon.last_busyness_values[i];
@@ -203,10 +203,13 @@ nxt:
}
}
struct uwsgi_plugin carbon_plugin = {
.name = "carbon",
.master_cleanup = carbon_master_cycle,
.options = carbon_options,
.master_cycle = carbon_master_cycle,
.post_init = carbon_post_init,
+69 -18
View File
@@ -22,8 +22,11 @@ struct uwsgi_cheaper_busyness_global {
int last_action; // 1 - spawn workers ; 2 - cheap worker
int verbose; // 1 - show debug logs, 0 - only important
uint64_t tolerance_counter; // used to keep track of what to do if min <= busyness <= max for few cycles in row
int emergency_workers; // counts the number of running emergency workers
#ifdef __linux__
int backlog_alert;
int backlog_step;
uint64_t backlog_multi; // multiplier used to cheap emergency workers
#endif
} uwsgi_cheaper_busyness_global;
@@ -52,6 +55,12 @@ struct uwsgi_option uwsgi_cheaper_busyness_options[] = {
{"cheaper-busyness-backlog-alert", required_argument, 0,
"spawn emergency worker if anytime listen queue is higher than this value (default 33)",
uwsgi_opt_set_int, &uwsgi_cheaper_busyness_global.backlog_alert, 0},
{"cheaper-busyness-backlog-multiplier", required_argument, 0,
"set cheaper multiplier used for emergency workers (default 3)",
uwsgi_opt_set_64bit, &uwsgi_cheaper_busyness_global.backlog_multi, 0},
{"cheaper-busyness-backlog-step", required_argument, 0,
"number of emergency workers to spawn at a time (default 1)",
uwsgi_opt_set_int, &uwsgi_cheaper_busyness_global.backlog_step, 0},
#endif
{0, 0, 0, 0, 0, 0 ,0},
@@ -63,10 +72,22 @@ struct uwsgi_option uwsgi_cheaper_busyness_options[] = {
void set_next_cheap_time(void) {
uint64_t now = uwsgi_micros();
// we will start workers now so we will set time when workers can be cheaped to
// some time in the future, so that workers are cheaped only if long term busyness
// is low enough
uwsgi_cheaper_busyness_global.next_cheap = now + uwsgi.cheaper_overload*uwsgi_cheaper_busyness_global.cheap_multi*1000000;
#ifdef __linux__
if (uwsgi_cheaper_busyness_global.emergency_workers > 0) {
// we have some emergency workers running, we will use minimum delay (2 cycles) to cheap workers
// to have quicker recovery from big but short load spikes
// otherwise we might wait a lot before cheaping all emergency workers
if (uwsgi_cheaper_busyness_global.verbose)
uwsgi_log("[busyness] %d emergency worker(s) running, using %d seconds cheaper timer\n",
uwsgi_cheaper_busyness_global.emergency_workers, uwsgi.cheaper_overload*uwsgi_cheaper_busyness_global.backlog_multi);
uwsgi_cheaper_busyness_global.next_cheap = now + uwsgi.cheaper_overload*uwsgi_cheaper_busyness_global.backlog_multi*1000000;
} else {
#endif
// no emergency workers running, we use normal math for setting timer
uwsgi_cheaper_busyness_global.next_cheap = now + uwsgi.cheaper_overload*uwsgi_cheaper_busyness_global.cheap_multi*1000000;
#ifdef __linux__
}
#endif
}
@@ -79,6 +100,36 @@ void decrease_multi(void) {
}
#ifdef __linux__
int spawn_emergency_worker(int backlog) {
// reset cheaper multiplier to minimum value so we can start cheaping workers sooner
// if this was just random spike
uwsgi_cheaper_busyness_global.cheap_multi = uwsgi_cheaper_busyness_global.min_multi;
// set last action to spawn
uwsgi_cheaper_busyness_global.last_action = 1;
int decheaped = 0;
int i;
for (i = 1; i <= uwsgi.numproc; i++) {
if (uwsgi.workers[i].cheaped == 1 && uwsgi.workers[i].pid == 0) {
decheaped++;
if (decheaped >= uwsgi_cheaper_busyness_global.backlog_step) break;
}
}
uwsgi_cheaper_busyness_global.emergency_workers += decheaped;
set_next_cheap_time();
uwsgi_log("[busyness] %d requests in listen queue, spawning %d emergency worker(s) (%d)!\n",
backlog, decheaped, uwsgi_cheaper_busyness_global.emergency_workers);
return decheaped;
}
#endif
int cheaper_busyness_algo(void) {
int i;
@@ -98,6 +149,8 @@ int cheaper_busyness_algo(void) {
#ifdef __linux__
if (!uwsgi_cheaper_busyness_global.backlog_alert) uwsgi_cheaper_busyness_global.backlog_alert = 33;
if (!uwsgi_cheaper_busyness_global.backlog_multi) uwsgi_cheaper_busyness_global.backlog_multi = 3;
if (!uwsgi_cheaper_busyness_global.backlog_step) uwsgi_cheaper_busyness_global.backlog_step = 1;
#endif
if (!uwsgi_cheaper_busyness_global.min_multi) {
@@ -108,7 +161,8 @@ int cheaper_busyness_algo(void) {
uwsgi_cheaper_busyness_global.busyness_min, uwsgi_cheaper_busyness_global.busyness_max,
uwsgi.cheaper_overload, uwsgi_cheaper_busyness_global.cheap_multi, uwsgi_cheaper_busyness_global.penalty);
#ifdef __linux__
uwsgi_log("[busyness] backlog alert is set to %d request(s)\n", uwsgi_cheaper_busyness_global.backlog_alert);
uwsgi_log("[busyness] backlog alert is set to %d request(s), step is %d\n",
uwsgi_cheaper_busyness_global.backlog_alert, uwsgi_cheaper_busyness_global.backlog_step);
#endif
}
@@ -195,12 +249,7 @@ int cheaper_busyness_algo(void) {
#ifdef __linux__
} else if (backlog > uwsgi_cheaper_busyness_global.backlog_alert && active_workers < uwsgi.numproc) {
// reset counters
set_next_cheap_time();
uwsgi_cheaper_busyness_global.last_action = 1;
uwsgi_log("[busyness] %d requests in listen queue, spawning emergency worker!\n", backlog);
return 1;
return spawn_emergency_worker(backlog);
#endif
} else if (avg_busyness < uwsgi_cheaper_busyness_global.busyness_min) {
@@ -226,6 +275,9 @@ int cheaper_busyness_algo(void) {
// store information that last action performed was cheaping worker
uwsgi_cheaper_busyness_global.last_action = 2;
if (uwsgi_cheaper_busyness_global.emergency_workers > 0)
uwsgi_cheaper_busyness_global.emergency_workers--;
return -1;
} else if (uwsgi_cheaper_busyness_global.verbose)
uwsgi_log("[busyness] need to wait %d more second(s) to cheap worker\n", (uwsgi_cheaper_busyness_global.next_cheap - now)/1000000);
@@ -235,6 +287,11 @@ int cheaper_busyness_algo(void) {
// with only 1 worker running there is no point in doing all that magic
if (active_workers == 1) return 0;
if (uwsgi_cheaper_busyness_global.emergency_workers > 0)
// we had emergency workers running and we went down to the busyness
// level that is high enough to slow down cheaping workers at extra speed
uwsgi_cheaper_busyness_global.emergency_workers--;
// we have min <= busyness <= max we need to check what happened before
uwsgi_cheaper_busyness_global.tolerance_counter++;
@@ -260,13 +317,7 @@ int cheaper_busyness_algo(void) {
#ifdef __linux__
} else if (backlog > uwsgi_cheaper_busyness_global.backlog_alert && active_workers < uwsgi.numproc) {
// we check for backlog overload every cycle
// reset counters
set_next_cheap_time();
uwsgi_cheaper_busyness_global.last_action = 1;
uwsgi_log("[busyness] %d requests in listen queue, spawning emergency worker!\n", backlog);
return 1;
return spawn_emergency_worker(backlog);
#endif
}
+1 -1
View File
@@ -2,5 +2,5 @@ NAME='clock_monotonic'
CFLAGS = []
LDFLAGS = []
LIBS = []
LIBS = ['-lrt']
GCC_LIST = ['clock_monotonic']
+1 -1
View File
@@ -2,5 +2,5 @@ NAME='clock_realtime'
CFLAGS = []
LDFLAGS = []
LIBS = []
LIBS = ['-lrt']
GCC_LIST = ['clock_realtime']
+5
View File
@@ -250,7 +250,12 @@ void uwsgi_corerouter_setup_sockets(struct uwsgi_corerouter *ucr) {
else if (ugs->subscription) {
if (ugs->fd == -1) {
if (strchr(ugs->name, ':')) {
#ifdef UWSGI_UDP
ugs->fd = bind_to_udp(ugs->name, 0, 0);
#else
uwsgi_log("uWSGI has been built without UDP support !!!\n");
exit(1);
#endif
}
else {
ugs->fd = bind_to_unix_dgram(ugs->name);
+4
View File
@@ -80,10 +80,14 @@ int uwsgi_cr_map_use_to(struct uwsgi_corerouter *ucr, struct corerouter_session
}
int uwsgi_cr_map_use_cluster(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
#ifdef UWSGI_MULTICAST
cr_session->instance_address = uwsgi_cluster_best_node();
if (cr_session->instance_address) {
cr_session->instance_address_len = strlen(cr_session->instance_address);
}
#else
uwsgi_log("uWSGI has been built without multicast/clustering support !!!\n");
#endif
return 0;
}
+4
View File
@@ -822,12 +822,16 @@ To have a reliable implementation, we need to reset a bunch of values
// writable ?
if (cs->fd_state) {
#ifdef UWSGI_SSL
if (!cs->ugs->mode == UWSGI_HTTP_SSL) {
len = cs->send(&uhttp.cr, cs, NULL, 0);
}
else {
#endif
len = cs->send(&uhttp.cr, cs, hs->buffer,hs->buffer_len);
#ifdef UWSGI_SSL
}
#endif
#ifdef UWSGI_EVENT_USE_PORT
event_queue_add_fd_write(ucr->queue, cs->fd);
#endif
+2
View File
@@ -498,6 +498,7 @@ PHP_FUNCTION(uwsgi_rpc) {
argvs[i] = Z_STRLEN_P(z_current_obj);
}
// response must always be freed
char *response = uwsgi_do_rpc(node, func, num_args - 2, argv, argvs, &size);
if (size > 0) {
@@ -506,6 +507,7 @@ PHP_FUNCTION(uwsgi_rpc) {
free(response);
RETURN_STRING(ret, 0);
}
free(response);
clear:
efree(varargs);
+49 -9
View File
@@ -297,6 +297,12 @@ SV *build_psgi_env(struct wsgi_request *wsgi_req) {
if (!hv_store(env, "psgix.harakiri", 14, newSViv(1), 0)) goto clear;
}
if (!hv_store(env, "psgix.cleanup", 13, newSViv(1), 0)) goto clear;
// cleanup handlers array
av = newAV();
if (!hv_store(env, "psgix.cleanup.handlers", 22, newRV_noinc((SV *)av ), 0)) goto clear;
SV *pe = uwsgi_perl_obj_new("uwsgi::error", 12);
if (!hv_store(env, "psgi.errors", 11, pe, 0)) goto clear;
@@ -375,8 +381,6 @@ int uwsgi_perl_init(){
int uwsgi_perl_request(struct wsgi_request *wsgi_req) {
SV **harakiri;
#ifdef UWSGI_ASYNC
if (wsgi_req->async_status == UWSGI_AGAIN) {
return psgi_response(wsgi_req, wsgi_req->async_placeholder);
@@ -457,13 +461,7 @@ int uwsgi_perl_request(struct wsgi_request *wsgi_req) {
}
clear2:
// check for psgix.harakiri
harakiri = hv_fetch((HV*)SvRV( (SV*)wsgi_req->async_environ), "psgix.harakiri.commit", 21, 0);
if (harakiri) {
if (SvTRUE(*harakiri)) wsgi_req->async_plagued = 1;
}
SvREFCNT_dec(wsgi_req->async_environ);
// clear response
SvREFCNT_dec(wsgi_req->async_result);
clear:
@@ -478,16 +476,58 @@ clear:
return UWSGI_OK;
}
static void psgi_call_cleanup_hook(SV *hook, SV *env) {
dSP;
ENTER;
SAVETMPS;
PUSHMARK(SP);
XPUSHs(env);
PUTBACK;
call_sv(hook, G_DISCARD);
if(SvTRUE(ERRSV)) {
uwsgi_log("[uwsgi-perl error] %s\n", SvPV_nolen(ERRSV));
}
FREETMPS;
LEAVE;
}
void uwsgi_perl_after_request(struct wsgi_request *wsgi_req) {
log_request(wsgi_req);
// dereference %env
SV *env = SvRV((SV *) wsgi_req->async_environ);
// check for cleanup handlers
if (hv_exists((HV *)env, "psgix.cleanup.handlers", 22)) {
SV **cleanup_handlers = hv_fetch((HV *)env, "psgix.cleanup.handlers", 22, 0);
if (SvROK(*cleanup_handlers)) {
if (SvTYPE(SvRV(*cleanup_handlers)) == SVt_PVAV) {
I32 n = av_len((AV *)SvRV(*cleanup_handlers));
I32 i;
for(i=0;i<=n;i++) {
SV **hook = av_fetch((AV *)SvRV(*cleanup_handlers), i, 0);
psgi_call_cleanup_hook(*hook, (SV *) wsgi_req->async_environ);
}
}
}
}
// check for psgix.harakiri
if (hv_exists((HV *)env, "psgix.harakiri.commit", 21)) {
SV **harakiri = hv_fetch((HV *)env, "psgix.harakiri.commit", 21, 0);
if (SvTRUE(*harakiri)) wsgi_req->async_plagued = 1;
}
// async plagued could be defined in other areas...
if (wsgi_req->async_plagued) {
uwsgi_log("*** psgix.harakiri.commit requested ***\n");
goodbye_cruel_world();
}
// clear the env
SvREFCNT_dec(wsgi_req->async_environ);
}
int uwsgi_perl_magic(char *mountpoint, char *lazy) {
+2
View File
@@ -190,6 +190,7 @@ XS(XS_call) {
argvs[i] = arg_len;
}
// response must be always freed
char *response = uwsgi_do_rpc(NULL, func, items-1, argv, argvs, &size);
if (size > 0) {
@@ -198,6 +199,7 @@ XS(XS_call) {
free(response);
XSRETURN(1);
}
free(response);
XSRETURN_UNDEF;
}
+2
View File
@@ -313,6 +313,7 @@ PyObject *py_uwsgi_call(PyObject * self, PyObject * args) {
}
UWSGI_RELEASE_GIL;
// response must always be freed
char *response = uwsgi_do_rpc(NULL, func, argc - 1, argv, argvs, &size);
UWSGI_GET_GIL;
@@ -322,6 +323,7 @@ PyObject *py_uwsgi_call(PyObject * self, PyObject * args) {
return ret;
}
free(response);
Py_INCREF(Py_None);
return Py_None;
+18 -15
View File
@@ -234,20 +234,6 @@ exception:
exit(UWSGI_EXCEPTION_CODE);
}
}
if (PyObject_HasAttrString((PyObject *)wsgi_req->async_result, "close")) {
PyObject *close_method = PyObject_GetAttrString((PyObject *)wsgi_req->async_result, "close");
PyObject *close_method_args = PyTuple_New(0);
#ifdef UWSGI_DEBUG
uwsgi_log("calling close() for %.*s %p %p\n", wsgi_req->uri_len, wsgi_req->uri, close_method, close_method_args);
#endif
PyObject *close_method_output = PyEval_CallObject(close_method, close_method_args);
if (PyErr_Occurred()) {
PyErr_Print();
}
Py_DECREF(close_method_args);
Py_XDECREF(close_method_output);
Py_DECREF(close_method);
}
goto clear;
}
@@ -293,13 +279,30 @@ clear:
if (wsgi_req->sendfile_fd != -1) {
Py_DECREF((PyObject *)wsgi_req->async_sendfile);
}
Py_XDECREF((PyObject *)wsgi_req->async_placeholder);
// send the headers if not already sent
if (!wsgi_req->headers_sent && wsgi_req->headers_hvec > 0) {
uwsgi_python_do_send_headers(wsgi_req);
}
if (wsgi_req->async_placeholder) {
// CALL close() ALWAYS if we are working with an iterator !!!
if (PyObject_HasAttrString((PyObject *)wsgi_req->async_result, "close")) {
PyObject *close_method = PyObject_GetAttrString((PyObject *)wsgi_req->async_result, "close");
PyObject *close_method_args = PyTuple_New(0);
#ifdef UWSGI_DEBUG
uwsgi_log("calling close() for %.*s %p %p\n", wsgi_req->uri_len, wsgi_req->uri, close_method, close_method_args);
#endif
PyObject *close_method_output = PyEval_CallObject(close_method, close_method_args);
if (PyErr_Occurred()) {
PyErr_Print();
}
Py_DECREF(close_method_args);
Py_XDECREF(close_method_output);
Py_DECREF(close_method);
}
Py_DECREF((PyObject *)wsgi_req->async_placeholder);
}
Py_DECREF((PyObject *)wsgi_req->async_result);
PyErr_Clear();
+2 -1
View File
@@ -654,6 +654,7 @@ VALUE uwsgi_ruby_do_rpc(int argc, VALUE *rpc_argv, VALUE *class) {
argvs[i] = RSTRING_LEN(rpc_str);
}
// response must always be freed
char *response = uwsgi_do_rpc(node, func, argc - 2, argv, argvs, &size);
if (size > 0) {
@@ -661,7 +662,7 @@ VALUE uwsgi_ruby_do_rpc(int argc, VALUE *rpc_argv, VALUE *class) {
free(response);
return ret;
}
free(response);
clear:
+1 -1
View File
@@ -41,7 +41,7 @@ int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_route *u
if (wsgi_req->post_cl > 0) {
int post_fd = wsgi_req->poll.fd;
if (wsgi_req->async_post) {
post_fd = fileno(wsgi_req->async_post);
post_fd = fileno((FILE *)wsgi_req->async_post);
}
ret = uwsgi_pipe_sized(post_fd, http_fd, wsgi_req->post_cl, 0);
if (ret < 0) {
+1 -1
View File
@@ -52,7 +52,7 @@ int uwsgi_routing_func_uwsgi_remote(struct wsgi_request *wsgi_req, struct uwsgi_
int post_fd = wsgi_req->poll.fd;
if (wsgi_req->async_post) {
post_fd = fileno(wsgi_req->async_post);
post_fd = fileno((FILE*)wsgi_req->async_post);
}
if (uwsgi_send_message(uwsgi_fd, uh->modifier1, uh->modifier2, wsgi_req->buffer, wsgi_req->uh.pktsize, post_fd, wsgi_req->post_cl, 0) < 0) {
+24
View File
@@ -7,8 +7,32 @@ if ($rpc_value) {
print "rpc value = ".$rpc_value."\n";
}
my $one = sub {
my $env = shift;
sleep(1);
print "one\n";
};
my $two = sub {
my $env = shift;
sleep(1);
print "two\n";
};
my $three = sub {
my $env = shift;
sleep(1);
print "three\n";
};
my $app = sub {
my $env = shift;
if ($env->{'psgix.cleanup'}) {
print "cleanup supported\n";
push @{$env->{'psgix.cleanup.handlers'}}, $one;
push @{$env->{'psgix.cleanup.handlers'}}, $two;
push @{$env->{'psgix.cleanup.handlers'}}, $three;
}
uwsgi::cache_set("key1", "val1");
if ($rpc_value) {
print uwsgi::call('hello')."\n";
+24
View File
@@ -368,6 +368,17 @@ struct uwsgi_dyn_dict {
struct uwsgi_dyn_dict *next;
};
#ifdef UWSGI_PCRE
struct uwsgi_regexp_list {
pcre *pattern;
pcre_extra *pattern_extra;
uint64_t custom;
struct uwsgi_regexp_list *next;
};
#endif
union uwsgi_sockaddr {
struct sockaddr sa;
@@ -702,6 +713,8 @@ struct uwsgi_plugin {
int (*mule)(char *);
int (*mule_msg)(char *, size_t);
void (*master_cleanup) (void);
};
#ifdef UWSGI_PCRE
@@ -1386,6 +1399,10 @@ struct uwsgi_server {
struct uwsgi_logger *choosen_logger;
struct uwsgi_string_list *requested_logger;
#ifdef UWSGI_PCRE
struct uwsgi_regexp_list *log_drain_rules;
#endif
int threaded_logger;
pthread_mutex_t threaded_logger_lock;
@@ -1504,6 +1521,8 @@ struct uwsgi_server {
int to_hell;
int to_outworld;
int cleaning;
int ready_to_die;
int ready_to_reload;
@@ -2657,6 +2676,9 @@ struct uwsgi_socket *uwsgi_del_socket(struct uwsgi_socket *);
void uwsgi_close_all_sockets(void);
struct uwsgi_string_list *uwsgi_string_new_list(struct uwsgi_string_list **, char *);
#ifdef UWSGI_PCRE
struct uwsgi_regexp_list *uwsgi_regexp_new_list(struct uwsgi_regexp_list **, char *);
#endif
void uwsgi_string_del_list(struct uwsgi_string_list **, struct uwsgi_string_list *);
@@ -2886,6 +2908,7 @@ void uwsgi_opt_add_string_list(char *, char *, void *);
void uwsgi_opt_add_dyn_dict(char *, char *, void *);
#ifdef UWSGI_PCRE
void uwsgi_opt_add_regexp_dyn_dict(char *, char *, void *);
void uwsgi_opt_add_regexp_list(char *, char *, void *);
#endif
void uwsgi_opt_set_int(char *, char *, void *);
void uwsgi_opt_set_rawint(char *, char *, void *);
@@ -3222,6 +3245,7 @@ ssize_t uwsgi_pipe(int, int, int);
ssize_t uwsgi_pipe_sized(int, int, size_t, int);
int uwsgi_buffer_send(struct uwsgi_buffer *, int);
void uwsgi_master_cleanup_hooks(void);
void uwsgi_check_emperor(void);
#ifdef UWSGI_AS_SHARED_LIBRARY
+8 -4
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '1.3-dev'
uwsgi_version = '1.3-rc3'
import os
import re
@@ -235,9 +235,13 @@ def build_uwsgi(uc, print_only=False):
print("Error: plugin '%s' not found" % p)
sys.exit(1)
sys.path.insert(0, path)
import uwsgiplugin as up
reload(up)
try:
import importlib
up = importlib.machinery.SourceFileLoader('uwsgiplugin', '%s/uwsgiplugin.py' % path).load_module()
except:
sys.path.insert(0, path)
import uwsgiplugin as up
reload(up)
p_cflags = cflags[:]
p_cflags += up.CFLAGS