Compare commits

...
27 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
36 changed files with 1015 additions and 349 deletions
+1
View File
@@ -53,3 +53,4 @@ e1568fd16b7b586cc72deb4dfccbbe64ae0b84df 1.1
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
+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;
}
+28
View File
@@ -252,3 +252,31 @@ void uwsgi_setup_workers() {
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;
}
+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);
+113 -5
View File
@@ -1193,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;
@@ -1203,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;
@@ -1212,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;
@@ -2149,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) {
@@ -2160,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);
@@ -2531,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);
@@ -4738,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
+43 -58
View File
@@ -480,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},
@@ -660,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 = '/';
@@ -1203,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()");
}
@@ -1522,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 = {
@@ -2026,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;
@@ -2064,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()");
@@ -2386,30 +2385,7 @@ next:
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) {
@@ -2421,9 +2397,12 @@ next:
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;
@@ -3320,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, ""), ",");
+216 -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,153 +79,198 @@ 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 || uwsgi.cleaning) {
for (i = 0; i < uwsgi.numproc; i++) {
u_carbon.current_busyness_values[i] = uwsgi.workers[i+1].running_time - u_carbon.last_busyness_values[i];
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);
}
}
+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
@@ -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);
}
+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;
+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;
+5 -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;
+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']
+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;
}
}
+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);
}
+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;
}
+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);
+21 -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)\
@@ -2179,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);
@@ -2240,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, ...);
@@ -2321,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);
@@ -2538,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 {
@@ -2584,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 *);
@@ -2610,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 *);
@@ -2923,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 *);
@@ -3147,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;
@@ -3247,6 +3238,15 @@ 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
int uwsgi_init(int, char **, char **);
+24 -1
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '1.3-rc3'
uwsgi_version = '1.3-rc4'
import os
import re
@@ -1020,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[:]
@@ -1027,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:
@@ -1090,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)
@@ -1102,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__":