mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-10-04 07:58:33 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c5dc74b021 | ||
|
|
97137488f9 | ||
|
|
8cff01a1c7 | ||
|
|
3c8452313e | ||
|
|
322b130fc9 | ||
|
|
6c001c2fda | ||
|
|
5d4050817d | ||
|
|
0a1147a4f6 | ||
|
|
f92cb0e6a2 | ||
|
|
8b79c76261 | ||
|
|
0015b51262 | ||
|
|
b1dc8cd9c6 | ||
|
|
3269e2d288 | ||
|
|
71a9747d96 | ||
|
|
455d041ce4 | ||
|
|
88205f8dd1 | ||
|
|
7d9dadddf6 | ||
|
|
57e78dd76f | ||
|
|
7d0d2b33f4 | ||
|
|
03c739102d | ||
|
|
392bb8d7b1 | ||
|
|
505b35be29 | ||
|
|
fdb002ad9b | ||
|
|
557c46373f | ||
|
|
3cde8be505 | ||
|
|
c2cf1f3a62 | ||
|
|
956cb1f13c |
@@ -53,3 +53,4 @@ e1568fd16b7b586cc72deb4dfccbbe64ae0b84df 1.1
|
||||
29de0fb320bc0a1ce84972f9f249360e315d3c10 1.2-rc2
|
||||
c3cdecbf2bac591336baddd14f9ac22e18d6e200 1.2
|
||||
14524da00a8b382dffb1d16a0e70cdf4a946f369 1.3-rc2
|
||||
883b946db9038372cb1d76cb1d328a3093140e7b 1.3-rc3
|
||||
|
||||
+1
-1
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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 = §ions[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 = §ions[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 = §ions[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 = §ions[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
@@ -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
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
@@ -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']
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -4,4 +4,6 @@ CFLAGS = []
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
|
||||
REQUIRES = ['corerouter']
|
||||
|
||||
GCC_LIST = ['fastrouter', 'fr_events']
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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 = []
|
||||
|
||||
|
||||
@@ -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
@@ -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;
|
||||
|
||||
@@ -4,4 +4,6 @@ CFLAGS = []
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
|
||||
REQUIRES = ['corerouter']
|
||||
|
||||
GCC_LIST = ['http']
|
||||
|
||||
@@ -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,
|
||||
|
||||
};
|
||||
|
||||
@@ -0,0 +1,6 @@
|
||||
NAME='mongodblog'
|
||||
|
||||
CFLAGS = []
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
GCC_LIST = ['mongodblog_plugin']
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
|
||||
@@ -72,6 +72,9 @@ struct uwsgi_rack {
|
||||
char *gemset;
|
||||
|
||||
int rb_threads;
|
||||
#ifdef RUBY19
|
||||
int patch_bodyproxy;
|
||||
#endif
|
||||
|
||||
};
|
||||
|
||||
|
||||
@@ -4,4 +4,6 @@ CFLAGS = []
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
|
||||
REQUIRES = ['corerouter']
|
||||
|
||||
GCC_LIST = ['rawrouter', 'rr_events']
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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,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
@@ -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__":
|
||||
|
||||
Reference in New Issue
Block a user