Compare commits

..
53 Commits
Author SHA1 Message Date
roberto@quantal64 c5dc74b021 uWSGI 1.3-rc4 2012-09-22 09:45:08 +02:00
roberto@quantal64 97137488f9 improved emperor mongodb tyrant mode 2012-09-22 09:38:15 +02:00
roberto@quantal64 8cff01a1c7 report imperial monitor in emperor statistics 2012-09-22 09:21:37 +02:00
roberto@quantal64 3c8452313e added emperor_mongodb plugin 2012-09-22 09:04:47 +02:00
roberto@quantal64 322b130fc9 refactored imperial monitors api 2012-09-22 08:58:17 +02:00
roberto@quantal64 6c001c2fda fixed usage of inlining 2012-09-22 07:57:42 +02:00
roberto@quantal64 5d4050817d added --dlopen 2012-09-22 05:52:28 +02:00
roberto@quantal64 0a1147a4f6 added --rb-patch-rack-bodyproxy for older ruby/rack versions 2012-09-21 11:02:29 +02:00
roberto@quantal64 f92cb0e6a2 fixed rack compilation 2012-09-21 09:51:36 +02:00
roberto@quantal64 8b79c76261 refactored daemonize2 2012-09-20 18:08:49 +02:00
roberto@quantal64 0015b51262 fixed POST-handling reports 2012-09-20 17:49:30 +02:00
roberto@quantal64 b1dc8cd9c6 improved async cores detection 2012-09-20 17:41:13 +02:00
roberto@quantal64 3269e2d288 added support for dependancies in plugins 2012-09-20 17:16:30 +02:00
roberto@quantal64 71a9747d96 applied latest carbon patches from Łukasz Mierzwa 2012-09-20 16:43:20 +02:00
roberto@quantal64 455d041ce4 fixed http router parser 2012-09-20 16:36:08 +02:00
roberto@arch 88205f8dd1 added mongodblog plugin 2012-09-19 15:15:03 +02:00
roberto@centos6 7d9dadddf6 support for loading file for elf sections 2012-09-18 00:42:53 +02:00
roberto@centos6 57e78dd76f implemented rack.input each 2012-09-17 16:26:44 +02:00
roberto@centos6 7d0d2b33f4 improved rack.input disk buffering (second part) 2012-09-17 15:51:31 +02:00
roberto@quantal64 03c739102d reimplementing rack.input without IO wrapper 2012-09-17 13:43:43 +02:00
roberto@quantal64 392bb8d7b1 allows building with python3 and without embedded module 2012-09-17 07:41:11 +02:00
roberto@quantal64 505b35be29 do not try to remove sockets in abstract namespace 2012-09-17 06:47:07 +02:00
roberto@quantal64 fdb002ad9b added %0 - %9 magic vars for splitting paths 2012-09-17 06:33:52 +02:00
roberto@quantal64 557c46373f improved %c 2012-09-15 16:15:47 +02:00
roberto@quantal64 3cde8be505 allows building without embedded and multiple interpreters 2012-09-15 12:14:12 +02:00
roberto@quantal64 c2cf1f3a62 fixed threading + lazy 2012-09-14 15:52:23 +02:00
roberto@quantal64 956cb1f13c Added tag 1.3-rc3 for changeset 883b946db903 2012-09-13 18:47:08 +02:00
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
56 changed files with 1393 additions and 421 deletions
+2
View File
@@ -52,3 +52,5 @@ e1568fd16b7b586cc72deb4dfccbbe64ae0b84df 1.1
24c8fe5a26e557db41cf1cee45b8964d5885e484 1.2-rc1
29de0fb320bc0a1ce84972f9f249360e315d3c10 1.2-rc2
c3cdecbf2bac591336baddd14f9ac22e18d6e200 1.2
14524da00a8b382dffb1d16a0e70cdf4a946f369 1.3-rc2
883b946db9038372cb1d76cb1d328a3093140e7b 1.3-rc3
+1 -1
View File
@@ -29,7 +29,7 @@ plugins =
bin_name = uwsgi
append_version =
plugin_dir = .
embedded_plugins = %(main_plugin)s, ping, cache, nagios, rrdtool, carbon, rpc, corerouter, fastrouter, http, ugreen, signal, syslog, rsyslog, logsocket, router_uwsgi, router_redirect, router_basicauth, zergpool, redislog, router_rewrite, router_http
embedded_plugins = %(main_plugin)s, ping, cache, nagios, rrdtool, carbon, rpc, corerouter, fastrouter, http, ugreen, signal, syslog, rsyslog, logsocket, router_uwsgi, router_redirect, router_basicauth, zergpool, redislog, mongodblog, router_rewrite, router_http
as_shared_library = false
locking = auto
+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);
+1 -1
View File
@@ -102,7 +102,7 @@ uint32_t djb33x_hash(char *key, int keylen) {
}
inline uint64_t uwsgi_cache_get_index(char *key, uint16_t keylen) {
static inline uint64_t uwsgi_cache_get_index(char *key, uint16_t keylen) {
uint32_t hash = djb33x_hash(key, keylen);
+37 -1
View File
@@ -498,7 +498,7 @@ void emperor_add(struct uwsgi_emperor_scanner *ues, char *name, time_t born, cha
if (uwsgi.emperor_tyrant) {
if (uid == 0 || gid == 0) {
uwsgi_log("[emperor-tyrant] invalid permissions for file %s\n", name);
uwsgi_log("[emperor-tyrant] invalid permissions for vassal %s\n", name);
return;
}
}
@@ -1208,6 +1208,9 @@ void emperor_send_stats(int fd) {
if (uwsgi_stats_keylong_comma(us, "gid", (unsigned long long) c_ui->gid))
goto end0;
if (uwsgi_stats_keyval_comma(us, "monitor", c_ui->scanner->arg))
goto end0;
if (uwsgi_stats_keylong(us, "respawns", (unsigned long long) c_ui->respawns))
goto end0;
@@ -1338,3 +1341,36 @@ void uwsgi_check_emperor() {
}
}
void uwsgi_emperor_simple_do(struct uwsgi_emperor_scanner *ues, char *name, char *config, time_t ts, uid_t uid, gid_t gid) {
if (!uwsgi_emperor_is_valid(name))
return;
struct uwsgi_instance *ui_current = emperor_get(name);
if (ui_current) {
// check if uid or gid are changed, in such case, stop the instance
if (uwsgi.emperor_tyrant) {
if (uid != ui_current->uid || gid != ui_current->gid) {
uwsgi_log("[emperor-tyrant] !!! permissions of vassal %s changed. stopping the instance... !!!\n", name);
emperor_stop(ui_current);
return;
}
}
// check if mtime is changed and the uWSGI instance must be reloaded
if (ts > ui_current->last_mod) {
// make a new config (free the old one)
free(ui_current->config);
ui_current->config = config;
ui_current->config_len = strlen(config);
// always respawn (no need for amqp-style rules)
emperor_respawn(ui_current, ts);
}
}
else {
// make a copy of the config as it will be freed
emperor_add(ues, name, ts, uwsgi_str(config), strlen((const char *)config), uid, gid);
}
}
+2 -2
View File
@@ -926,10 +926,10 @@ struct uwsgi_timer *event_queue_ack_timer(int id) {
}
#endif
inline int event_queue_read() {
int event_queue_read() {
return UWSGI_EVENT_IN;
}
inline int event_queue_write() {
int event_queue_write() {
return UWSGI_EVENT_OUT;
}
+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)
+52 -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,41 @@ 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);
}
pid_t uwsgi_daemonize2() {
if (uwsgi.has_emperor) {
logto(uwsgi.daemonize2);
}
else {
if (!uwsgi.is_a_reload) {
uwsgi_log("*** daemonizing uWSGI ***\n");
daemonize(uwsgi.daemonize2);
}
else if (uwsgi.log_reopen) {
logto(uwsgi.daemonize2);
}
}
uwsgi.mypid = getpid();
uwsgi.workers[0].pid = uwsgi.mypid;
if (uwsgi.pidfile && !uwsgi.is_a_reload) {
uwsgi_write_pidfile(uwsgi.pidfile);
}
if (uwsgi.pidfile2 && !uwsgi.is_a_reload) {
uwsgi_write_pidfile(uwsgi.pidfile2);
}
return uwsgi.mypid;
}
+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();
+31 -1
View File
@@ -2,6 +2,27 @@
extern struct uwsgi_server uwsgi;
#ifdef UWSGI_ELF
static void uwsgi_plugin_parse_section(char *filename) {
size_t s_len = 0;
char *buf = uwsgi_elf_section(filename, "uwsgi", &s_len);
if (buf) {
char *p = strtok(buf, "\n");
while(p) {
char *equal = strchr(p, '=');
if (equal) {
*equal = 0;
if (!strcmp(p, "requires")) {
uwsgi_load_plugin(-1, equal+1, NULL);
}
}
p = strtok(NULL, "\n");
}
free(buf);
}
}
#endif
static int plugin_already_loaded(const char *plugin) {
int i;
@@ -79,6 +100,9 @@ void *uwsgi_load_plugin(int modifier, char *plugin, char *has_option) {
// step 1: check for absolute plugin (stop if it fails)
if (strchr(plugin_name, '/')) {
#ifdef UWSGI_ELF
uwsgi_plugin_parse_section(plugin_name);
#endif
plugin_handle = dlopen(plugin_name, RTLD_NOW | RTLD_GLOBAL);
if (!plugin_handle) {
if (!has_option)
@@ -95,6 +119,9 @@ void *uwsgi_load_plugin(int modifier, char *plugin, char *has_option) {
struct uwsgi_string_list *pdir = uwsgi.plugins_dir;
while(pdir) {
plugin_filename = uwsgi_concat3(pdir->value, "/", plugin_name);
#ifdef UWSGI_ELF
uwsgi_plugin_parse_section(plugin_filename);
#endif
plugin_handle = dlopen(plugin_filename, RTLD_NOW | RTLD_GLOBAL);
if (plugin_handle) {
plugin_abs_path = plugin_filename;
@@ -109,6 +136,9 @@ void *uwsgi_load_plugin(int modifier, char *plugin, char *has_option) {
// last step: search in compile-time plugin_dir
if (!plugin_handle) {
plugin_filename = uwsgi_concat3(UWSGI_PLUGIN_DIR, "/", plugin_name);
#ifdef UWSGI_ELF
uwsgi_plugin_parse_section(plugin_filename);
#endif
plugin_handle = dlopen(plugin_filename, RTLD_NOW | RTLD_GLOBAL);
plugin_abs_path = plugin_filename;
//free(plugin_filename);
@@ -117,7 +147,7 @@ void *uwsgi_load_plugin(int modifier, char *plugin, char *has_option) {
success:
if (!plugin_handle) {
if (!has_option)
uwsgi_log( "%s\n", dlerror());
uwsgi_log( "!!! UNABLE to load uWSGI plugin: %s !!!\n", dlerror());
}
else {
char *plugin_entry_symbol = uwsgi_concat2n(plugin_symbol_name_start, strlen(plugin_symbol_name_start)-3, "", 0);
+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, ",");
+147 -7
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);
@@ -1192,7 +1193,7 @@ char *uwsgi_str_contains(char *str, int slen, char what) {
}
// fast compare 2 sized strings
inline int uwsgi_strncmp(char *src, int slen, char *dst, int dlen) {
int uwsgi_strncmp(char *src, int slen, char *dst, int dlen) {
if (slen != dlen)
return 1;
@@ -1202,7 +1203,7 @@ inline int uwsgi_strncmp(char *src, int slen, char *dst, int dlen) {
}
// fast sized check of initial part of a string
inline int uwsgi_starts_with(char *src, int slen, char *dst, int dlen) {
int uwsgi_starts_with(char *src, int slen, char *dst, int dlen) {
if (slen < dlen)
return -1;
@@ -1211,7 +1212,7 @@ inline int uwsgi_starts_with(char *src, int slen, char *dst, int dlen) {
}
// unsized check
inline int uwsgi_startswith(char *src, char *what, int wlen) {
int uwsgi_startswith(char *src, char *what, int wlen) {
int i;
@@ -2148,7 +2149,7 @@ int uwsgi_waitfd_event(int fd, int timeout, int event) {
return ret;
}
inline void *uwsgi_malloc(size_t size) {
void *uwsgi_malloc(size_t size) {
char *ptr = malloc(size);
if (ptr == NULL) {
@@ -2159,7 +2160,7 @@ inline void *uwsgi_malloc(size_t size) {
return ptr;
}
inline void *uwsgi_calloc(size_t size) {
void *uwsgi_calloc(size_t size) {
char *ptr = uwsgi_malloc(size);
memset(ptr, 0, size);
@@ -2530,6 +2531,18 @@ char *uwsgi_open_and_read(char *url, int *size, int add_zero, char *magic_table[
memcpy(buffer, sym_start_ptr, sym_end_ptr - sym_start_ptr);
}
#ifdef UWSGI_ELF
else if (!strncmp("section://", url, 10)) {
size_t s_len = 0;
buffer = uwsgi_elf_section(uwsgi.binary_path, url+10, &s_len);
if (!buffer) {
uwsgi_log("unable to find section %s in %s\n", url+10, uwsgi.binary_path);
exit(1);
}
*size = s_len;
if (add_zero) *size += 1;
}
#endif
// fallback to file
else {
fd = open(url, O_RDONLY);
@@ -3133,7 +3146,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 +3191,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 +4646,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) {
@@ -4706,3 +4750,99 @@ void uwsgi_set_cpu_affinity() {
}
}
#ifdef UWSGI_ELF
#if defined(__linux__)
#include <elf.h>
#endif
char *uwsgi_elf_section(char *filename, char *s, size_t *len) {
struct stat st;
char *output = NULL;
int fd = open(filename, O_RDONLY);
if (fd < 0) {
uwsgi_error_open(filename);
return NULL;
}
if (fstat(fd, &st)) {
uwsgi_error("stat()");
close(fd);
return NULL;
}
if (st.st_size < EI_NIDENT) {
uwsgi_log("invalid elf file: %s\n", filename);
close(fd);
return NULL;
}
char *addr = mmap(NULL, st.st_size , PROT_READ, MAP_PRIVATE, fd, 0);
if (addr == MAP_FAILED) {
uwsgi_error("mmap()");
close(fd);
return NULL;
}
if (addr[0] != ELFMAG0) goto clear;
if (addr[1] != ELFMAG1) goto clear;
if (addr[2] != ELFMAG2) goto clear;
if (addr[3] != ELFMAG3) goto clear;
if (addr[4] == ELFCLASS32) {
// elf header
Elf32_Ehdr *elfh = (Elf32_Ehdr *) addr;
// first section
Elf32_Shdr *sections = ((Elf32_Shdr *) (addr + elfh->e_shoff));
// number of sections
int ns = elfh->e_shnum;
// the names table
Elf32_Shdr *table = &sections[elfh->e_shstrndx];
// string table session pointer
char *names = addr + table->sh_offset;
Elf32_Shdr *ss = NULL; int i;
for(i=0;i<ns;i++) {
char *name = names + sections[i].sh_name;
if (!strcmp(name, s)) {
ss = &sections[i];
break;
}
}
if (ss) {
*len = ss->sh_size;
output = uwsgi_concat2n(addr + ss->sh_offset, ss->sh_size, "", 0);
}
}
else if (addr[4] == ELFCLASS64) {
// elf header
Elf64_Ehdr *elfh = (Elf64_Ehdr *) addr;
// first section
Elf64_Shdr *sections = ((Elf64_Shdr *) (addr + elfh->e_shoff));
// number of sections
int ns = elfh->e_shnum;
// the names table
Elf64_Shdr *table = &sections[elfh->e_shstrndx];
// string table session pointer
char *names = addr + table->sh_offset;
Elf64_Shdr *ss = NULL; int i;
for(i=0;i<ns;i++) {
char *name = names + sections[i].sh_name;
if (!strcmp(name, s)) {
ss = &sections[i];
break;
}
}
if (ss) {
*len = ss->sh_size;
output = uwsgi_concat2n(addr + ss->sh_offset, ss->sh_size, "", 0);
}
}
clear:
close(fd);
munmap(addr, st.st_size);
return output;
}
#endif
+82 -69
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
@@ -477,6 +480,7 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"plugins-list", no_argument, 0, "list enabled plugins", uwsgi_opt_true, &uwsgi.plugins_list, 0},
{"plugin-list", no_argument, 0, "list enabled plugins", uwsgi_opt_true, &uwsgi.plugins_list, 0},
{"autoload", no_argument, 0, "try to automatically load plugins when unknown options are found", uwsgi_opt_true, &uwsgi.autoload, UWSGI_OPT_IMMEDIATE},
{"dlopen", required_argument, 0, "blindly load a shared library", uwsgi_opt_load_dl, NULL, UWSGI_OPT_IMMEDIATE},
{"allowed-modifiers", required_argument, 0, "comma separated list of allowed modifiers", uwsgi_opt_set_str, &uwsgi.allowed_modifiers, 0},
{"remap-modifier", required_argument, 0, "remap request modifier from one id to another", uwsgi_opt_set_str, &uwsgi.remap_modifier, 0},
@@ -657,8 +661,21 @@ void config_magic_table_fill(char *filename, char **magic_table) {
#endif
*tmp = 0;
}
if (uwsgi_get_last_char(magic_table['d'], '/'))
magic_table['c'] = uwsgi_get_last_char(magic_table['d'], '/') + 1;
if (uwsgi_get_last_char(magic_table['d'], '/')) {
magic_table['c'] = uwsgi_str(uwsgi_get_last_char(magic_table['d'], '/') + 1);
if (magic_table['c'][strlen(magic_table['c']) - 1] == '/') {
magic_table['c'][strlen(magic_table['c']) - 1] = 0;
}
}
int base = '0';
char *to_split = uwsgi_str(magic_table['d']);
char *p = strtok(to_split,"/");
while(p && base <= '9') {
magic_table[base] = p;
base++;
p = strtok(NULL, "/");
}
if (tmp)
*tmp = '/';
@@ -1200,7 +1217,7 @@ static void vacuum(void) {
}
}
while (uwsgi_sock) {
if (uwsgi_sock->family == AF_UNIX) {
if (uwsgi_sock->family == AF_UNIX && uwsgi_sock->name[0] != '@') {
if (unlink(uwsgi_sock->name)) {
uwsgi_error("unlink()");
}
@@ -1519,7 +1536,7 @@ static time_t uwsgi_unix_seconds() {
static uint64_t uwsgi_unix_microseconds() {
struct timeval tv;
gettimeofday(&tv, NULL);
return (tv.tv_sec * 1000000) + tv.tv_usec;
return ((uint64_t)tv.tv_sec * 1000000) + tv.tv_usec;
}
static struct uwsgi_clock uwsgi_unix_clock = {
@@ -1938,6 +1955,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 +1997,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) {
@@ -2014,34 +2040,6 @@ int uwsgi_start(void *v_argv) {
pthread_mutex_init(&uwsgi.static_offload_thread_lock, NULL);
}
#ifdef UWSGI_ASYNC
// TODO rewrite to use uwsgi.max_fd
if (uwsgi.async > 1) {
if (!getrlimit(RLIMIT_NOFILE, &uwsgi.rl)) {
if ((unsigned long) uwsgi.rl.rlim_cur < (unsigned long) uwsgi.async) {
uwsgi_log("- your current max open files limit is %lu, this is lower than requested async cores !!! -\n", (unsigned long) uwsgi.rl.rlim_cur);
if (uwsgi.rl.rlim_cur < uwsgi.rl.rlim_max && (unsigned long) uwsgi.rl.rlim_max > (unsigned long) uwsgi.async) {
unsigned long tmp_nofile = (unsigned long) uwsgi.rl.rlim_cur;
uwsgi.rl.rlim_cur = uwsgi.async;
if (!setrlimit(RLIMIT_NOFILE, &uwsgi.rl)) {
uwsgi_log("max open files limit reset to %lu\n", (unsigned long) uwsgi.rl.rlim_cur);
uwsgi.async = uwsgi.rl.rlim_cur;
}
else {
uwsgi.async = (int) tmp_nofile;
}
}
else {
uwsgi.async = uwsgi.rl.rlim_cur;
}
uwsgi_log("- async cores set to %d -\n", uwsgi.async);
}
}
}
#endif
if (uwsgi.requested_max_fd) {
uwsgi.rl.rlim_cur = uwsgi.requested_max_fd;
uwsgi.rl.rlim_max = uwsgi.requested_max_fd;
@@ -2052,11 +2050,24 @@ int uwsgi_start(void *v_argv) {
if (!getrlimit(RLIMIT_NOFILE, &uwsgi.rl)) {
uwsgi.max_fd = uwsgi.rl.rlim_cur;
uwsgi_log_initial("detected max file descriptor number: %d\n", (int) uwsgi.max_fd);
uwsgi_log_initial("detected max file descriptor number: %lu\n", (unsigned long) uwsgi.max_fd);
}
if (uwsgi.async > 1) {
uwsgi_log("async fd table size: %d\n", uwsgi.max_fd);
if ((unsigned long) uwsgi.max_fd < (unsigned long) uwsgi.async) {
uwsgi_log("- your current max open files limit is %lu, this is lower than requested async cores !!! -\n", (unsigned long) uwsgi.max_fd);
uwsgi.rl.rlim_cur = uwsgi.async;
uwsgi.rl.rlim_max = uwsgi.async;
if (!setrlimit(RLIMIT_NOFILE, &uwsgi.rl)) {
uwsgi_log("max open files limit raised to %lu\n", (unsigned long) uwsgi.rl.rlim_cur);
uwsgi.async = uwsgi.rl.rlim_cur;
uwsgi.max_fd = uwsgi.rl.rlim_cur;
}
else {
uwsgi.async = (int) uwsgi.max_fd;
}
}
uwsgi_log("- async cores set to %d - fd table size: %d\n", uwsgi.async, (int) uwsgi.max_fd);
uwsgi.async_waiting_fd_table = malloc(sizeof(struct wsgi_request *) * uwsgi.max_fd);
if (!uwsgi.async_waiting_fd_table) {
uwsgi_error("malloc()");
@@ -2075,10 +2086,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 +2335,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) {
@@ -2361,30 +2385,7 @@ int uwsgi_start(void *v_argv) {
if (uwsgi.daemonize2) {
if (uwsgi.has_emperor) {
logto(uwsgi.daemonize2);
}
else {
if (!uwsgi.is_a_reload) {
uwsgi_log("*** daemonizing uWSGI ***\n");
daemonize(uwsgi.daemonize2);
}
else if (uwsgi.log_reopen) {
logto(uwsgi.daemonize2);
}
}
uwsgi.mypid = getpid();
masterpid = uwsgi.mypid;
uwsgi.workers[0].pid = masterpid;
if (uwsgi.pidfile && !uwsgi.is_a_reload) {
uwsgi_write_pidfile(uwsgi.pidfile);
}
if (uwsgi.pidfile2 && !uwsgi.is_a_reload) {
uwsgi_write_pidfile(uwsgi.pidfile2);
}
masterpid = uwsgi_daemonize2();
}
if (uwsgi.no_server) {
@@ -2396,9 +2397,12 @@ int uwsgi_start(void *v_argv) {
if (!uwsgi.master_process && uwsgi.numproc == 0) {
exit(0);
}
#ifdef UWSGI_MINTERPRETERS
if (!uwsgi.single_interpreter && uwsgi.numproc > 0) {
uwsgi_log("*** uWSGI is running in multiple interpreter mode ***\n");
}
#endif
// check for request plugins, and eventually print a warning
int rp_available = 0;
@@ -2460,15 +2464,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 +3155,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));
}
@@ -3292,6 +3299,12 @@ void uwsgi_opt_pidfile_signal(char *opt, char *pidfile, void *sig) {
exit(0);
}
void uwsgi_opt_load_dl(char *opt, char *value, void *none) {
if (!dlopen(value, RTLD_NOW | RTLD_GLOBAL)) {
uwsgi_log( "%s\n", dlerror());
}
}
void uwsgi_opt_load_plugin(char *opt, char *value, void *none) {
char *p = strtok(uwsgi_concat2(value, ""), ",");
+219 -145
View File
@@ -2,22 +2,37 @@
extern struct uwsgi_server uwsgi;
struct carbon_server_list {
char *value; // server address
int healthy;
int errors;
struct carbon_server_list *next;
};
struct uwsgi_carbon {
struct uwsgi_string_list *servers;
struct carbon_server_list *servers_data;
int freq;
int timeout;
char *id;
int no_workers;
unsigned long long *last_busyness_values;
unsigned long long *current_busyness_values;
int need_retry;
time_t last_update;
time_t next_retry;
int max_retries;
int retry_delay;
} u_carbon;
struct uwsgi_option carbon_options[] = {
{"carbon", required_argument, 0, "push statistics to the specified carbon server", uwsgi_opt_add_string_list, &u_carbon.servers, UWSGI_OPT_MASTER},
{"carbon-timeout", required_argument, 0, "set carbon connection timeout", uwsgi_opt_set_int, &u_carbon.timeout, 0},
{"carbon-freq", required_argument, 0, "set carbon push frequency", uwsgi_opt_set_int, &u_carbon.freq, 0},
{"carbon-timeout", required_argument, 0, "set carbon connection timeout in seconds (default 3)", uwsgi_opt_set_int, &u_carbon.timeout, 0},
{"carbon-freq", required_argument, 0, "set carbon push frequency in seconds (default 60)", uwsgi_opt_set_int, &u_carbon.freq, 0},
{"carbon-id", required_argument, 0, "set carbon id", uwsgi_opt_set_str, &u_carbon.id, 0},
{"carbon-no-workers", no_argument, 0, "disable generation of single worker metrics", uwsgi_opt_true, &u_carbon.no_workers, 0},
{"carbon-max-retry", required_argument, 0, "set maximum number of retries in case of connection errors (default 1)", uwsgi_opt_set_int, &u_carbon.max_retries, 0},
{"carbon-retry-delay", required_argument, 0, "set connection retry delay in seconds (default 7)", uwsgi_opt_set_int, &u_carbon.retry_delay, 0},
{0, 0, 0, 0, 0, 0, 0},
};
@@ -31,12 +46,23 @@ void carbon_post_init() {
if (!u_carbon.servers) return;
while(usl) {
uwsgi_log("added carbon server %s\n", usl->value);
struct carbon_server_list *u_server = uwsgi_calloc(sizeof(struct carbon_server_list));
u_server->value = usl->value;
u_server->healthy = 1;
u_server->errors = 0;
if (u_carbon.servers_data) {
u_server->next = u_carbon.servers_data;
}
u_carbon.servers_data = u_server;
uwsgi_log("[carbon] added server %s\n", usl->value);
usl = usl->next;
}
if (u_carbon.freq < 1) u_carbon.freq = 60;
if (u_carbon.timeout < 1) u_carbon.timeout = 3;
if (u_carbon.max_retries <= 0) u_carbon.max_retries = 1;
if (u_carbon.retry_delay <= 0) u_carbon.retry_delay = 7;
if (!u_carbon.id) {
u_carbon.id = uwsgi_str(uwsgi.sockets->name);
@@ -53,160 +79,208 @@ void carbon_post_init() {
u_carbon.current_busyness_values = uwsgi_calloc(sizeof(unsigned long long) * uwsgi.numproc);
}
// set next update to now()+retry_delay, this way we will have first flush just after start
u_carbon.last_update = uwsgi_now() - u_carbon.freq + u_carbon.retry_delay;
uwsgi_log("[carbon] carbon plugin started, %is frequency, %is timeout, max retries %i, retry delay %is",
u_carbon.freq, u_carbon.timeout, u_carbon.max_retries, u_carbon.retry_delay);
}
int carbon_write(int *fd, char *fmt,...) {
va_list ap;
va_start(ap, fmt);
char ptr[4096];
int rlen;
rlen = vsnprintf(ptr, 4096, fmt, ap);
if (rlen < 1) return 0;
if (write(*fd, ptr, rlen) <= 0) {
uwsgi_error("write()");
return 0;
}
return 1;
}
void carbon_push_stats(int retry_cycle) {
struct carbon_server_list *usl = u_carbon.servers_data;
int i;
int fd;
int wok;
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];
u_carbon.last_busyness_values[i] = uwsgi.workers[i+1].running_time;
}
u_carbon.need_retry = 0;
while(usl) {
if (retry_cycle && usl->healthy)
// skip healthy servers during retry cycle
goto nxt;
if (retry_cycle && usl->healthy == 0)
uwsgi_log("[carbon] Retrying failed server at %s (%d)\n", usl->value, usl->errors);
if (!retry_cycle) {
usl->healthy = 1;
usl->errors = 0;
}
fd = uwsgi_connect(usl->value, u_carbon.timeout, 0);
if (fd < 0) {
uwsgi_log("[carbon] Could not connect to carbon server at %s\n", usl->value);
if (usl->errors < u_carbon.max_retries) {
u_carbon.need_retry = 1;
u_carbon.next_retry = uwsgi_now() + u_carbon.retry_delay;
} else {
uwsgi_log("[carbon] Maximum number of retries for %s (1)\n",
usl->value, u_carbon.max_retries);
usl->healthy = 0;
usl->errors = 0;
}
usl->healthy = 0;
usl->errors++;
goto nxt;
}
// put the socket in non-blocking mode
uwsgi_socket_nb(fd);
unsigned long long total_rss = 0;
unsigned long long total_vsz = 0;
unsigned long long total_tx = 0;
unsigned long long total_avg_rt = 0; // total avg_rt
unsigned long long avg_rt = 0; // per worker avg_rt reported to carbon
unsigned long long active_workers = 0; // number of workers used to calculate total avg_rt
unsigned long long total_busyness = 0;
unsigned long long total_avg_busyness = 0;
unsigned long long worker_busyness = 0;
unsigned long long total_harakiri = 0;
wok = carbon_write(&fd, "uwsgi.%s.%s.requests %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) uwsgi.workers[0].requests, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
for(i=1;i<=uwsgi.numproc;i++) {
total_tx += uwsgi.workers[i].tx;
if (uwsgi.workers[i].cheaped) {
// also if worker is cheaped than we report its average response time as zero, sending last value might be confusing
avg_rt = 0;
worker_busyness = 0;
}
else {
// global average response time is calculated from active/idle workers, cheaped workers are excluded, otherwise it is not accurate
avg_rt = uwsgi.workers[i].avg_response_time;
active_workers++;
total_avg_rt += uwsgi.workers[i].avg_response_time;
// calculate worker busyness
worker_busyness = ((u_carbon.current_busyness_values[i-1]*100) / (u_carbon.freq*1000000));
if (worker_busyness > 100) worker_busyness = 100;
total_busyness += worker_busyness;
// only running workers are counted in total memory stats
total_rss += uwsgi.workers[i].rss_size;
total_vsz += uwsgi.workers[i].vsz_size;
total_harakiri += uwsgi.workers[i].harakiri_count/2;
}
//skip per worker metrics when disabled
if (u_carbon.no_workers) continue;
wok = carbon_write(&fd, "uwsgi.%s.%s.worker%d.requests %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].requests, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.worker%d.rss_size %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].rss_size, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.worker%d.vsz_size %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].vsz_size, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.worker%d.avg_rt %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long) avg_rt, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.worker%d.tx %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].tx, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.worker%d.busyness %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long) worker_busyness, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.worker%d.harakiri %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].harakiri_count/2, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
}
wok = carbon_write(&fd, "uwsgi.%s.%s.rss_size %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) total_rss, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.vsz_size %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) total_vsz, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.avg_rt %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) (active_workers ? total_avg_rt / active_workers : 0), (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.tx %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) total_tx, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
if (active_workers > 0) {
total_avg_busyness = total_busyness / active_workers;
if (total_avg_busyness > 100) total_avg_busyness = 100;
} else {
total_avg_busyness = 0;
}
wok = carbon_write(&fd, "uwsgi.%s.%s.busyness %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) total_avg_busyness, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
wok = carbon_write(&fd, "uwsgi.%s.%s.active_workers %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) active_workers, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
if (uwsgi.cheaper) {
wok = carbon_write(&fd, "uwsgi.%s.%s.cheaped_workers %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) uwsgi.numproc - active_workers, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
}
wok = carbon_write(&fd, "uwsgi.%s.%s.harakiri %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long) total_harakiri, (unsigned long long) uwsgi.current_time);
if (!wok) goto clear;
usl->healthy = 1;
usl->errors = 0;
clear:
close(fd);
nxt:
usl = usl->next;
}
if (!retry_cycle) u_carbon.last_update = uwsgi_now();
if (u_carbon.need_retry)
// timeouts and retries might cause additional lags in carbon cycles, we will adjust timer to fix that
u_carbon.last_update -= u_carbon.timeout;
}
void carbon_master_cycle() {
static time_t last_update = 0;
char ptr[4096];
int rlen, i;
int fd;
struct uwsgi_string_list *usl = u_carbon.servers;
if (!u_carbon.servers) return;
if (!u_carbon.servers) return ;
if (last_update == 0) last_update = uwsgi_now();
// update
if (uwsgi.current_time - last_update >= u_carbon.freq) {
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];
u_carbon.last_busyness_values[i] = uwsgi.workers[i+1].running_time;
}
while(usl) {
fd = uwsgi_connect(usl->value, u_carbon.timeout, 0);
if (fd < 0) goto nxt;
// put the socket in non-blocking mode
uwsgi_socket_nb(fd);
unsigned long long total_rss = 0;
unsigned long long total_vsz = 0;
unsigned long long total_tx = 0;
unsigned long long total_avg_rt = 0; // total avg_rt
unsigned long long avg_rt = 0; // per worker avg_rt reported to carbon
unsigned long long active_workers = 0; // number of workers used to calculate total avg_rt
unsigned long long total_busyness = 0;
unsigned long long total_avg_busyness = 0;
unsigned long long worker_busyness = 0;
unsigned long long total_harakiri = 0;
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.requests %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) uwsgi.workers[0].requests, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
for(i=1;i<=uwsgi.numproc;i++) {
total_tx += uwsgi.workers[i].tx;
if (uwsgi.workers[i].cheaped) {
// also if worker is cheaped than we report its average response time as zero, sending last value might be confusing
avg_rt = 0;
worker_busyness = 0;
}
else {
// global average response time is calcucalted from active/idle workers, cheaped workers are excluded, otherwise it is not accurate
avg_rt = uwsgi.workers[i].avg_response_time;
active_workers++;
total_avg_rt += uwsgi.workers[i].avg_response_time;
// calculate worker busyness
worker_busyness = ((u_carbon.current_busyness_values[i-1]*100) / (u_carbon.freq*1000000));
if (worker_busyness > 100) worker_busyness = 100;
total_busyness += worker_busyness;
// only running workers are counted in total memory stats
total_rss += uwsgi.workers[i].rss_size;
total_vsz += uwsgi.workers[i].vsz_size;
total_harakiri += uwsgi.workers[i].harakiri_count/2;
}
//skip per worker metrics when disabled
if (u_carbon.no_workers) continue;
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.worker%d.requests %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long ) uwsgi.workers[i].requests, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.worker%d.rss_size %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long ) uwsgi.workers[i].rss_size, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.worker%d.vsz_size %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long ) uwsgi.workers[i].vsz_size, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.worker%d.avg_rt %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long ) avg_rt, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.worker%d.tx %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long ) uwsgi.workers[i].tx, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.worker%d.busyness %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long ) worker_busyness, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.worker%d.harakiri %llu %llu\n", uwsgi.hostname, u_carbon.id, i, (unsigned long long ) uwsgi.workers[i].harakiri_count/2, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.rss_size %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) total_rss, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.vsz_size %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) total_vsz, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.avg_rt %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) (active_workers ? total_avg_rt / active_workers : 0), (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.tx %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) total_tx, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
if (active_workers > 0) {
total_avg_busyness = total_busyness / active_workers;
if (total_avg_busyness > 100) total_avg_busyness = 100;
} else {
total_avg_busyness = 0;
}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.busyness %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) total_avg_busyness, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.active_workers %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) active_workers, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
if (uwsgi.cheaper) {
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.cheaped_workers %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) uwsgi.numproc - active_workers, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
}
rlen = snprintf(ptr, 4096, "uwsgi.%s.%s.harakiri %llu %llu\n", uwsgi.hostname, u_carbon.id, (unsigned long long ) total_harakiri, (unsigned long long ) uwsgi.current_time);
if (rlen < 1) goto clear;
if (write(fd, ptr, rlen) <= 0) { uwsgi_error("write()"); goto clear;}
clear:
close(fd);
nxt:
usl = usl->next;
}
last_update = uwsgi_now();
if (uwsgi.current_time - u_carbon.last_update >= u_carbon.freq || uwsgi.cleaning) {
// update
u_carbon.need_retry = 0;
carbon_push_stats(0);
} else if (u_carbon.need_retry && (uwsgi.current_time >= u_carbon.next_retry)) {
// retry failed servers
carbon_push_stats(1);
}
}
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
@@ -11,7 +11,7 @@ time_t uwsgi_monotonic_seconds() {
uint64_t uwsgi_monotonic_microseconds() {
struct timespec ts;
clock_gettime(CLOCK_MONOTONIC, &ts);
return (ts.tv_sec * 1000000) + (ts.tv_nsec/1000);
return ((uint64_t)ts.tv_sec * 1000000) + (ts.tv_nsec/1000);
}
+1 -1
View File
@@ -2,5 +2,5 @@ NAME='clock_monotonic'
CFLAGS = []
LDFLAGS = []
LIBS = []
LIBS = ['-lrt']
GCC_LIST = ['clock_monotonic']
+1 -1
View File
@@ -11,7 +11,7 @@ time_t uwsgi_realtime_seconds() {
uint64_t uwsgi_realtime_microseconds() {
struct timespec ts;
clock_gettime(CLOCK_REALTIME, &ts);
return (ts.tv_sec * 1000000) + (ts.tv_nsec/1000);
return ((uint64_t) ts.tv_sec * 1000000) + (ts.tv_nsec/1000);
}
+1 -1
View File
@@ -2,5 +2,5 @@ NAME='clock_realtime'
CFLAGS = []
LDFLAGS = []
LIBS = []
LIBS = ['-lrt']
GCC_LIST = ['clock_realtime']
+2
View File
@@ -89,7 +89,9 @@ struct corerouter_session {
int instance_fd;
int instance_stopped;
int status;
uint8_t h_pos;
uint16_t pos;
struct uwsgi_gateway_socket *ugs;
+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;
}
+121
View File
@@ -0,0 +1,121 @@
#include "../../uwsgi.h"
#include "client/dbclient.h"
extern struct uwsgi_server uwsgi;
extern struct uwsgi_instance *ui;
// one for each mongodb imperial monitor
struct uwsgi_emperor_mongodb_state {
char *address;
char *collection;
char *json;
};
extern "C" void uwsgi_imperial_monitor_mongodb(struct uwsgi_emperor_scanner *ues) {
struct uwsgi_emperor_mongodb_state *uems = (struct uwsgi_emperor_mongodb_state *) ues->data;
try {
// requested fields
mongo::BSONObj p = BSON( "name" << 1 << "config" << 1 << "ts" << 1 << "uid" << 1 << "gid" << 1 );
mongo::BSONObj q = mongo::fromjson(uems->json);
// the connection object (will be automatically destroyed at each cycle)
mongo::DBClientConnection c;
// set the socket timeout
c.setSoTimeout(uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]);
// connect
c.connect(uems->address);
// run the query
std::auto_ptr<mongo::DBClientCursor> cursor = c.query(uems->collection, q, 0, 0, &p);
while( cursor->more() ) {
mongo::BSONObj p = cursor->next();
// checking for an empty string is not required, but we reduce the load
// in case of badly strctured databases
const char *name = p.getStringField("name");
if (strlen(name) == 0) continue;
const char *config = p.getStringField("config");
time_t vassal_ts = 0;
// ts must be a Date object !!!
mongo::BSONElement ts = p.getField("ts");
if (ts.type() == mongo::Date) {
vassal_ts = ts.date();
}
uid_t vassal_uid = 0;
gid_t vassal_gid = 0;
// check for tyrant mode
if (uwsgi.emperor_tyrant) {
int tmp_uid = p.getIntField("uid");
int tmp_gid = p.getIntField("gid");
if (tmp_uid < 0) tmp_uid = 0;
if (tmp_gid < 0) tmp_gid = 0;
vassal_uid = tmp_uid;
vassal_gid = tmp_gid;
}
uwsgi_emperor_simple_do(ues, (char *) name, (char *) config, vassal_ts/1000, vassal_uid, vassal_gid);
}
// now check for removed instances
struct uwsgi_instance *c_ui = ui->ui_next;
while (c_ui) {
if (c_ui->scanner == ues) {
mongo::BSONObjBuilder b;
b.appendElements(q);
b.append("name", c_ui->name);
mongo::BSONObj q2 = b.obj();
cursor = c.query(uems->collection, q2, 0, 0, &p);
#ifdef UWSGI_DEBUG
uwsgi_log("JSON: %s\n", q2.toString().c_str());
#endif
if (!cursor->more()) {
emperor_stop(c_ui);
}
}
c_ui = c_ui->ui_next;
}
}
catch ( mongo::DBException &e ) {
uwsgi_log("[emperor-mongodb] ERROR(%s/%s): %s\n", uems->address, uems->collection, e.what());
}
}
// setup a new mongodb imperial monitor
extern "C" void uwsgi_imperial_monitor_mongodb_init(struct uwsgi_emperor_scanner *ues) {
// allocate a new state
ues->data = uwsgi_calloc(sizeof(struct uwsgi_emperor_mongodb_state));
size_t arg_len = strlen(ues->arg);
struct uwsgi_emperor_mongodb_state *uems = (struct uwsgi_emperor_mongodb_state *) ues->data;
// parse args/ set defaults
uems->address = (char *) "127.0.0.1:27017";
uems->collection = (char *) "uwsgi.emperor.vassals";
uems->json = (char *) "";
if (arg_len > 10) {
uems->address = uwsgi_str(ues->arg+10);
char *comma = strchr(uems->address, ',');
if (!comma) goto done;
*comma = 0;
uems->collection = comma+1;
comma = strchr(uems->collection, ',');
if (!comma) goto done;
*comma = 0;
uems->json = comma+1;
}
done:
uwsgi_log("[emperor] enabled emperor MongoDB monitor for %s on collection %s\n", uems->address, uems->collection);
}
+14
View File
@@ -0,0 +1,14 @@
#include "../../uwsgi.h"
void uwsgi_imperial_monitor_mongodb(struct uwsgi_emperor_scanner *);
void uwsgi_imperial_monitor_mongodb_init(struct uwsgi_emperor_scanner *);
void emperor_mongodb_init(void) {
uwsgi_register_imperial_monitor("mongodb", uwsgi_imperial_monitor_mongodb_init, uwsgi_imperial_monitor_mongodb);
}
struct uwsgi_plugin emperor_mongodb_plugin = {
.name = "emperor_mongodb",
.on_load = emperor_mongodb_init,
};
+7
View File
@@ -0,0 +1,7 @@
NAME='emperor_mongodb'
CFLAGS = ['-I/usr/include/mongo','-I/usr/local/include/mongo']
LDFLAGS = []
LIBS = ['-lmongoclient', '-lboost_thread','-lboost_filesystem']
GCC_LIST = ['plugin', 'emperor_mongodb.cc']
+1 -35
View File
@@ -7,7 +7,6 @@ extern struct uwsgi_instance *ui;
void uwsgi_imperial_monitor_pg_init(struct uwsgi_emperor_scanner *);
void uwsgi_imperial_monitor_pg(struct uwsgi_emperor_scanner *);
void emperor_pg_init(void);
void emperor_pg_do(struct uwsgi_emperor_scanner *, char *, char *, time_t, uid_t, gid_t);
void emperor_pg_init(void) {
uwsgi_register_imperial_monitor("pg", uwsgi_imperial_monitor_pg_init, uwsgi_imperial_monitor_pg);
@@ -17,39 +16,6 @@ void uwsgi_imperial_monitor_pg_init(struct uwsgi_emperor_scanner *ues) {
uwsgi_log("[emperor] enabled emperor PostgreSQL monitor\n");
}
void emperor_pg_do(struct uwsgi_emperor_scanner *ues, char *name, char *config, time_t ts, uid_t uid, gid_t gid) {
if (!uwsgi_emperor_is_valid(name))
return;
struct uwsgi_instance *ui_current = emperor_get(name);
if (ui_current) {
// check if uid or gid are changed, in such case, stop the instance
if (uwsgi.emperor_tyrant) {
if (uid != ui_current->uid || gid != ui_current->gid) {
uwsgi_log("[emperor-tyrant] !!! permissions of vassal %s changed. stopping the instance... !!!\n", name);
emperor_stop(ui_current);
return;
}
}
// check if mtime is changed and the uWSGI instance must be reloaded
if (ts > ui_current->last_mod) {
// make a new config (free the old one)
free(ui_current->config);
ui_current->config = config;
ui_current->config_len = strlen(config);
// always respawn (no need for amqp-style rules)
emperor_respawn(ui_current, ts);
}
}
else {
// make a copy of the config as it will be freed
emperor_add(ues, name, ts, uwsgi_str(config), strlen(config), uid, gid);
}
}
void uwsgi_imperial_monitor_pg(struct uwsgi_emperor_scanner *ues) {
PGconn *conn = NULL;
@@ -102,7 +68,7 @@ void uwsgi_imperial_monitor_pg(struct uwsgi_emperor_scanner *ues) {
vassal_uid = uwsgi_str_num(q_uid, strlen(q_uid));
vassal_gid = uwsgi_str_num(q_gid, strlen(q_gid));
}
emperor_pg_do(ues, name, config, uwsgi_str_num(ts, len), vassal_uid, vassal_gid);
uwsgi_emperor_simple_do(ues, name, config, uwsgi_str_num(ts, len), vassal_uid, vassal_gid);
}
}
+2
View File
@@ -4,4 +4,6 @@ CFLAGS = []
LDFLAGS = []
LIBS = []
REQUIRES = ['corerouter']
GCC_LIST = ['fastrouter', 'fr_events']
+2 -2
View File
@@ -15,7 +15,7 @@ VALUE uwsgi_fiber_request() {
return Qnil;
}
inline static void fiber_schedule_to_req() {
static inline void fiber_schedule_to_req() {
int id = uwsgi.wsgi_req->async_id;
@@ -33,7 +33,7 @@ inline static void fiber_schedule_to_req() {
}
inline static void fiber_schedule_to_main(struct wsgi_request *wsgi_req) {
static inline void fiber_schedule_to_main(struct wsgi_request *wsgi_req) {
rb_fiber_yield(0, NULL);
uwsgi.wsgi_req = wsgi_req;
+10 -7
View File
@@ -1,27 +1,30 @@
import os
NAME='fiber'
try:
RUBYPATH = os.environ['UWSGICONFIG_RUBYPATH']
except:
RUBYPATH = 'ruby'
NAME='fiber'
CFLAGS = os.popen(RUBYPATH + " -e \"require 'rbconfig';print Config::CONFIG['CFLAGS']\"").read().rstrip().split()
CFLAGS.append('-Wno-unused-parameter')
CFLAGS = os.popen(RUBYPATH + " -e \"require 'rbconfig';print RbConfig::CONFIG['CFLAGS']\"").read().rstrip().split()
CFLAGS.append('-DRUBY19')
CFLAGS.append('-Wno-unused-parameter')
rbconfig = 'RbConfig'
includedir = os.popen(RUBYPATH + " -e \"require 'rbconfig';print Config::CONFIG['rubyhdrdir']\"").read().rstrip()
includedir = os.popen(RUBYPATH + " -e \"require 'rbconfig';print %s::CONFIG['rubyhdrdir']\"" % rbconfig).read().rstrip()
if includedir == 'nil':
includedir = os.popen(RUBYPATH + " -e \"require 'rbconfig';print Config::CONFIG['archdir']\"").read().rstrip()
includedir = os.popen(RUBYPATH + " -e \"require 'rbconfig';print %s::CONFIG['archdir']\"" % rbconfig).read().rstrip()
CFLAGS.append('-I' + includedir)
else:
CFLAGS.append('-I' + includedir)
archdir = os.popen(RUBYPATH + " -e \"require 'rbconfig';print Config::CONFIG['archdir']\"").read().rstrip()
arch = os.popen(RUBYPATH + " -e \"require 'rbconfig';print Config::CONFIG['arch']\"").read().rstrip()
archdir = os.popen(RUBYPATH + " -e \"require 'rbconfig';print %s::CONFIG['archdir']\"" % rbconfig).read().rstrip()
arch = os.popen(RUBYPATH + " -e \"require 'rbconfig';print %s::CONFIG['arch']\"" % rbconfig).read().rstrip()
CFLAGS.append('-I' + archdir)
CFLAGS.append('-I' + archdir + '/' + arch)
CFLAGS.append('-I' + includedir + '/' + arch)
LDFLAGS = []
LIBS = []
+2 -2
View File
@@ -28,7 +28,7 @@ PyObject *py_uwsgi_greenlet_request(PyObject * self, PyObject *args) {
PyMethodDef uwsgi_greenlet_request_method[] = {{"uwsgi_greenlet_request", py_uwsgi_greenlet_request, METH_VARARGS, ""}};
inline static void greenlet_schedule_to_req() {
static inline void greenlet_schedule_to_req() {
int id = uwsgi.wsgi_req->async_id;
@@ -45,7 +45,7 @@ inline static void greenlet_schedule_to_req() {
}
inline static void greenlet_schedule_to_main(struct wsgi_request *wsgi_req) {
static inline void greenlet_schedule_to_main(struct wsgi_request *wsgi_req) {
PyGreenlet_Switch(ugl.main, NULL, NULL);
uwsgi.wsgi_req = wsgi_req;
+9 -5
View File
@@ -506,7 +506,8 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se
goto choose_node;
}
len = cs->recv(&uhttp.cr, cs, hs->buffer + cs->h_pos, UMAX16 - cs->h_pos);
len = cs->recv(&uhttp.cr, cs, hs->buffer + cs->pos, UMAX16 - cs->pos);
#ifdef UWSGI_EVENT_USE_PORT
event_queue_add_fd_read(ucr->queue, cs->fd);
#endif
@@ -517,10 +518,9 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se
break;
}
cs->h_pos += len;
cs->pos += len;
for (j = 0; j < len; j++) {
//uwsgi_log("%d %d %d\n", j, *cs->ptr, cs->rnrn);
if (*hs->ptr == '\r' && (hs->rnrn == 0 || hs->rnrn == 2)) {
hs->rnrn++;
}
@@ -531,6 +531,7 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se
hs->rnrn = 2;
}
else if (*hs->ptr == '\n' && hs->rnrn == 3) {
hs->ptr++;
cs->post_remains = len - (j + 1);
hs->iov_len = http_parse(hs);
@@ -597,7 +598,6 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se
}
cs->pass_fd = is_unix(cs->instance_address, cs->instance_address_len);
cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1);
@@ -620,6 +620,7 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se
else {
hs->rnrn = 0;
}
hs->ptr++;
}
@@ -767,7 +768,6 @@ To have a reliable implementation, we need to reset a bunch of values
hs->ptr = hs->buffer;
hs->rnrn = 0;
cs->pos = 0;
cs->h_pos = 0;
hs->received_body = 0;
cs->post_cl = 0;
cs->instance_fd = -1;
@@ -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
@@ -4,4 +4,6 @@ CFLAGS = []
LDFLAGS = []
LIBS = []
REQUIRES = ['corerouter']
GCC_LIST = ['http']
+164
View File
@@ -0,0 +1,164 @@
#include "../../uwsgi.h"
extern struct uwsgi_server uwsgi;
struct uwsgi_mongodb_header {
int32_t len;
int32_t request_id;
int32_t response_id;
int32_t opcode;
};
struct uwsgi_mongodb_state {
int fd;
char *address;
int32_t base_len;
struct uwsgi_mongodb_header header;
int32_t flags;
char *collection;
int32_t bson_base_len;
int32_t bson_len;
int64_t ts;
int32_t bson_node_len;
char *bson_node;
int32_t bson_msg_len;
struct iovec iovec[13];
};
ssize_t uwsgi_mongodb_logger(struct uwsgi_logger *ul, char *message, size_t len) {
struct uwsgi_mongodb_state *ums = NULL;
if (!ul->configured) {
ul->data = uwsgi_calloc(sizeof(struct uwsgi_mongodb_state));
ums = (struct uwsgi_mongodb_state *) ul->data;
// full default
if (ul->arg == NULL) {
ums->address = uwsgi_str("127.0.0.1:27017");
ums->collection = "uwsgi.logs";
ums->bson_node = uwsgi.hostname;
ums->bson_node_len = uwsgi.hostname_len;
goto done;
}
ums->address = uwsgi_str(ul->arg);
char *collection = strchr(ums->address, ',');
// default to uwsgi.logs
if (!collection) {
ums->collection = "uwsgi.logs";
ums->bson_node = uwsgi.hostname;
ums->bson_node_len = uwsgi.hostname_len;
goto done;
}
*collection = 0;
ums->collection = collection+1;
char *node = strchr(ums->collection, ',');
// default to hostname
if (!node) {
ums->bson_node = uwsgi.hostname;
ums->bson_node_len = uwsgi.hostname_len;
goto done;
}
*node = 0;
ums->bson_node = node+1;
ums->bson_node_len = strlen(ums->bson_node)+1;
done:
ums->fd = -1;
// header
ums->iovec[0].iov_base = &ums->header;
ums->iovec[0].iov_len = sizeof(struct uwsgi_mongodb_header);
// OPCODE INSERT
ums->header.opcode = 2002;
// flags
ums->iovec[1].iov_base = &ums->flags;
ums->iovec[1].iov_len = sizeof(int32_t);
// collection name
ums->iovec[2].iov_base = ums->collection;
ums->iovec[2].iov_len = strlen(ums->collection)+1;
// BSON len
ums->iovec[3].iov_base = &ums->bson_len;
ums->iovec[3].iov_len = sizeof(int32_t);
// BSON node
ums->iovec[4].iov_base = "\x02node\0";
ums->iovec[4].iov_len = 6;
ums->iovec[5].iov_base = &ums->bson_node_len;
ums->iovec[5].iov_len = sizeof(int32_t);
ums->iovec[6].iov_base = ums->bson_node;
ums->iovec[6].iov_len = ums->bson_node_len;
// BSON timestamp (ts)
ums->iovec[7].iov_base = "\x09ts\0";
ums->iovec[7].iov_len = 4;
ums->iovec[8].iov_base = &ums->ts;
ums->iovec[8].iov_len = sizeof(int64_t);
// BSON msg
ums->iovec[9].iov_base = "\2msg\0";
ums->iovec[9].iov_len = 5;
ums->iovec[10].iov_base = &ums->bson_msg_len;
ums->iovec[10].iov_len = sizeof(int32_t);
// iov 11 is reset at each cycle
// ...
// BSON end (msg_zero + bson_zero);
ums->iovec[12].iov_base = "\0\0";
ums->iovec[12].iov_len = 2;
ums->bson_base_len = ums->iovec[3].iov_len + ums->iovec[4].iov_len + ums->iovec[5].iov_len + ums->iovec[6].iov_len +
ums->iovec[7].iov_len + ums->iovec[8].iov_len + ums->iovec[9].iov_len + ums->iovec[10].iov_len +
ums->iovec[12].iov_len;
ums->base_len = ums->iovec[0].iov_len + ums->iovec[1].iov_len + ums->iovec[2].iov_len + ums->bson_base_len;
ul->configured = 1;
}
ums = (struct uwsgi_mongodb_state *) ul->data;
if (ums->fd == -1) {
ums->fd = uwsgi_connect(ums->address, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], 0);
}
if (ums->fd == -1) return -1;
// fix the packet
ums->bson_msg_len = len+1;
ums->bson_len = ums->bson_base_len + len;
ums->header.len = ums->base_len + len;
ums->header.request_id++;
// get milliseconds time
ums->ts = uwsgi_micros()/1000;
ums->iovec[11].iov_base = message;
ums->iovec[11].iov_len = len;
ssize_t ret = writev(ums->fd, ums->iovec, 13);
if (ret <= 0) {
close(ums->fd);
ums->fd = -1;
return -1;
}
return ret;
}
void uwsgi_mongodblog_register() {
uwsgi_register_logger("mongodblog", uwsgi_mongodb_logger);
}
struct uwsgi_plugin mongodblog_plugin = {
.name = "mongodblog",
.on_load = uwsgi_mongodblog_register,
};
+6
View File
@@ -0,0 +1,6 @@
NAME='mongodblog'
CFLAGS = []
LDFLAGS = []
LIBS = []
GCC_LIST = ['mongodblog_plugin']
+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;
}
+4
View File
@@ -483,7 +483,9 @@ PyObject *uwsgi_uwsgi_loader(void *arg1) {
PyObject *tmp_callable;
PyObject *applications;
#ifdef UWSGI_EMBEDDED
PyObject *uwsgi_dict = get_uwsgi_pydict("uwsgi");
#endif
char *module = (char *) arg1;
@@ -506,8 +508,10 @@ PyObject *uwsgi_uwsgi_loader(void *arg1) {
return NULL;
}
#ifdef UWSGI_EMBEDDED
applications = PyDict_GetItemString(uwsgi_dict, "applications");
if (applications && PyDict_Check(applications)) return applications;
#endif
applications = PyDict_GetItemString(wsgi_dict, "applications");
if (applications && PyDict_Check(applications)) return applications;
+14 -3
View File
@@ -282,6 +282,7 @@ realstuff:
if (uwsgi.has_threads)
PyGILState_Ensure();
// no need to worry about freeing memory
#ifdef UWSGI_EMBEDDED
PyObject *uwsgi_dict = get_uwsgi_pydict("uwsgi");
if (uwsgi_dict) {
PyObject *ae = PyDict_GetItemString(uwsgi_dict, "atexit");
@@ -289,6 +290,7 @@ realstuff:
python_call(ae, PyTuple_New(0), 0, NULL);
}
}
#endif
// this part is a 1:1 copy of mod_wsgi 3.x
// it is required to fix some atexit bug with python 3
@@ -528,8 +530,6 @@ next:
#ifdef UWSGI_EMBEDDED
PyDoc_STRVAR(uwsgi_py_doc, "uWSGI api module.");
#endif
#ifdef PYTHREE
static PyModuleDef uwsgi_module3 = {
@@ -544,7 +544,6 @@ PyObject *init_uwsgi3(void) {
}
#endif
#ifdef UWSGI_EMBEDDED
void init_uwsgi_embedded_module() {
PyObject *new_uwsgi_module, *zero;
int i;
@@ -1050,6 +1049,11 @@ void uwsgi_python_init_apps() {
struct http_status_codes *http_sc;
// lazy ?
if (uwsgi.mywid > 0) {
UWSGI_GET_GIL;
}
#ifndef UWSGI_PYPY
// prepare for stack suspend/resume
if (uwsgi.async > 1) {
@@ -1160,6 +1164,7 @@ next:
}
#endif
#ifdef UWSGI_EMBEDDED
PyObject *uwsgi_dict = get_uwsgi_pydict("uwsgi");
if (uwsgi_dict) {
up.after_req_hook = PyDict_GetItemString(uwsgi_dict, "after_req_hook");
@@ -1169,6 +1174,12 @@ next:
Py_INCREF(up.after_req_hook_args);
}
}
#endif
// lazy ?
if (uwsgi.mywid > 0) {
UWSGI_RELEASE_GIL;
}
}
+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;
+4 -2
View File
@@ -152,14 +152,14 @@ static PyObject *uwsgi_Input_read(uwsgi_Input *self, PyObject *args) {
if (uwsgi_waitfd(self->wsgi_req->poll.fd, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]) <= 0) {
free(tmp_buf);
UWSGI_GET_GIL
return PyErr_Format(PyExc_IOError, "error waiting for wsgi.input data: Content-Length %llu received %llu", (unsigned long long) self->wsgi_req->post_cl, (unsigned long long) self->wsgi_req->post_cl - remains);
return PyErr_Format(PyExc_IOError, "error waiting for wsgi.input data: Content-Length %llu requested %llu received %llu", (unsigned long long) self->wsgi_req->post_cl, (unsigned long long) remains + tmp_pos, (unsigned long long) tmp_pos);
}
rlen = read(self->wsgi_req->poll.fd, tmp_buf+tmp_pos, remains);
if (rlen <= 0) {
free(tmp_buf);
UWSGI_GET_GIL
return PyErr_Format(PyExc_IOError, "error reading wsgi.input data: Content-Length %llu received %llu", (unsigned long long) self->wsgi_req->post_cl, (unsigned long long) self->wsgi_req->post_cl - remains);
return PyErr_Format(PyExc_IOError, "error reading wsgi.input data: Content-Length %llu requested %llu received %llu", (unsigned long long) self->wsgi_req->post_cl, (unsigned long long) remains + tmp_pos, (unsigned long long) tmp_pos);
}
tmp_pos += rlen;
remains -= rlen;
@@ -510,7 +510,9 @@ int uwsgi_request_wsgi(struct wsgi_request *wsgi_req) {
}
// this object must be freed/cleared always
#ifdef UWSGI_ASYNC
end:
#endif
if (wsgi_req->async_input) {
Py_DECREF((PyObject *)wsgi_req->async_input);
}
+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:
+84 -16
View File
@@ -29,6 +29,7 @@ struct uwsgi_option uwsgi_rack_options[] = {
{"rb-threads", required_argument, 0, "set the number of ruby threads to run", uwsgi_opt_set_int, &ur.rb_threads, 0},
{"rbthreads", required_argument, 0, "set the number of ruby threads to run", uwsgi_opt_set_int, &ur.rb_threads, 0},
{"ruby-threads", required_argument, 0, "set the number of ruby threads to run", uwsgi_opt_set_int, &ur.rb_threads, 0},
{"rb-patch-rack-bodyproxy", no_argument, 0, "some specific (old) combos of ruby 1.9+rack could require that hack...", uwsgi_opt_true, &ur.patch_bodyproxy, 0},
#endif
{0, 0, 0, 0, 0, 0 ,0},
@@ -79,6 +80,14 @@ VALUE rb_uwsgi_io_gets(VALUE obj, VALUE args) {
struct wsgi_request *wsgi_req;
VALUE line;
Data_Get_Struct(obj, struct wsgi_request, wsgi_req);
char linebuf[4096];
if (wsgi_req->async_post) {
if (fgets(linebuf, 4096, (FILE *) wsgi_req->async_post) == NULL) {
return Qnil;
}
return rb_str_new2(linebuf);
}
// return a line of body
for(i=wsgi_req->buf_pos;i<wsgi_req->post_cl;i++) {
@@ -100,12 +109,18 @@ VALUE rb_uwsgi_io_gets(VALUE obj, VALUE args) {
VALUE rb_uwsgi_io_each(VALUE obj, VALUE args) {
struct wsgi_request *wsgi_req;
Data_Get_Struct(obj, struct wsgi_request, wsgi_req);
if (!rb_block_given_p())
rb_raise(rb_eArgError, "Expected block on rack.input 'each' method");
// yield strings chunks
rb_raise(rb_eRuntimeError, "rack.input::each is not implemented (req %p)\n", wsgi_req);
for(;;) {
VALUE chunk = rb_uwsgi_io_gets(obj, Qnil);
if (chunk == Qnil) {
return Qnil;
}
rb_yield(chunk);
}
// never here
return Qnil;
}
@@ -116,11 +131,56 @@ VALUE rb_uwsgi_io_read(VALUE obj, VALUE args) {
VALUE chunk;
unsigned int chunk_size;
if (!wsgi_req->post_cl || wsgi_req->buf_pos >= wsgi_req->post_cl) {
/*
When EOF is reached, this method returns nil if length is given and not nil, or "" if length is not given or is nil.
If buffer is given, then the read data will be placed into buffer instead of a newly created String object.
*/
// --- disk buffering ---
if (wsgi_req->async_post) {
// 0 size, read the whole body from the file...
if (RARRAY_LEN(args) == 0) {
char *tmp_chunk = uwsgi_malloc(wsgi_req->post_cl);
size_t rlen = fread(tmp_chunk, 1, wsgi_req->post_cl, (FILE *) wsgi_req->async_post);
if (rlen == 0) {
free(tmp_chunk);
return rb_str_new("", 0);
}
// return a new string
chunk = rb_str_new(tmp_chunk, rlen);
free(tmp_chunk);
return chunk;
}
// size specified
else if (RARRAY_LEN(args) > 0) {
chunk_size = NUM2UINT(RARRAY_PTR(args)[0]);
char *tmp_chunk = uwsgi_malloc(chunk_size);
size_t rlen = fread(tmp_chunk, 1, chunk_size, (FILE *) wsgi_req->async_post);
// error, return Qnil
if (rlen == 0) {
free(tmp_chunk);
return Qnil;
}
// push in the specified buffer
if (RARRAY_LEN(args) > 1) {
rb_str_cat(RARRAY_PTR(args)[1], tmp_chunk, rlen);
free(tmp_chunk);
return RARRAY_PTR(args)[1];
}
// return a new string
chunk = rb_str_new(tmp_chunk, rlen);
free(tmp_chunk);
return chunk;
}
// never happend...
return Qnil;
}
// --- memory buffering ---
// first check for virtual EOF
if (!wsgi_req->post_cl || wsgi_req->buf_pos >= wsgi_req->post_cl) {
if (RARRAY_LEN(args) > 0) {
if (RARRAY_PTR(args)[0] == Qnil) {
return rb_str_new("", 0);
@@ -163,8 +223,14 @@ VALUE rb_uwsgi_io_rewind(VALUE obj, VALUE args) {
return Qnil;
}
wsgi_req->buf_pos = 0;
// buffered to disk ?
if (wsgi_req->async_post) {
rewind((FILE *) wsgi_req->async_post);
}
// or memory ???
else {
wsgi_req->buf_pos = 0;
}
return Qnil;
}
@@ -673,6 +739,11 @@ int uwsgi_rack_request(struct wsgi_request *wsgi_req) {
struct http_status_codes *http_sc;
if (!ur.call) {
internal_server_error(wsgi_req, "Ruby application not found");
return -1;
}
/* Standard RACK request */
if (!wsgi_req->uh.pktsize) {
uwsgi_log("Invalid RACK request. skip.\n");
@@ -748,12 +819,7 @@ int uwsgi_rack_request(struct wsgi_request *wsgi_req) {
VALUE dws_wr = Data_Wrap_Struct(ur.rb_uwsgi_io_class, 0, 0, wsgi_req);
if (wsgi_req->async_post) {
rb_hash_aset(env, rb_str_new2("rack.input"), rb_funcall( rb_const_get(rb_cObject, rb_intern("IO")), rb_intern("new"), 2, INT2NUM(fileno((FILE*)wsgi_req->async_post)), rb_str_new("r",1) ));
}
else {
rb_hash_aset(env, rb_str_new2("rack.input"), rb_funcall(ur.rb_uwsgi_io_class, rb_intern("new"), 1, dws_wr ));
}
rb_hash_aset(env, rb_str_new2("rack.input"), rb_funcall(ur.rb_uwsgi_io_class, rb_intern("new"), 1, dws_wr ));
rb_hash_aset(env, rb_str_new2("rack.errors"), rb_funcall( rb_const_get(rb_cObject, rb_intern("IO")), rb_intern("new"), 2, INT2NUM(2), rb_str_new("w",1) ));
@@ -963,9 +1029,11 @@ VALUE init_rack_app( VALUE script ) {
VALUE rack = rb_const_get(rb_cObject, rb_intern("Rack"));
#ifdef RUBY19
VALUE ret = rb_protect(uwsgi_rack_patch_body_proxy, rack, &error);
if (!error && ret != Qnil) {
uwsgi_log("Rack::BodyProxy successfully patched for ruby 1.9.x\n");
if (ur.patch_bodyproxy) {
VALUE ret = rb_protect(uwsgi_rack_patch_body_proxy, rack, &error);
if (!error && ret != Qnil) {
uwsgi_log("Rack::BodyProxy successfully patched for ruby 1.9.x\n");
}
}
#endif
+3
View File
@@ -72,6 +72,9 @@ struct uwsgi_rack {
char *gemset;
int rb_threads;
#ifdef RUBY19
int patch_bodyproxy;
#endif
};
+2
View File
@@ -4,4 +4,6 @@ CFLAGS = []
LDFLAGS = []
LIBS = []
REQUIRES = ['corerouter']
GCC_LIST = ['rawrouter', 'rr_events']
+42 -35
View File
@@ -10,7 +10,7 @@ struct uwsgi_redislog_state {
char msgsize[11];
struct iovec iovec[7];
char response[8];
} uredislog;
};
static char *uwsgi_redis_logger_build_command(char *src) {
ssize_t len = 4096;
@@ -53,95 +53,102 @@ static char *uwsgi_redis_logger_build_command(char *src) {
ssize_t uwsgi_redis_logger(struct uwsgi_logger *ul, char *message, size_t len) {
ssize_t ret,ret2;
struct uwsgi_redislog_state *uredislog = NULL;
if (!ul->configured) {
if (!ul->data) {
ul->data = uwsgi_calloc(sizeof(struct uwsgi_redislog_state));
uredislog = (struct uwsgi_redislog_state *) ul->data;
}
if (ul->arg != NULL) {
char *logarg = uwsgi_str(ul->arg);
char *comma1 = strchr(logarg, ',');
if (!comma1) {
uredislog.address = logarg;
uredislog->address = logarg;
goto done;
}
*comma1 = 0;
uredislog.address = logarg;
uredislog->address = logarg;
comma1++;
if (*comma1 == 0) goto done;
char *comma2 = strchr(comma1,',');
if (!comma2) {
uredislog.command = uwsgi_redis_logger_build_command(comma1);
uredislog->command = uwsgi_redis_logger_build_command(comma1);
goto done;
}
*comma2 = 0;
uredislog.command = uwsgi_redis_logger_build_command(comma1);
uredislog->command = uwsgi_redis_logger_build_command(comma1);
comma2++;
if (*comma2 == 0) goto done;
uredislog.prefix = comma2;
uredislog->prefix = comma2;
}
done:
if (!uredislog.address) uredislog.address = uwsgi_str("127.0.0.1:6379");
if (!uredislog.command) uredislog.command = "*3\r\n$7\r\npublish\r\n$5\r\nuwsgi\r\n";
if (!uredislog.prefix) uredislog.prefix = "";
if (!uredislog->address) uredislog->address = uwsgi_str("127.0.0.1:6379");
if (!uredislog->command) uredislog->command = "*3\r\n$7\r\npublish\r\n$5\r\nuwsgi\r\n";
if (!uredislog->prefix) uredislog->prefix = "";
uredislog.fd = -1;
uredislog->fd = -1;
uredislog.iovec[0].iov_base = uredislog.command;
uredislog.iovec[0].iov_len = strlen(uredislog.command);
uredislog.iovec[1].iov_base = "$";
uredislog.iovec[1].iov_len = 1;
uredislog->iovec[0].iov_base = uredislog->command;
uredislog->iovec[0].iov_len = strlen(uredislog->command);
uredislog->iovec[1].iov_base = "$";
uredislog->iovec[1].iov_len = 1;
uredislog.iovec[2].iov_base = uredislog.msgsize;
uredislog->iovec[2].iov_base = uredislog->msgsize;
uredislog.iovec[3].iov_base = "\r\n";
uredislog.iovec[3].iov_len = 2;
uredislog->iovec[3].iov_base = "\r\n";
uredislog->iovec[3].iov_len = 2;
uredislog.iovec[4].iov_base = uredislog.prefix;
uredislog.iovec[4].iov_len = strlen(uredislog.prefix);
uredislog->iovec[4].iov_base = uredislog->prefix;
uredislog->iovec[4].iov_len = strlen(uredislog->prefix);
uredislog.iovec[6].iov_base = "\r\n";
uredislog.iovec[6].iov_len = 2;
uredislog->iovec[6].iov_base = "\r\n";
uredislog->iovec[6].iov_len = 2;
ul->configured = 1;
}
if (uredislog.fd == -1) {
uredislog.fd = uwsgi_connect(uredislog.address, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], 0);
uredislog = (struct uwsgi_redislog_state *) ul->data;
if (uredislog->fd == -1) {
uredislog->fd = uwsgi_connect(uredislog->address, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], 0);
}
if (uredislog.fd == -1) return -1;
if (uredislog->fd == -1) return -1;
// drop newline
if (message[len-1] == '\n') len--;
uwsgi_num2str2(len + uredislog.iovec[4].iov_len, uredislog.msgsize);
uredislog.iovec[2].iov_len = strlen(uredislog.msgsize);
uwsgi_num2str2(len + uredislog->iovec[4].iov_len, uredislog->msgsize);
uredislog->iovec[2].iov_len = strlen(uredislog->msgsize);
uredislog.iovec[5].iov_base = message;
uredislog.iovec[5].iov_len = len;
uredislog->iovec[5].iov_base = message;
uredislog->iovec[5].iov_len = len;
ret = writev(uredislog.fd, uredislog.iovec, 7);
ret = writev(uredislog->fd, uredislog->iovec, 7);
if (ret <= 0) {
close(uredislog.fd);
uredislog.fd = -1;
close(uredislog->fd);
uredislog->fd = -1;
return -1;
}
again:
// read til a \n is found (ugly but fast)
ret2 = read(uredislog.fd, uredislog.response, 8);
ret2 = read(uredislog->fd, uredislog->response, 8);
if (ret2 <= 0) {
close(uredislog.fd);
uredislog.fd = -1;
close(uredislog->fd);
uredislog->fd = -1;
return -1;
}
if (!memchr(uredislog.response, '\n', ret2)) {
if (!memchr(uredislog->response, '\n', ret2)) {
goto again;
}
+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) {
+2 -2
View File
@@ -27,7 +27,7 @@ PyObject *py_uwsgi_stackless_request(PyObject * self, PyObject *args) {
PyMethodDef uwsgi_stackless_request_method[] = {{"uwsgi_stackless_request", py_uwsgi_stackless_request, METH_VARARGS, ""}};
inline static void stackless_schedule_to_req() {
static inline static void stackless_schedule_to_req() {
int id = uwsgi.wsgi_req->async_id;
@@ -47,7 +47,7 @@ inline static void stackless_schedule_to_req() {
}
inline static void stackless_schedule_to_main(struct wsgi_request *wsgi_req) {
static inline static void stackless_schedule_to_main(struct wsgi_request *wsgi_req) {
PyStackless_Schedule(Py_None, 1);
uwsgi.wsgi_req = wsgi_req;
+2 -2
View File
@@ -32,7 +32,7 @@ void u_green_request() {
uwsgi.wsgi_req->suspended = 0;
}
inline static void u_green_schedule_to_req() {
static inline void u_green_schedule_to_req() {
int id = uwsgi.wsgi_req->async_id;
@@ -58,7 +58,7 @@ inline static void u_green_schedule_to_req() {
}
inline static void u_green_schedule_to_main(struct wsgi_request *wsgi_req) {
static inline void u_green_schedule_to_main(struct wsgi_request *wsgi_req) {
if (uwsgi.p[wsgi_req->uh.modifier1]->suspend) {
uwsgi.p[wsgi_req->uh.modifier1]->suspend(wsgi_req);
+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";
+45 -21
View File
@@ -21,7 +21,7 @@ extern "C" {
wsgi_req->method_len, wsgi_req->method, wsgi_req->uri_len, wsgi_req->uri, wsgi_req->remote_addr_len, wsgi_req->remote_addr); else uwsgi_log_verbose("%s %s [%s line %d] \n",x, strerror(errno), __FILE__, __LINE__);
#define uwsgi_debug(x, ...) uwsgi_log("[uWSGI DEBUG] " x, __VA_ARGS__);
#define uwsgi_rawlog(x) if (write(2, x, strlen(x)) != strlen(x)) uwsgi_error("write()")
#define uwsgi_str(x) uwsgi_concat2(x, "")
#define uwsgi_str(x) uwsgi_concat2(x, (char *)"")
#define uwsgi_notify(x) if (uwsgi.notify) uwsgi.notify(x)
#define uwsgi_notify_ready() uwsgi.shared->ready = 1 ; if (uwsgi.notify_ready) uwsgi.notify_ready()
@@ -34,7 +34,7 @@ extern "C" {
#define thunder_lock if (uwsgi.threads > 1 && !uwsgi.is_et) {pthread_mutex_lock(&uwsgi.thunder_mutex);}
#define thunder_unlock if (uwsgi.threads > 1 && !uwsgi.is_et) {pthread_mutex_unlock(&uwsgi.thunder_mutex);}
#define uwsgi_check_scheme(file) (!uwsgi_startswith(file, "http://", 7) || !uwsgi_startswith(file, "data://", 7) || !uwsgi_startswith(file, "sym://", 6) || !uwsgi_startswith(file, "fd://", 5) || !uwsgi_startswith(file, "exec://", 7))
#define uwsgi_check_scheme(file) (!uwsgi_startswith(file, "http://", 7) || !uwsgi_startswith(file, "data://", 7) || !uwsgi_startswith(file, "sym://", 6) || !uwsgi_startswith(file, "fd://", 5) || !uwsgi_startswith(file, "exec://", 7) || !uwsgi_startswith(file, "section://", 10))
#define ushared uwsgi.shared
@@ -73,10 +73,6 @@ extern char UWSGI_EMBED_CONFIG;
extern char UWSGI_EMBED_CONFIG_END;
#endif
#ifdef __clang__
#define inline
#endif
#define UDEP(pname) extern struct uwsgi_plugin pname##_plugin;
#define ULEP(pname)\
@@ -368,6 +364,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 +709,8 @@ struct uwsgi_plugin {
int (*mule)(char *);
int (*mule_msg)(char *, size_t);
void (*master_cleanup) (void);
};
#ifdef UWSGI_PCRE
@@ -1386,6 +1395,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 +1517,8 @@ struct uwsgi_server {
int to_hell;
int to_outworld;
int cleaning;
int ready_to_die;
int ready_to_reload;
@@ -2160,7 +2175,7 @@ struct wsgi_request *find_wsgi_req_by_id(int);
void async_add_fd_write(struct wsgi_request *, int, int);
void async_add_fd_read(struct wsgi_request *, int, int);
inline struct wsgi_request *next_wsgi_req(struct wsgi_request *);
struct wsgi_request *next_wsgi_req(struct wsgi_request *);
void async_add_timeout(struct wsgi_request *, int);
@@ -2221,8 +2236,8 @@ void uwsgi_opt_ldap_dump_ldif(char *, char *, void *);
void uwsgi_ldap_config(char *);
#endif
inline int uwsgi_strncmp(char *, int, char *, int);
inline int uwsgi_startswith(char *, char *, int);
int uwsgi_strncmp(char *, int, char *, int);
int uwsgi_startswith(char *, char *, int);
char *uwsgi_concat(int, ...);
@@ -2302,9 +2317,8 @@ char *uwsgi_cache_get(char *, uint16_t, uint64_t *);
uint32_t uwsgi_cache_exists(char *, uint16_t);
inline void *uwsgi_malloc(size_t);
inline void *uwsgi_calloc(size_t);
void *uwsgi_malloc(size_t);
void *uwsgi_calloc(size_t);
int event_queue_init(void);
@@ -2519,14 +2533,8 @@ struct rb_root
void rb_insert_color(struct rb_node *, struct rb_root *);
void rb_erase(struct rb_node *, struct rb_root *);
#ifdef __clang__
void rb_link_node(struct rb_node *, struct rb_node *,
struct rb_node **);
#else
inline void rb_link_node(struct rb_node *, struct rb_node *,
struct rb_node **);
#endif
struct uwsgi_rb_timer {
@@ -2565,8 +2573,8 @@ struct uwsgi_async_request {
struct uwsgi_async_request *next;
};
inline int event_queue_read(void);
inline int event_queue_write(void);
int event_queue_read(void);
int event_queue_write(void);
void uwsgi_help(char *opt, char *val, void *);
@@ -2591,7 +2599,7 @@ int uwsgi_amqp_consume_queue(int, char *, char *, char *, char *, char *, char *
char *uwsgi_amqp_consume(int, uint64_t *, char **);
int uwsgi_file_serve(struct wsgi_request *, char *, uint16_t, char *, uint16_t, int);
inline int uwsgi_starts_with(char *, int, char *, int);
int uwsgi_starts_with(char *, int, char *, int);
#ifdef __sun__
time_t timegm(struct tm *);
@@ -2657,6 +2665,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 +2897,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 *);
@@ -2900,6 +2912,7 @@ void uwsgi_opt_add_socket(char *, char *, void *);
void uwsgi_opt_add_lazy_socket(char *, char *, void *);
void uwsgi_opt_add_cron(char *, char *, void *);
void uwsgi_opt_load_plugin(char *, char *, void *);
void uwsgi_opt_load_dl(char *, char *, void *);
void uwsgi_opt_load(char *, char *, void *);
void uwsgi_opt_cluster_log(char *, char *, void *);
void uwsgi_opt_cluster_reload(char *, char *, void *);
@@ -3124,6 +3137,7 @@ void uwsgi_logvar_add(struct wsgi_request *, char *, uint8_t, char *, uint8_t);
struct uwsgi_emperor_scanner {
char *arg;
int fd;
void *data;
void (*event_func)(struct uwsgi_emperor_scanner *);
struct uwsgi_imperial_monitor *monitor;
struct uwsgi_emperor_scanner *next;
@@ -3222,6 +3236,16 @@ 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);
pid_t uwsgi_daemonize2();
void uwsgi_emperor_simple_do(struct uwsgi_emperor_scanner *, char *, char *, time_t, uid_t, gid_t);
#if defined(__linux__)
#define UWSGI_ELF
char *uwsgi_elf_section(char *, char *, size_t *);
#endif
void uwsgi_check_emperor(void);
#ifdef UWSGI_AS_SHARED_LIBRARY
+31 -4
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '1.3-dev'
uwsgi_version = '1.3-rc4'
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
@@ -1016,6 +1020,8 @@ def build_plugin(path, uc, cflags, ldflags, libs, name = None):
import uwsgiplugin as up
reload(up)
requires = []
p_cflags = cflags[:]
p_ldflags = ldflags[:]
@@ -1023,6 +1029,11 @@ def build_plugin(path, uc, cflags, ldflags, libs, name = None):
p_ldflags += up.LDFLAGS
p_libs = up.LIBS
try:
requires = up.REQUIRES
except:
pass
p_cflags.insert(0, '-I.')
if name is None:
@@ -1086,6 +1097,11 @@ def build_plugin(path, uc, cflags, ldflags, libs, name = None):
except:
pass
try:
p_cflags.remove('-pie')
except:
pass
#for ofile in up.OBJ_LIST:
# gcc_list.insert(0,ofile)
@@ -1098,6 +1114,17 @@ def build_plugin(path, uc, cflags, ldflags, libs, name = None):
print("*** unable to build %s plugin ***" % name)
sys.exit(1)
try:
if requires:
f = open('.uwsgi_plugin_section', 'w')
for rp in requires:
f.write("requires=%s\n" % rp)
f.close()
os.system("objcopy %s.so --add-section uwsgi=.uwsgi_plugin_section %s.so" % (plugin_dest, plugin_dest))
os.unlink('.uwsgi_plugin_section')
except:
pass
print("*** %s plugin built and available in %s ***" % (name, plugin_dest + '.so'))
if __name__ == "__main__":