Compare commits

...
50 Commits
Author SHA1 Message Date
Roberto De Ioris 593aba2b19 improved mod_proxy_uwsgi 2012-12-30 11:05:20 +01:00
Roberto De Ioris 47c9e5e3f0 fixed /proc check 2012-12-30 09:10:10 +01:00
Roberto De Ioris bca8bc0c4d fixed socket naming of fd 0 2012-12-30 09:07:38 +01:00
Roberto De Ioris a67dfdd542 fixed a https/corerouter corner-case bug 2012-12-21 11:32:33 +01:00
Roberto De Ioris d1bbe7e3f9 set HTTPS to on when behind stud 2012-12-21 09:01:10 +01:00
Roberto De Ioris 442c982aff backported --http-stud-prefix 2012-12-19 18:41:48 +01:00
Roberto De Ioris a2f50ccf2e report pos in failed uploads 2012-12-16 12:58:12 +01:00
Roberto De Ioris b1c5fef0dc fixed computation of post data 2012-12-16 12:18:13 +01:00
Roberto De Ioris f6d8ac52f4 applied more debug to https 2012-12-16 10:59:56 +01:00
Roberto De Ioris 94e4fcbfe1 added additional uwgsi-gevent error report 2012-12-15 20:36:26 +01:00
Roberto De Ioris c9e7b0f48e fixed wrong typecasting in yaml and fixed subscription system on 32 bit 2012-12-15 14:32:25 +01:00
Roberto De Ioris 2f932d8ca0 backported http internal router fixes 2012-12-13 15:34:47 +01:00
Roberto De Ioris c517320f7a prepare for next release (1.4.4) 2012-12-13 11:07:55 +01:00
Roberto De Ioris b0dae66264 applied 0001-allows-embedding-.cc-and-.m-files.patch 2012-12-13 11:07:37 +01:00
Roberto De Ioris 0d87ed7298 fix dup2() usage 2012-12-12 15:23:04 +01:00
Roberto De Ioris fd3291b11f uWSGI 1.4.3 2012-12-10 15:26:09 +01:00
Roberto De Ioris e5ac5d6352 added emperor:// to supported scheme 2012-12-10 14:16:13 +01:00
Roberto De Ioris 403bd604e2 backported fix for bug #66 2012-12-10 14:14:01 +01:00
Roberto De Ioris 045d8947ae applied 0001-missing-new-line-in-carbon-plugin-logs.patch 2012-12-10 14:12:29 +01:00
Roberto De Ioris a9f2c71863 applied 0001-fallback-to-REMOTE_ADDR-if-HTTP_X_FORWARDED_FOR-is-m.patch 2012-12-10 14:11:42 +01:00
Roberto De Ioris b504a1ebe8 applied 0001-fixed-log-format-in-cheaper_busyness.patch 2012-12-10 14:10:28 +01:00
Roberto De Ioris f05ce0d2fc applied 0001-improved-stdin-management-in-cron.patch 2012-12-10 14:08:55 +01:00
Roberto De Ioris a2f3c716b1 applied 0001-added-offload-thread-alias.patch 2012-12-10 14:06:39 +01:00
Roberto De Ioris 2d7ca36f60 preparing for 1.4.3 2012-12-10 14:05:30 +01:00
unbit 5370e535dd core: make smart-attach-daemon smart enough
Avoid call storming of pidfile check that leads to multiple daemons
running instead of one. Fix #65.
2012-11-28 20:52:59 +01:00
unbit 092767c29f Merge pull request #63 from xrmx/cpucount14
uwsgiconfig: permit to override CPUCOUNT with environment variable
2012-11-27 10:32:45 -08:00
Riccardo Magliocchetti 4e119e6838 uwsgiconfig: permit to override CPUCOUNT with environment variable
CPUCOUNT=X make

will force X cores instead of the autodetected ones
2012-11-27 19:25:49 +01:00
Roberto De Ioris 87ef444ea6 uWSGI 1.4.2 2012-11-24 11:27:06 +01:00
Roberto De Ioris 714be76f9b improved QUERY_STRING handling in router_rewrite, added --php-var 2012-11-24 11:26:25 +01:00
Roberto De Ioris af0f0da8c2 added --php-app-qs 2012-11-24 10:11:56 +01:00
Roberto De Ioris 8407372801 fixed PATH_TRANSLATED in php and improved offloading of bigger files 2012-11-24 10:11:12 +01:00
Roberto De Ioris d19c4e84fc applied 0001-fixed-corner-case-ssl-body-reading.patch 2012-11-24 10:10:08 +01:00
Roberto De Ioris 3342852a9a backported 0001-do-not-crash-on-non-existent-config-files.patch 2012-11-21 07:49:13 +01:00
Roberto De Ioris 586a02d849 backported fixes from 1.5 2012-11-21 07:42:32 +01:00
Roberto De Ioris 188dad744a uWSGI 1.4.1 2012-11-13 04:53:32 +01:00
Roberto De Ioris 6cf3bbea69 backported fixed for offloading, static maps and non-blocking connect 2012-11-13 04:41:14 +01:00
Roberto De Ioris 04736d371e backported 0001-corerouter-NOT-courerouter.patch 2012-11-12 13:33:14 +01:00
Roberto De Ioris cc9151f89e 'wsgi' option is an alias for 'module' not 'wsgi-file' 2012-11-12 11:57:55 +01:00
Roberto De Ioris e13734ae6c uWSGI 1.4 2012-11-12 10:45:05 +01:00
Roberto De Ioris cf62b03e33 endiness fixes and added support for not waiting for thread cancellation 2012-11-12 10:36:18 +01:00
Roberto De Ioris b2f292c2da some Changelong improvement 2012-11-12 09:46:47 +01:00
Roberto De Ioris b96b6a24a6 support for rack official commit afc9b0313c 2012-11-12 09:40:38 +01:00
Roberto De Ioris b2a1237781 added --undeferred-shared-socket 2012-11-11 19:41:00 +01:00
Roberto De Ioris a48af2e226 workaround missing IP_FREEBIND 2012-11-11 19:07:28 +01:00
Roberto De Ioris 33446c3179 better cap management 2012-11-11 18:59:05 +01:00
Roberto De Ioris 242e01e37a added 'wsgi' alias option 2012-11-11 17:32:18 +01:00
Roberto De Ioris 41eba2b486 report static nodes in corerouters stats 2012-11-10 11:51:19 +01:00
Roberto De Ioris 24763fdbd2 fixed xclient 2012-11-10 11:23:39 +01:00
Roberto De Ioris 951cbf3eba implemented retries and fallback on the rawrouter 2012-11-10 10:25:40 +01:00
Roberto De Ioris 8759b66206 ported rawrouter to the new api with xclient support 2012-11-10 09:33:33 +01:00
42 changed files with 1046 additions and 401 deletions
+12 -1
View File
@@ -1,5 +1,14 @@
*** current ***
* 1.4.1
- fixed typos in corerouter plugins
- fixed offloading when the number of threads is higher than 1
- fixed static_maps for non-existent paths
- fixed uwsgi_connect() on modern Linux systems to reset the socket to blocking mode
*** november 2012 ***
* 1.4
- gevent improvements
@@ -9,11 +18,13 @@
- improved systemd support
- log filtering and routing
- improved tracebacker
- offload transfer for static files
- offload transfer for static files, and network transfers
- matheval support
- plugins can be written in Obj-C
- smart attach daemon
- added support for PEP 405 virtualenvs
- rawrouter with xclient support
- internal routing plugin for cache
*** semptember 2012 ***
+18 -2
View File
@@ -103,8 +103,16 @@ static int uwsgi_send_headers(request_rec *r, proxy_conn_rec *conn)
const char *script_name = apr_table_get(r->subprocess_env, "SCRIPT_NAME");
const char *path_info = apr_table_get(r->subprocess_env, "PATH_INFO");
if (script_name && path_info) {
apr_table_set(r->subprocess_env, "SCRIPT_NAME", apr_pstrndup(r->pool, script_name, strlen(script_name)-strlen(path_info)));
if (strcmp(path_info, "/")) {
apr_table_set(r->subprocess_env, "SCRIPT_NAME", apr_pstrndup(r->pool, script_name, strlen(script_name)-strlen(path_info)));
}
else {
if (!strcmp(script_name, "/")) {
apr_table_set(r->subprocess_env, "SCRIPT_NAME", "");
}
}
}
env_table = apr_table_elts(r->subprocess_env);
@@ -326,7 +334,15 @@ static int uwsgi_handler(request_rec *r, proxy_worker *worker,
}
// ADD PATH_INFO
apr_table_add(r->subprocess_env, "PATH_INFO", url+strlen(worker->name));
size_t w_len = strlen(worker->name);
char *u_path_info = r->filename + 6 + w_len;
ap_log_rerror(APLOG_MARK, APLOG_ERR, 0, r,
"URL %s: %s %s %s", url, worker->name, r->filename, u_path_info);
int delta = 0;
if (u_path_info[0] != '/') {
delta = 1;
}
apr_table_add(r->subprocess_env, "PATH_INFO", url+w_len-delta);
/* Create space for state information */
+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, mongodblog, router_rewrite, router_http, logfile, router_cache
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, logfile, router_cache, rawrouter
as_shared_library = false
locking = auto
+10
View File
@@ -114,3 +114,13 @@ int uwsgi_buffer_send(struct uwsgi_buffer *ub, int fd) {
return 0;
}
int uwsgi_buffer_num64(struct uwsgi_buffer *ub, int64_t num) {
char buf[sizeof(UMAX64_STR)+1];
int ret = snprintf(buf, sizeof(UMAX64_STR)+1, "%lld", (long long) num);
if (ret <= 0 || ret > (int) (sizeof(UMAX64_STR)+1)) {
return -1;
}
return uwsgi_buffer_append(ub, buf, ret);
}
+2 -2
View File
@@ -368,7 +368,7 @@ void manage_cluster_message(char *cluster_opt_buf, int cluster_opt_size) {
memset(&nucn, 0, sizeof(struct uwsgi_cluster_node));
#ifdef __BIG_ENDIAN__
uwsgi.workers[0].cores[0].req.uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize);
uwsgi.workers[0].cores[0].req.uh.pktsize = uwsgi_swap16(uwsgi.workers[0].cores[0].req.uh.pktsize);
#endif
uwsgi_hooked_parse(uwsgi.workers[0].cores[0].req.buffer, uwsgi.workers[0].cores[0].req.uh.pktsize, manage_cluster_announce, &nucn);
if (nucn.name[0] != 0) {
@@ -377,7 +377,7 @@ void manage_cluster_message(char *cluster_opt_buf, int cluster_opt_size) {
break;
case 96:
#ifdef __BIG_ENDIAN__
uwsgi.workers[0].cores[0].req.uh.pktsize = uwsgi_swap16(uwsgi.wsgi_requests[0]->uh.pktsize);
uwsgi.workers[0].cores[0].req.uh.pktsize = uwsgi_swap16(uwsgi.workers[0].cores[0].req.uh.pktsize);
#endif
uwsgi_log_verbose("%.*s\n", uwsgi.workers[0].cores[0].req.uh.pktsize, uwsgi.workers[0].cores[0].req.buffer);
break;
+14 -3
View File
@@ -30,6 +30,16 @@ extern struct uwsgi_server uwsgi;
*/
void uwsgi_daemons_smart_check() {
static time_t last_run = 0;
time_t now = uwsgi_now();
if (now - last_run <= 0) {
return;
}
last_run = now;
struct uwsgi_daemon *ud = uwsgi.daemons;
while (ud) {
if (ud->pidfile) {
@@ -41,14 +51,14 @@ void uwsgi_daemons_smart_check() {
}
else {
ud->pidfile_checks++;
if (ud->pidfile_checks >= (uint64_t) ud->freq) {
if (ud->pidfile_checks >= (unsigned int) ud->freq) {
uwsgi_log("[uwsgi-daemons] found changed pidfile for \"%s\" (old_pid: %d new_pid: %d)\n", ud->command, (int) ud->pid, (int) checked_pid);
uwsgi_spawn_daemon(ud);
}
}
}
else if (checked_pid != ud->pid) {
uwsgi_log("[uwsgi-daemons] found changed pidfile for \"%s\" (old_pid: %d new_pid: %d)\n", ud->command, (int) ud->pid, (int) checked_pid);
uwsgi_log("[uwsgi-daemons] found changed pid for \"%s\" (old_pid: %d new_pid: %d)\n", ud->command, (int) ud->pid, (int) checked_pid);
ud->pid = checked_pid;
}
// all ok, pidfile and process found
@@ -183,6 +193,7 @@ void uwsgi_spawn_daemon(struct uwsgi_daemon *ud) {
else {
// close uwsgi sockets
uwsgi_close_all_sockets();
uwsgi_close_all_fds();
if (ud->daemonize) {
/* refork... */
@@ -204,7 +215,7 @@ void uwsgi_spawn_daemon(struct uwsgi_daemon *ud) {
exit(1);
}
if (devnull != 0) {
if (dup2(devnull, 0)) {
if (dup2(devnull, 0) < 0) {
uwsgi_error("dup2()");
exit(1);
}
+3 -1
View File
@@ -709,7 +709,7 @@ void emperor_add(struct uwsgi_emperor_scanner *ues, char *name, time_t born, cha
exit(1);
}
if (stdin_fd != 0) {
if (dup2(stdin_fd, 0)) {
if (dup2(stdin_fd, 0) < 0) {
uwsgi_error("dup2()");
exit(1);
}
@@ -784,6 +784,8 @@ void uwsgi_imperial_monitor_directory_init(struct uwsgi_emperor_scanner *ues) {
exit(1);
}
ues->arg = uwsgi.emperor_absolute_dir;
}
struct uwsgi_imperial_monitor *imperial_monitor_get_by_id(char *scheme) {
+1
View File
@@ -92,6 +92,7 @@ void uwsgi_ini_config(char *file, char *magic_table[]) {
char *section_asked = "uwsgi";
char *colon;
if (uwsgi_check_scheme(file)) {
colon = uwsgi_get_last_char(file, '/');
colon = uwsgi_get_last_char(colon, ':');
+20 -19
View File
@@ -107,7 +107,7 @@ error:
return -1;
}
static void uwsgi_offload_close(struct uwsgi_offload_request *uor) {
static void uwsgi_offload_close(struct uwsgi_thread *ut, struct uwsgi_offload_request *uor) {
// close the socket and the file descriptor
close(uor->s);
close(uor->fd);
@@ -115,12 +115,12 @@ static void uwsgi_offload_close(struct uwsgi_offload_request *uor) {
struct uwsgi_offload_request *prev = uor->prev;
struct uwsgi_offload_request *next = uor->next;
if (uor == uwsgi.offload_requests_head) {
uwsgi.offload_requests_head = next;
if (uor == ut->offload_requests_head) {
ut->offload_requests_head = next;
}
if (uor == uwsgi.offload_requests_tail) {
uwsgi.offload_requests_tail = prev;
if (uor == ut->offload_requests_tail) {
ut->offload_requests_tail = prev;
}
if (prev) {
@@ -142,22 +142,22 @@ static void uwsgi_offload_close(struct uwsgi_offload_request *uor) {
free(uor);
}
static void uwsgi_offload_append(struct uwsgi_offload_request *uor) {
static void uwsgi_offload_append(struct uwsgi_thread *ut, struct uwsgi_offload_request *uor) {
if (!uwsgi.offload_requests_head) {
uwsgi.offload_requests_head = uor;
if (!ut->offload_requests_head) {
ut->offload_requests_head = uor;
}
if (uwsgi.offload_requests_tail) {
uwsgi.offload_requests_tail->next = uor;
uor->prev = uwsgi.offload_requests_tail;
if (ut->offload_requests_tail) {
ut->offload_requests_tail->next = uor;
uor->prev = ut->offload_requests_tail;
}
uwsgi.offload_requests_tail = uor;
ut->offload_requests_tail = uor;
}
static struct uwsgi_offload_request *uwsgi_offload_get_by_fd(int s) {
struct uwsgi_offload_request *uor = uwsgi.offload_requests_head;
static struct uwsgi_offload_request *uwsgi_offload_get_by_fd(struct uwsgi_thread *ut, int s) {
struct uwsgi_offload_request *uor = ut->offload_requests_head;
while (uor) {
if (uor->s == s || uor->fd == s) {
return uor;
@@ -187,20 +187,20 @@ static void uwsgi_offload_loop(struct uwsgi_thread *ut) {
}
// start monitoring socket for write
if (uor->func(ut, uor, -1)) {
uwsgi_offload_close(uor);
uwsgi_offload_close(ut, uor);
continue;
}
uwsgi_offload_append(uor);
uwsgi_offload_append(ut, uor);
continue;
}
// get the task from the interesting fd
struct uwsgi_offload_request *uor = uwsgi_offload_get_by_fd(interesting_fd);
struct uwsgi_offload_request *uor = uwsgi_offload_get_by_fd(ut, interesting_fd);
if (!uor)
continue;
// run the hook
if (uor->func(ut, uor, interesting_fd)) {
uwsgi_offload_close(uor);
uwsgi_offload_close(ut, uor);
}
}
}
@@ -234,6 +234,7 @@ static int uwsgi_offload_sendfile_transfer(struct uwsgi_thread *ut, struct uwsgi
if (uor->written >= uor->len) {
return -1;
}
return 0;
}
else if (len < 0) {
uwsgi_offload_retry
@@ -330,7 +331,7 @@ static int uwsgi_offload_net_transfer(struct uwsgi_thread *ut, struct uwsgi_offl
}
}
else if (fd == uor->s) {
rlen = read(uor->fd, uor->buf, 4096);
rlen = read(uor->s, uor->buf, 4096);
if (rlen > 0) {
uor->to_write = rlen;
uor->pos = 0;
+18 -2
View File
@@ -784,7 +784,8 @@ int uwsgi_parse_vars(struct wsgi_request *wsgi_req) {
wsgi_req->method = ptrbuf;
wsgi_req->method_len = strsize;
}
else if (!uwsgi.log_x_forwarded_for && !uwsgi_strncmp("REMOTE_ADDR", 11, wsgi_req->hvec[wsgi_req->var_cnt].iov_base, wsgi_req->hvec[wsgi_req->var_cnt].iov_len)) {
else if ((!uwsgi.log_x_forwarded_for || uwsgi_strncmp("HTTP_X_FORWARDED_FOR", 20, wsgi_req->hvec[wsgi_req->var_cnt].iov_base, wsgi_req->hvec[wsgi_req->var_cnt].iov_len))
&& !uwsgi_strncmp("REMOTE_ADDR", 11, wsgi_req->hvec[wsgi_req->var_cnt].iov_base, wsgi_req->hvec[wsgi_req->var_cnt].iov_len)) {
wsgi_req->remote_addr = ptrbuf;
wsgi_req->remote_addr_len = strsize;
}
@@ -1082,7 +1083,7 @@ nextcs:
udd = uwsgi.static_maps;
while (udd) {
#ifdef UWSGI_DEBUG
uwsgi_log("checking for %.*s <-> %.*s\n", wsgi_req->path_info_len, wsgi_req->path_info, udd->keylen, udd->key);
uwsgi_log("checking for %.*s <-> %.*s %.*s\n", (int)wsgi_req->path_info_len, wsgi_req->path_info, (int)udd->keylen, udd->key, (int) udd->vallen, udd->value);
#endif
if (udd->status == 0) {
#ifdef UWSGI_THREADING
@@ -1092,6 +1093,7 @@ nextcs:
char *real_docroot = uwsgi_malloc(PATH_MAX + 1);
if (!realpath(udd->value, real_docroot)) {
free(real_docroot);
real_docroot = NULL;
udd->value = NULL;
}
#ifdef UWSGI_THREADING
@@ -1128,6 +1130,7 @@ nextsm:
char *real_docroot = uwsgi_malloc(PATH_MAX + 1);
if (!realpath(udd->value, real_docroot)) {
free(real_docroot);
real_docroot = NULL;
udd->value = NULL;
}
#ifdef UWSGI_THREADING
@@ -1662,18 +1665,31 @@ char *uwsgi_req_append(struct wsgi_request *wsgi_req, char *key, uint16_t keylen
return NULL;
}
if (wsgi_req->var_cnt >= uwsgi.vec_size - (4 + 2)) {
uwsgi_log("max vec size reached. skip this header.\n");
return NULL;
}
char *ptr = wsgi_req->buffer + wsgi_req->uh.pktsize;
*ptr++ = (uint8_t) (keylen & 0xff);
*ptr++ = (uint8_t) ((keylen >> 8) & 0xff);
memcpy(ptr, key, keylen);
wsgi_req->hvec[wsgi_req->var_cnt].iov_base = ptr;
wsgi_req->hvec[wsgi_req->var_cnt].iov_len = keylen;
wsgi_req->var_cnt++;
ptr += keylen;
*ptr++ = (uint8_t) (vallen & 0xff);
*ptr++ = (uint8_t) ((vallen >> 8) & 0xff);
memcpy(ptr, val, vallen);
wsgi_req->hvec[wsgi_req->var_cnt].iov_base = ptr;
wsgi_req->hvec[wsgi_req->var_cnt].iov_len = vallen;
wsgi_req->var_cnt++;
wsgi_req->uh.pktsize += (2 + keylen + 2 + vallen);
+18 -8
View File
@@ -564,14 +564,15 @@ int bind_to_tcp(char *socket_name, int listen_queue, char *tcp_port) {
}
#ifdef __linux__
#ifdef IP_FREEBIND
#ifndef IP_FREEBIND
#define IP_FREEBIND 15
#endif
if (uwsgi.freebind) {
if (setsockopt(serverfd, SOL_IP, IP_FREEBIND, (const void *) &uwsgi.freebind, sizeof(int)) < 0) {
uwsgi_error("IP_FREEBIND setsockopt()");
uwsgi_nuclear_blast();
}
}
#endif
#endif
if (uwsgi.reuse_port) {
@@ -629,7 +630,7 @@ int bind_to_tcp(char *socket_name, int listen_queue, char *tcp_port) {
}
#ifdef __linux__
long somaxconn = uwsgi_num_from_file("/proc/sys/net/core/somaxconn");
long somaxconn = uwsgi_num_from_file("/proc/sys/net/core/somaxconn", 1);
if (somaxconn > 0 && uwsgi.listen_queue > somaxconn) {
uwsgi_log("Listen queue size is greater than the system max net.core.somaxconn (%i).\n", somaxconn);
uwsgi_nuclear_blast();
@@ -717,7 +718,7 @@ int timed_connect(struct pollfd *fdpoll, const struct sockaddr *addr, int addr_s
#if defined(__linux__) && defined(SOCK_NONBLOCK) && !defined(OBSOLETE_LINUX_KERNEL)
// hmm, nothing to do, as we are already non-blocking
uwsgi_socket_b(fdpoll->fd);
#else
/* re-set blocking socket */
arg &= (~O_NONBLOCK);
@@ -1483,6 +1484,10 @@ void uwsgi_setup_shared_sockets() {
while (shared_sock) {
if (!uwsgi.is_a_reload) {
char *tcp_port = strrchr(shared_sock->name, ':');
int current_defer_accept = uwsgi.no_defer_accept;
if (shared_sock->no_defer) {
uwsgi.no_defer_accept = 1;
}
if (tcp_port == NULL) {
shared_sock->fd = bind_to_unix(shared_sock->name, uwsgi.listen_queue, uwsgi.chmod_socket, uwsgi.abstract_socket);
shared_sock->family = AF_UNIX;
@@ -1513,6 +1518,11 @@ void uwsgi_setup_shared_sockets() {
uwsgi_log("unable to create shared socket on: %s\n", shared_sock->name);
exit(1);
}
if (shared_sock->no_defer) {
uwsgi.no_defer_accept = current_defer_accept;
}
}
else {
for (i = 3; i < (int) uwsgi.max_fd; i++) {
@@ -1601,7 +1611,7 @@ void uwsgi_map_sockets() {
exit(1);
}
if (fd != uwsgi_sock->fd) {
if (dup2(fd, uwsgi_sock->fd)) {
if (dup2(fd, uwsgi_sock->fd) < 0) {
uwsgi_error("dup2()");
exit(1);
}
@@ -1692,14 +1702,14 @@ void uwsgi_bind_sockets() {
gsa.sa = (struct sockaddr *) &usa;
if (!uwsgi.skip_zero && !getsockname(0, gsa.sa, &socket_type_len)) {
if (gsa.sa->sa_family == AF_UNIX) {
uwsgi_sock = uwsgi_new_socket(usa.sa_un.sun_path);
uwsgi_sock = uwsgi_new_socket(uwsgi_getsockname(0));
uwsgi_sock->family = AF_UNIX;
uwsgi_sock->fd = 0;
uwsgi_sock->bound = 1;
uwsgi_log("uwsgi socket %d inherited UNIX address %s fd 0\n", uwsgi_get_socket_num(uwsgi_sock), uwsgi_sock->name);
}
else {
uwsgi_sock = uwsgi_new_socket(uwsgi_concat2("::", ""));
uwsgi_sock = uwsgi_new_socket(uwsgi_getsockname(0));
uwsgi_sock->family = AF_INET;
uwsgi_sock->fd = 0;
uwsgi_sock->bound = 1;
@@ -1713,7 +1723,7 @@ void uwsgi_bind_sockets() {
exit(1);
}
if (fd != 0) {
if (dup2(fd, 0)) {
if (dup2(fd, 0) < 0) {
uwsgi_error("dup2()");
exit(1);
}
+73 -12
View File
@@ -385,10 +385,17 @@ void uwsgi_as_root() {
cap_free(caps);
#ifdef __linux__
#ifdef SECBIT_KEEP_CAPS
if (prctl(SECBIT_KEEP_CAPS, 1, 0, 0, 0) < 0) {
uwsgi_error("prctl()");
exit(1);
}
#else
if (prctl(PR_SET_KEEPCAPS, 1, 0, 0, 0) < 0) {
uwsgi_error("prctl()");
exit(1);
}
#endif
#endif
}
#endif
@@ -841,17 +848,19 @@ void uwsgi_linux_ksm_map(void) {
#endif
#ifdef __linux__
long uwsgi_num_from_file(char *filename) {
long uwsgi_num_from_file(char *filename, int quiet) {
char buf[16];
ssize_t len;
int fd = open(filename, O_RDONLY);
if (fd < 0) {
uwsgi_error_open(filename);
if (!quiet)
uwsgi_error_open(filename);
return -1L;
}
len = read(fd, buf, sizeof(buf));
if (len == 0) {
uwsgi_log("read error %s\n", filename);
if (!quiet)
uwsgi_log("read error %s\n", filename);
close(fd);
return -1L;
}
@@ -2684,15 +2693,10 @@ void init_magic_table(char *magic_table[]) {
}
char *uwsgi_get_last_char(char *what, char c) {
int i, j = 0;
int i;
char *ptr = NULL;
if (!strncmp("http://", what, 7))
j = 7;
if (!strncmp("emperor://", what, 10))
j = 10;
for (i = j; i < (int) strlen(what); i++) {
for (i = 0; i < (int) strlen(what); i++) {
if (what[i] == c) {
ptr = what + i;
}
@@ -2892,7 +2896,7 @@ size_t uwsgi_str_num(char *str, int len) {
int i;
size_t num = 0;
size_t delta = pow(10, len);
uint64_t delta = pow(10, len);
for (i = 0; i < len; i++) {
delta = delta / 10;
@@ -3101,6 +3105,7 @@ struct uwsgi_string_list *uwsgi_string_new_list(struct uwsgi_string_list **list,
}
uwsgi_string->next = NULL;
uwsgi_string->custom = 0;
uwsgi_string->custom2 = 0;
return uwsgi_string;
}
@@ -3271,10 +3276,46 @@ pid_t uwsgi_run_command(char *command, int *stdin_fd, int stdout_fd) {
}
uwsgi_close_all_sockets();
//uwsgi_close_all_fds();
int i;
for (i = 3; i < (int) uwsgi.max_fd; i++) {
if (stdin_fd) {
if (i == stdin_fd[0] || i == stdin_fd[1]) {
continue;
}
}
if (stdout_fd > -1) {
if (i == stdout_fd) {
continue;
}
}
#ifdef __APPLE__
fcntl(i, F_SETFD, FD_CLOEXEC);
#else
close(i);
#endif
}
if (stdin_fd) {
close(stdin_fd[1]);
}
else {
if (!uwsgi_valid_fd(0)) {
int in_fd = open("/dev/null", O_RDONLY);
if (in_fd < 0) {
uwsgi_error_open("/dev/null");
}
else {
if (in_fd != 0) {
if (dup2(in_fd, 0) < 0) {
uwsgi_error("dup2()");
}
}
}
}
}
if (stdout_fd > -1 && stdout_fd != 1) {
if (dup2(stdout_fd, 1) < 0) {
@@ -4798,7 +4839,7 @@ static void *uwsgi_thread_run(void *arg) {
struct uwsgi_thread *uwsgi_thread_new(void (*func) (struct uwsgi_thread *)) {
struct uwsgi_thread *ut = uwsgi_malloc(sizeof(struct uwsgi_thread));
struct uwsgi_thread *ut = uwsgi_calloc(sizeof(struct uwsgi_thread));
#if defined(SOCK_SEQPACKET) && defined(__linux__)
if (socketpair(AF_UNIX, SOCK_SEQPACKET, 0, ut->pipe)) {
@@ -5036,3 +5077,23 @@ end:
free(symbol_name);
return ret;
}
int uwsgi_valid_fd(int fd) {
int ret = fcntl(fd, F_GETFL);
if (ret == 0) {
return 1;
}
return 0;
}
void uwsgi_close_all_fds(void) {
int i;
for (i = 3; i < (int) uwsgi.max_fd; i++) {
#ifdef __APPLE__
fcntl(i, F_SETFD, FD_CLOEXEC);
#else
close(i);
#endif
}
}
+58 -12
View File
@@ -42,6 +42,7 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"protocol", required_argument, 0, "force the specified protocol for default sockets", uwsgi_opt_set_str, &uwsgi.protocol, 0},
{"socket-protocol", required_argument, 0, "force the specified protocol for default sockets", uwsgi_opt_set_str, &uwsgi.protocol, 0},
{"shared-socket", required_argument, 0, "create a shared sacket for advanced jailing or ipc", uwsgi_opt_add_shared_socket, NULL, 0},
{"undeferred-shared-socket", required_argument, 0, "create a shared sacket for advanced jailing or ipc (undeferred mode)", uwsgi_opt_add_shared_socket, NULL, 0},
{"processes", required_argument, 'p', "spawn the specified number of workers/processes", uwsgi_opt_set_int, &uwsgi.numproc, 0},
{"workers", required_argument, 'p', "spawn the specified number of workers/processes", uwsgi_opt_set_int, &uwsgi.numproc, 0},
{"harakiri", required_argument, 't', "set harakiri timeout", uwsgi_opt_set_dyn, (void *) UWSGI_OPTION_HARAKIRI, 0},
@@ -124,6 +125,7 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"map-socket", required_argument, 0, "map sockets to specific workers", uwsgi_opt_add_string_list, &uwsgi.map_socket, 0},
#ifdef UWSGI_THREADING
{"enable-threads", no_argument, 'T', "enable threads", uwsgi_opt_true, &uwsgi.has_threads, 0},
{"no-threads-wait", no_argument, 0, "do not wait for threads cancellation on quit/reload", uwsgi_opt_true, &uwsgi.no_threads_wait, 0},
#endif
{"auto-procname", no_argument, 0, "automatically set processes name to something meaningful", uwsgi_opt_true, &uwsgi.auto_procname, 0},
@@ -469,6 +471,7 @@ static struct uwsgi_option uwsgi_base_options[] = {
#endif
{"offload-threads", required_argument, 0, "set the number of offload threads to spawn (per-worker, default 0)", uwsgi_opt_set_int, &uwsgi.offload_threads, 0},
{"offload-thread", required_argument, 0, "set the number of offload threads to spawn (per-worker, default 0)", uwsgi_opt_set_int, &uwsgi.offload_threads, 0},
{"file-serve-mode", required_argument, 0, "set static file serving mode", uwsgi_opt_fileserve_mode, NULL, UWSGI_OPT_MIME},
{"fileserve-mode", required_argument, 0, "set static file serving mode", uwsgi_opt_fileserve_mode, NULL, UWSGI_OPT_MIME},
@@ -675,6 +678,17 @@ void config_magic_table_fill(char *filename, char **magic_table) {
char *tmp = NULL;
char *fullname = filename;
magic_table['o'] = filename;
if (uwsgi_check_scheme(filename) || !strcmp(filename, "-")) {
return;
}
char *section = uwsgi_get_last_char(filename, ':');
if (section) {
*section = 0;
}
// we have a special case for symlinks
if (uwsgi_is_link(filename)) {
if (filename[0] != '/') {
@@ -684,19 +698,16 @@ void config_magic_table_fill(char *filename, char **magic_table) {
else {
fullname = uwsgi_expand_path(filename, strlen(filename), NULL);
if (fullname) {
char *minimal_name = uwsgi_malloc(strlen(fullname) + 1);
memcpy(minimal_name, fullname, strlen(fullname));
minimal_name[strlen(fullname)] = 0;
free(fullname);
fullname = minimal_name;
}
else {
fullname = filename;
if (!fullname) {
exit(1);
}
char *minimal_name = uwsgi_malloc(strlen(fullname) + 1);
memcpy(minimal_name, fullname, strlen(fullname));
minimal_name[strlen(fullname)] = 0;
free(fullname);
fullname = minimal_name;
}
magic_table['o'] = filename;
magic_table['p'] = fullname;
magic_table['s'] = uwsgi_get_last_char(fullname, '/') + 1;
magic_table['d'] = uwsgi_concat2n(magic_table['p'], magic_table['s'] - magic_table['p'], "", 0);
@@ -730,6 +741,11 @@ void config_magic_table_fill(char *filename, char **magic_table) {
magic_table['e'] = uwsgi_get_last_char(filename, '.') + 1;
if (uwsgi_get_last_char(magic_table['s'], '.'))
magic_table['n'] = uwsgi_concat2n(magic_table['s'], uwsgi_get_last_char(magic_table['s'], '.') - magic_table['s'], "", 0);
if (section) {
magic_table['x'] = section+1;
*section = ':';
}
}
int find_worker_id(pid_t pid) {
@@ -759,6 +775,9 @@ void warn_pipe() {
void wait_for_threads() {
int i, ret;
// on some platform thread cancellation is REALLY flaky
if (uwsgi.no_threads_wait) return;
int sudden_death = 0;
pthread_mutex_lock(&uwsgi.six_feet_under_lock);
@@ -2671,7 +2690,7 @@ next2:
}
if (uwsgi.sockets->fd != 0) {
if (dup2(uwsgi.sockets->fd, 0)) {
if (dup2(uwsgi.sockets->fd, 0) < 0) {
uwsgi_error("dup2()");
}
}
@@ -3275,6 +3294,30 @@ void uwsgi_opt_add_string_list(char *opt, char *value, void *list) {
uwsgi_string_new_list(ptr, value);
}
void uwsgi_opt_add_addr_list(char *opt, char *value, void *list) {
struct uwsgi_string_list **ptr = (struct uwsgi_string_list **) list;
int af = AF_INET;
#ifdef UWSGI_IPV6
void *ip = uwsgi_malloc(16);
if (strchr(value, ':')) {
af = AF_INET6;
}
#else
void *ip = uwsgi_malloc(4);
#endif
if (inet_pton(af, value, ip) <= 0) {
uwsgi_log("%s: invalid address\n", opt);
uwsgi_error("uwsgi_opt_add_addr_list()");
exit(1);
}
struct uwsgi_string_list *usl = uwsgi_string_new_list(ptr, ip);
usl->custom = af;
usl->custom_ptr = value;
}
#ifdef UWSGI_PCRE
void uwsgi_opt_add_regexp_list(char *opt, char *value, void *list) {
struct uwsgi_regexp_list **ptr = (struct uwsgi_regexp_list **) list;
@@ -3294,7 +3337,10 @@ void uwsgi_opt_add_regexp_custom_list(char *opt, char *value, void *list) {
#endif
void uwsgi_opt_add_shared_socket(char *opt, char *value, void *protocol) {
uwsgi_new_shared_socket(generate_socket_name(value));
struct uwsgi_socket *us = uwsgi_new_shared_socket(generate_socket_name(value));
if (!strcmp(opt, "undeferred-shared-socket")) {
us->no_defer = 1;
}
}
void uwsgi_opt_add_socket(char *opt, char *value, void *protocol) {
+1 -1
View File
@@ -137,7 +137,7 @@ void uwsgi_yaml_config(char *file, char *magic_table[]) {
exit(1);
}
yaml_parser_set_input_string(&parser, (const unsigned char *) yaml, (size_t) len - 1);
yaml_parser_set_input_string(&parser, (unsigned char *) yaml, (size_t) len - 1);
while (parsing) {
if (!yaml_parser_scan(&parser, &token)) {
+1 -1
View File
@@ -82,7 +82,7 @@ void carbon_post_init() {
// 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",
uwsgi_log("[carbon] carbon plugin started, %is frequency, %is timeout, max retries %i, retry delay %is\n",
u_carbon.freq, u_carbon.timeout, u_carbon.max_retries, u_carbon.retry_delay);
}
+9 -3
View File
@@ -838,9 +838,15 @@ clear2:
uwsgi_error("setenv()");
}
if (setenv("PATH_TRANSLATED", uwsgi_concat3n(docroot, docroot_len, path_info, pi_len, "", 0) , 1)) {
uwsgi_error("setenv()");
if (wsgi_req->document_root_len > 0) {
if (setenv("PATH_TRANSLATED", uwsgi_concat3n(wsgi_req->document_root, wsgi_req->document_root_len, path_info, pi_len, "", 0) , 1)) {
uwsgi_error("setenv()");
}
}
else {
if (setenv("PATH_TRANSLATED", uwsgi_concat3n(docroot, docroot_len, path_info, pi_len, "", 0) , 1)) {
uwsgi_error("setenv()");
}
}
}
+14 -14
View File
@@ -78,7 +78,7 @@ void set_next_cheap_time(void) {
// to have quicker recovery from big but short load spikes
// otherwise we might wait a lot before cheaping all emergency workers
if (uwsgi_cheaper_busyness_global.verbose)
uwsgi_log("[busyness] %d emergency worker(s) running, using %d seconds cheaper timer\n",
uwsgi_log("[busyness] %d emergency worker(s) running, using %llu seconds cheaper timer\n",
uwsgi_cheaper_busyness_global.emergency_workers, uwsgi.cheaper_overload*uwsgi_cheaper_busyness_global.backlog_multi);
uwsgi_cheaper_busyness_global.next_cheap = now + uwsgi.cheaper_overload*uwsgi_cheaper_busyness_global.backlog_multi*1000000;
} else {
@@ -95,7 +95,7 @@ void decrease_multi(void) {
// will decrease multiplier but only down to initial value
if (uwsgi_cheaper_busyness_global.cheap_multi > uwsgi_cheaper_busyness_global.min_multi) {
uwsgi_cheaper_busyness_global.cheap_multi--;
uwsgi_log("[busyness] decreasing cheaper multiplier to %d\n", uwsgi_cheaper_busyness_global.cheap_multi);
uwsgi_log("[busyness] decreasing cheaper multiplier to %llu\n", uwsgi_cheaper_busyness_global.cheap_multi);
}
}
@@ -157,7 +157,7 @@ int cheaper_busyness_algo(void) {
// store initial multiplier so we don't loose its initial value
uwsgi_cheaper_busyness_global.min_multi = uwsgi_cheaper_busyness_global.cheap_multi;
// since this is first run we will print current values
uwsgi_log("[busyness] settings: min=%d%%, max=%d%%, overload=%d, multiplier=%d, respawn penalty=%d\n",
uwsgi_log("[busyness] settings: min=%llu%%, max=%llu%%, overload=%llu, multiplier=%llu, respawn penalty=%llu\n",
uwsgi_cheaper_busyness_global.busyness_min, uwsgi_cheaper_busyness_global.busyness_max,
uwsgi.cheaper_overload, uwsgi_cheaper_busyness_global.cheap_multi, uwsgi_cheaper_busyness_global.penalty);
#ifdef __linux__
@@ -171,7 +171,7 @@ int cheaper_busyness_algo(void) {
if (uwsgi_cheaper_busyness_global.next_cheap == 0) set_next_cheap_time();
int64_t active_workers = 0;
int active_workers = 0;
uint64_t total_busyness = 0;
uint64_t avg_busyness = 0;
@@ -195,7 +195,7 @@ int cheaper_busyness_algo(void) {
if (percent > 100) percent = 100;
total_busyness += percent;
if (uwsgi_cheaper_busyness_global.verbose && active_workers > 1)
uwsgi_log("[busyness] worker nr %d %ds average busyness is at %d%%\n",
uwsgi_log("[busyness] worker nr %d %llus average busyness is at %llu%%\n",
i+1, uwsgi.cheaper_overload, percent);
}
uwsgi_cheaper_busyness_global.last_values[i] = uwsgi.workers[i+1].running_time;
@@ -204,7 +204,7 @@ int cheaper_busyness_algo(void) {
avg_busyness = (active_workers ? total_busyness / active_workers : 0);
if (uwsgi_cheaper_busyness_global.verbose)
uwsgi_log("[busyness] %ds average busyness of %d worker(s) is at %d%%\n",
uwsgi.cheaper_overload, active_workers, avg_busyness);
(int) uwsgi.cheaper_overload, (int) active_workers, (int) avg_busyness);
if (avg_busyness > uwsgi_cheaper_busyness_global.busyness_max) {
@@ -229,7 +229,7 @@ int cheaper_busyness_algo(void) {
// worker was cheaped and then spawned back in less than current multiplier*cheaper_overload seconds
// we will increase the multiplier so that next time worker will need to wait longer before being cheaped
uwsgi_cheaper_busyness_global.cheap_multi += uwsgi_cheaper_busyness_global.penalty;
uwsgi_log("[busyness] worker(s) respawned to fast, increasing chpeaper multiplier to %d (+%d)\n",
uwsgi_log("[busyness] worker(s) respawned to fast, increasing chpeaper multiplier to %llu (+%llu)\n",
uwsgi_cheaper_busyness_global.cheap_multi, uwsgi_cheaper_busyness_global.penalty);
} else {
decrease_multi();
@@ -237,10 +237,10 @@ int cheaper_busyness_algo(void) {
set_next_cheap_time();
uwsgi_log("[busyness] %ds average busyness is at %d%%, will spawn %d new worker(s)\n",
uwsgi_log("[busyness] %llus average busyness is at %llu%%, will spawn %d new worker(s)\n",
uwsgi.cheaper_overload, avg_busyness, decheaped);
} else {
uwsgi_log("[busyness] %ds average busyness is at %d%% but we already started maximum number of workers (%d)\n",
uwsgi_log("[busyness] %llus average busyness is at %llu%% but we already started maximum number of workers (%d)\n",
uwsgi.cheaper_overload, avg_busyness, uwsgi.numproc);
}
@@ -267,8 +267,8 @@ int cheaper_busyness_algo(void) {
if (uwsgi_cheaper_busyness_global.last_action == 2) decrease_multi();
set_next_cheap_time();
uwsgi_log("[busyness] %ds average busyness is at %d%%, cheap one of %d running workers\n",
uwsgi.cheaper_overload, avg_busyness, active_workers);
uwsgi_log("[busyness] %llus average busyness is at %llu%%, cheap one of %d running workers\n",
uwsgi.cheaper_overload, avg_busyness, (int) active_workers);
// store timestamp
uwsgi_cheaper_busyness_global.last_cheaped = uwsgi_micros();
@@ -280,7 +280,7 @@ int cheaper_busyness_algo(void) {
return -1;
} else if (uwsgi_cheaper_busyness_global.verbose)
uwsgi_log("[busyness] need to wait %d more second(s) to cheap worker\n", (uwsgi_cheaper_busyness_global.next_cheap - now)/1000000);
uwsgi_log("[busyness] need to wait %llu more second(s) to cheap worker\n", (uwsgi_cheaper_busyness_global.next_cheap - now)/1000000);
}
} else {
@@ -301,14 +301,14 @@ int cheaper_busyness_algo(void) {
// time needed to cheap them, than a lot min<busy<max when we do not reset timer
// and then another idle cycle than would trigger cheaping
if (uwsgi_cheaper_busyness_global.verbose)
uwsgi_log("[busyness] %ds average busyness is at %d%%, %d non-idle cycle(s), reseting cheaper timer\n",
uwsgi_log("[busyness] %llus average busyness is at %llu%%, %llu non-idle cycle(s), reseting cheaper timer\n",
uwsgi.cheaper_overload, avg_busyness, uwsgi_cheaper_busyness_global.tolerance_counter);
set_next_cheap_time();
} else {
// we had < 3 idle cycles in a row so we won't reset idle timer yet since this might be just short load spike
// but we need to add cheaper-overload seconds to the cheaper timer so this cycle isn't counted as idle
if (uwsgi_cheaper_busyness_global.verbose)
uwsgi_log("[busyness] %ds average busyness is at %d%%, %d non-idle cycle(s), adjusting cheaper timer\n",
uwsgi_log("[busyness] %llus average busyness is at %llu%%, %llu non-idle cycle(s), adjusting cheaper timer\n",
uwsgi.cheaper_overload, avg_busyness, uwsgi_cheaper_busyness_global.tolerance_counter);
uwsgi_cheaper_busyness_global.next_cheap += uwsgi.cheaper_overload*1000000;
}
+94 -57
View File
@@ -22,6 +22,13 @@ void uwsgi_opt_corerouter(char *opt, char *value, void *cr) {
ucr->has_sockets++;
}
void uwsgi_opt_undeferred_corerouter(char *opt, char *value, void *cr) {
struct uwsgi_corerouter *ucr = (struct uwsgi_corerouter *) cr;
struct uwsgi_gateway_socket *ugs = uwsgi_new_gateway_socket(value, ucr->name);
ugs->no_defer = 1;
ucr->has_sockets++;
}
void uwsgi_opt_corerouter_use_socket(char *opt, char *value, void *cr) {
struct uwsgi_corerouter *ucr = (struct uwsgi_corerouter *) cr;
ucr->use_socket = 1;
@@ -195,19 +202,14 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se
if (cr_session->soopt) {
if (!ucr->quiet)
uwsgi_log("unable to connect() to uwsgi instance \"%.*s\": %s\n", (int) cr_session->instance_address_len, cr_session->instance_address, strerror(cr_session->soopt));
uwsgi_log("[uwsgi-%s] unable to connect() to node \"%.*s\": %s\n", ucr->short_name, (int) cr_session->instance_address_len, cr_session->instance_address, strerror(cr_session->soopt));
}
else if (cr_session->timed_out) {
if (cr_session->instance_address_len > 0) {
/*
if (cr_session->status == COREROUTER_STATUS_CONNECTING) {
if (cr_session->connecting) {
if (!ucr->quiet)
uwsgi_log("unable to connect() to uwsgi instance \"%.*s\": timeout\n", (int) cr_session->instance_address_len, cr_session->instance_address);
uwsgi_log("[uwsgi-%s] unable to connect() to node \"%.*s\": timeout\n", ucr->short_name, (int) cr_session->instance_address_len, cr_session->instance_address);
}
else if (cr_session->status == COREROUTER_STATUS_RESPONSE) {
uwsgi_log("timeout waiting for instance \"%.*s\"\n", (int) cr_session->instance_address_len, cr_session->instance_address);
}
*/
}
}
@@ -239,6 +241,29 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se
cr_session->tmp_socket_name = NULL;
}
if (!cr_session->retry) goto end;
// check for max retries
if (cr_session->retries >= (size_t) ucr->max_retries) goto end;
cr_session->retries++;
// reset error and timeout
cr_session->instance_failed = 0;
cr_session->timeout = corerouter_reset_timeout(ucr, cr_session);
cr_session->timed_out = 0;
cr_session->soopt = 0;
// reset nodes
cr_session->un = NULL;
cr_session->static_node = NULL;
cr_session->instance_fd = -1;
// reset hooks (safe as fd is closed)
cr_session->event_hook_read = NULL;
cr_session->event_hook_write = NULL;
cr_session->event_hook_instance_read = NULL;
cr_session->event_hook_instance_write = NULL;
if (ucr->fallback) {
// ok let's try with the fallback nodes
if (!cr_session->fallback) {
@@ -252,35 +277,18 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se
cr_session->instance_address = cr_session->fallback->value;
cr_session->instance_address_len = cr_session->fallback->len;
// reset error and timeout
cr_session->timeout = corerouter_reset_timeout(ucr, cr_session);
cr_session->timed_out = 0;
cr_session->soopt = 0;
// reset nodes
cr_session->un = NULL;
cr_session->static_node = NULL;
cr_session->pass_fd = is_unix(cr_session->instance_address, cr_session->instance_address_len);
cr_session->instance_fd = uwsgi_connectn(cr_session->instance_address, cr_session->instance_address_len, 0, 1);
if (cr_session->instance_fd < 0) {
cr_session->instance_failed = 1;
cr_session->soopt = errno;
corerouter_close_session(ucr, cr_session);
return;
if (cr_session->retry(ucr, cr_session)) {
if (!cr_session->instance_failed) goto end;
}
ucr->cr_table[cr_session->instance_fd] = cr_session;
//cr_session->status = COREROUTER_STATUS_CONNECTING;
ucr->cr_table[cr_session->instance_fd] = cr_session;
event_queue_add_fd_write(ucr->queue, cr_session->instance_fd);
return;
}
cr_session->instance_address = NULL;
cr_session->instance_address_len = 0;
if (cr_session->retry(ucr, cr_session)) {
if (!cr_session->instance_failed) goto end;
}
return;
}
end:
@@ -331,23 +339,10 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) {
if (urbt->key <= current) {
cr_session = (struct corerouter_session *) urbt->data;
cr_session->timed_out = 1;
if (cr_session->retry) {
cr_session->retry = 0;
/*
TODO allows retry
ucr->switch_events(ucr, cr_session, -1);
*/
if (cr_session->retry) {
cr_del_timeout(ucr, cr_session);
cr_session->timeout = cr_add_fake_timeout(ucr, cr_session);
}
else {
cr_session->timeout = corerouter_reset_timeout(ucr, cr_session);
}
}
else {
corerouter_close_session(ucr, cr_session);
if (cr_session->connecting) {
cr_session->instance_failed = 1;
}
corerouter_close_session(ucr, cr_session);
continue;
}
@@ -572,7 +567,7 @@ struct corerouter_session *corerouter_alloc_session(struct uwsgi_corerouter *ucr
ucr->cr_table[new_connection]->fd = new_connection;
ucr->cr_table[new_connection]->instance_fd = -1;
// map courerouter and socket
// map corerouter and socket
ucr->cr_table[new_connection]->corerouter = ucr;
ucr->cr_table[new_connection]->ugs = ugs;
@@ -771,7 +766,11 @@ void uwsgi_corerouter_loop(int id, void *data) {
#endif
#endif
corerouter_alloc_session(ucr, ugs, new_connection, (struct sockaddr *) &cr_addr, cr_addr_len);
struct corerouter_session *cr = corerouter_alloc_session(ucr, ugs, new_connection, (struct sockaddr *) &cr_addr, cr_addr_len);
//something wrong in the allocation
if (cr->instance_failed) {
corerouter_close_session(ucr, cr);
}
}
else if (ugs->subscription) {
uwsgi_corerouter_manage_subscription(ucr, id, ugs);
@@ -834,9 +833,21 @@ void uwsgi_corerouter_loop(int id, void *data) {
}
}
// not having a hook could mean a previous event in the loop cleared it...
if (!hook) {
uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received !!!\n");
corerouter_close_session(ucr, cr_session);
// a single event cannot be unexpected..
if (nevents == 1) {
if (interesting_fd == cr_session->instance_fd) {
uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received from backend instance (fd: %d nevents: %d) !!!\n", interesting_fd, nevents);
}
else if (interesting_fd == cr_session->fd) {
uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received from client (fd: %d nevents: %d)!!!\n", interesting_fd, nevents);
}
else {
uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received !!!\n");
}
corerouter_close_session(ucr, cr_session);
}
continue;
}
@@ -860,7 +871,7 @@ void uwsgi_corerouter_loop(int id, void *data) {
}
int uwsgi_courerouter_has_has_backends(struct uwsgi_corerouter *ucr) {
int uwsgi_corerouter_has_backends(struct uwsgi_corerouter *ucr) {
if (ucr->has_backends) return 1;
@@ -896,8 +907,11 @@ int uwsgi_corerouter_init(struct uwsgi_corerouter *ucr) {
if (!ucr->nevents)
ucr->nevents = 64;
if (!ucr->max_retries)
ucr->max_retries = 3;
ucr->has_backends = uwsgi_courerouter_has_has_backends(ucr);
ucr->has_backends = uwsgi_corerouter_has_backends(ucr);
uwsgi_corerouter_setup_sockets(ucr);
@@ -921,7 +935,7 @@ int uwsgi_corerouter_init(struct uwsgi_corerouter *ucr) {
struct uwsgi_plugin corerouter_plugin = {
.name = "courerouter",
.name = "corerouter",
};
void corerouter_send_stats(struct uwsgi_corerouter *ucr) {
@@ -968,6 +982,29 @@ void corerouter_send_stats(struct uwsgi_corerouter *ucr) {
if (uwsgi_stats_list_close(us)) goto end0;
if (uwsgi_stats_comma(us)) goto end0;
if (ucr->static_nodes) {
if (uwsgi_stats_key(us , "static_nodes")) goto end0;
if (uwsgi_stats_list_open(us)) goto end0;
struct uwsgi_string_list *usl = ucr->static_nodes;
while(usl) {
if (uwsgi_stats_object_open(us)) goto end0;
if (uwsgi_stats_keyvaln_comma(us, "name", usl->value, usl->len)) goto end0;
if (uwsgi_stats_keylong_comma(us, "hits", (unsigned long long) usl->custom2)) goto end0;
if (uwsgi_stats_keylong(us, "grace", (unsigned long long) usl->custom)) goto end0;
if (uwsgi_stats_object_close(us)) goto end0;
usl = usl->next;
if (usl) {
if (uwsgi_stats_comma(us)) goto end0;
}
}
if (uwsgi_stats_list_close(us)) goto end0;
if (uwsgi_stats_comma(us)) goto end0;
}
if (ucr->has_subscription_sockets) {
if (uwsgi_stats_key(us , "subscriptions")) goto end0;
if (uwsgi_stats_list_open(us)) goto end0;
+7 -3
View File
@@ -38,6 +38,8 @@ struct uwsgi_corerouter {
int use_cache;
int nevents;
int max_retries;
char *magic_table[256];
int queue;
@@ -102,14 +104,13 @@ struct corerouter_session {
uint16_t hostname_len;
int has_key;
int retry;
int connecting;
char *instance_address;
uint64_t instance_address_len;
struct uwsgi_subscribe_node *un;
struct uwsgi_string_list *static_node;
int pass_fd;
int soopt;
int timed_out;
@@ -139,6 +140,8 @@ struct corerouter_session {
ssize_t (*event_hook_instance_write)(struct corerouter_session *);
void (*close)(struct corerouter_session *);
int (*retry)(struct uwsgi_corerouter *, struct corerouter_session *);
size_t retries;
struct uwsgi_buffer *buffer;
size_t buffer_len;
@@ -151,6 +154,7 @@ struct corerouter_session {
};
void uwsgi_opt_corerouter(char *, char *, void *);
void uwsgi_opt_undeferred_corerouter(char *, char *, void *);
void uwsgi_opt_corerouter_use_socket(char *, char *, void *);
void uwsgi_opt_corerouter_use_base(char *, char *, void *);
void uwsgi_opt_corerouter_use_pattern(char *, char *, void *);
@@ -181,7 +185,7 @@ int uwsgi_cr_map_use_cs(struct uwsgi_corerouter *, struct corerouter_session *);
int uwsgi_cr_map_use_to(struct uwsgi_corerouter *, struct corerouter_session *);
int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *, struct corerouter_session *);
int uwsgi_courerouter_has_has_backends(struct uwsgi_corerouter *);
int uwsgi_corerouter_has_backends(struct uwsgi_corerouter *);
int uwsgi_cr_hook_read(struct corerouter_session *, ssize_t (*)(struct corerouter_session *));
int uwsgi_cr_hook_write(struct corerouter_session *, ssize_t (*)(struct corerouter_session *));
+7
View File
@@ -41,6 +41,10 @@ void uwsgi_corerouter_setup_sockets(struct uwsgi_corerouter *ucr) {
}
else {
ugs->port = strchr(ugs->name, ':');
int current_defer_accept = uwsgi.no_defer_accept;
if (ugs->no_defer) {
uwsgi.no_defer_accept = 1;
}
if (ugs->fd == -1) {
if (ugs->port) {
ugs->fd = bind_to_tcp(ugs->name, uwsgi.listen_queue, ugs->port);
@@ -51,6 +55,9 @@ void uwsgi_corerouter_setup_sockets(struct uwsgi_corerouter *ucr) {
ugs->fd = bind_to_unix(ugs->name, uwsgi.listen_queue, uwsgi.chmod_socket, uwsgi.abstract_socket);
}
}
if (ugs->no_defer) {
uwsgi.no_defer_accept = current_defer_accept;
}
}
// put socket in non-blocking mode
uwsgi_socket_nb(ugs->fd);
+4
View File
@@ -233,6 +233,8 @@ ssize_t fr_instance_send_request_header(struct corerouter_session * cs) {
ssize_t fr_instance_connected(struct corerouter_session * cs) {
cs->connecting = 0;
socklen_t solen = sizeof(int);
// first check for errors
@@ -250,6 +252,7 @@ ssize_t fr_instance_connected(struct corerouter_session * cs) {
cs->buffer_pos = 0;
// ok instance is connected, wait for write again
if (cs->static_node) cs->static_node->custom2++;
if (cs->un) cs->un->requests++;
uwsgi_cr_hook_instance_write(cs, fr_instance_send_request_header);
// return a value > 0
@@ -304,6 +307,7 @@ ssize_t fr_recv_uwsgi_vars(struct corerouter_session * cs) {
// map the instance
cs->corerouter->cr_table[cs->instance_fd] = cs;
// wait for connection
cs->connecting = 1;
uwsgi_cr_hook_instance_write(cs, fr_instance_connected);
}
+6
View File
@@ -253,6 +253,8 @@ ssize_t uwsgi_gevent_hook_input_read(struct wsgi_request *wsgi_req, char *tmp_bu
UWSGI_RELEASE_GIL;
ssize_t rlen = read(wsgi_req->poll.fd, tmp_buf+*tmp_pos, remains);
if (rlen <= 0) {
if (rlen < 0)
uwsgi_error("[uwsgi-gevent] read()");
UWSGI_GET_GIL
stop_the_watchers_and_clear
return -1;
@@ -328,6 +330,10 @@ void uwsgi_gevent_nb_write(struct wsgi_request *wsgi_req, PyObject *str) {
PyObject *ret;
char *content = PyString_AsString(str);
size_t content_len = PyString_Size(str);
// do not try to write empty chunks
if (content_len == 0) return;
/// create a watcher for writes
PyObject *watcher = PyObject_CallMethod(ugevent.hub_loop, "io", "ii", wsgi_req->poll.fd, 2);
if (!watcher) goto error;
+88 -20
View File
@@ -24,6 +24,7 @@ struct uwsgi_http {
#ifdef UWSGI_SSL
int https_export_cert;
#endif
struct uwsgi_string_list *stud_prefix;
} uhttp;
@@ -134,6 +135,7 @@ struct uwsgi_option http_options[] = {
{"http-stats-server", required_argument, 0, "run the http router stats server", uwsgi_opt_set_str, &uhttp.cr.stats_server, 0},
{"http-ss", required_argument, 0, "run the http router stats server", uwsgi_opt_set_str, &uhttp.cr.stats_server, 0},
{"http-harakiri", required_argument, 0, "enable http router harakiri", uwsgi_opt_set_int, &uhttp.cr.harakiri, 0},
{"http-stud-prefix", required_argument, 0, "expect a stud prefix (1byte family + 4/16 bytes address) on connections from the specified address", uwsgi_opt_add_addr_list, &uhttp.stud_prefix, 0},
{0, 0, 0, 0, 0, 0, 0},
};
@@ -175,6 +177,11 @@ struct http_session {
size_t post_buf_len;
off_t post_buf_pos;
// 1 (family) + 4/16 (addr)
char stud_prefix[17];
size_t stud_prefix_remains;
size_t stud_prefix_pos;
};
@@ -354,6 +361,11 @@ int http_parse(struct http_session *h_session, size_t http_req_len) {
// UWSGI_ROUTER
if (http_add_uwsgi_var(h_session, "UWSGI_ROUTER", 12, "http", 4)) return -1;
// stud HTTPS
if (h_session->stud_prefix_pos > 0) {
if (http_add_uwsgi_var(h_session, "HTTPS", 5, "on", 2)) return -1;
}
#ifdef UWSGI_SSL
// HTTPS (adapted from nginx)
if (h_session->cs.ugs->mode == UWSGI_HTTP_SSL) {
@@ -399,9 +411,9 @@ int http_parse(struct http_session *h_session, size_t http_req_len) {
while (ptr < watermark) {
if (*ptr == '\r') {
if (ptr + 1 >= watermark)
return 0;
break;
if (*(ptr + 1) != '\n')
return 0;
break;
// multiline header ?
if (ptr + 2 < watermark) {
if (*(ptr + 2) == ' ' || *(ptr + 2) == '\t') {
@@ -514,9 +526,17 @@ ssize_t hr_read_ssl_body(struct corerouter_session * cs) {
struct http_session *hs = (struct http_session *) cs;
int ret = SSL_read(hs->ssl, hs->post_buf->buf, hs->post_buf_max);
if (ret > 0) {
// fix waiting
if (cs->event_hook_write) {
uwsgi_cr_hook_write(cs, NULL);
int ret2 = SSL_pending(hs->ssl);
if (ret2 > 0) {
if (uwsgi_buffer_fix(hs->post_buf, hs->post_buf->len + ret2 )) {
uwsgi_log("[uwsgi-https] cannot fix the buffer to %d\n", hs->post_buf->len + ret2);
return -1;
}
if (SSL_read(hs->ssl, hs->post_buf->buf + ret, ret2) != ret2) {
uwsgi_log("[uwsgi-https] SSL_read() on %d bytes of pending data failed\n", ret2);
return -1;
}
ret += ret2;
}
hs->post_buf_len = ret;
hs->post_buf_pos = 0;
@@ -591,9 +611,6 @@ ssize_t hr_write_ssl_response(struct corerouter_session * cs) {
if (ret > 0) {
cs->buffer_pos += ret;
if (cs->event_hook_read) {
uwsgi_cr_hook_read(cs, NULL);
}
// could be a partial write
uwsgi_cr_hook_write(cs, hr_write_ssl_response);
// ok this response chunk is sent, let's wait for another one
@@ -762,6 +779,12 @@ ssize_t hr_instance_send_request(struct corerouter_session * cs) {
ssize_t hr_instance_send_request_header(struct corerouter_session * cs) {
#ifdef __BIG_ENDIAN__
// on the first round fix endianess
if (cs->buffer_pos == 0) {
cs->uh.pktsize = uwsgi_swap16(cs->uh.pktsize);
}
#endif
ssize_t len = write(cs->instance_fd, &cs->uh + cs->buffer_pos, 4 - cs->buffer_pos);
if (len < 0) {
cr_try_again;
@@ -775,6 +798,10 @@ ssize_t hr_instance_send_request_header(struct corerouter_session * cs) {
// for response
if (cs->buffer_pos == 4) {
cs->buffer_pos = 0;
#ifdef __BIG_ENDIAN__
// on the last round restore endianess
cs->uh.pktsize = uwsgi_swap16(cs->uh.pktsize);
#endif
uwsgi_cr_hook_instance_write(cs, hr_instance_send_request);
}
@@ -804,9 +831,6 @@ ssize_t hr_send_expect_continue(struct corerouter_session * cs) {
if (ret > 0) {
len = ret;
cs->buffer_pos += ret;
if (cs->event_hook_read) {
uwsgi_cr_hook_read(cs, NULL);
}
// could be a partial write
uwsgi_cr_hook_write(cs, hr_send_expect_continue);
goto done;
@@ -856,6 +880,8 @@ done:
}
ssize_t hr_instance_connected(struct corerouter_session * cs) {
cs->connecting = 0;
socklen_t solen = sizeof(int);
// first check for errors
@@ -882,6 +908,7 @@ ssize_t hr_instance_connected(struct corerouter_session * cs) {
return 1;
}
// ok instance is connected, wait for write again
if (cs->static_node) cs->static_node->custom2++;
if (cs->un) cs->un->requests++;
uwsgi_cr_hook_instance_write(cs, hr_instance_send_request_header);
// return a value > 0
@@ -895,19 +922,21 @@ ssize_t hr_recv_http_ssl(struct corerouter_session * cs) {
// be sure buffer does not grow over 64k
cs->buffer->limit = UMAX16;
// try to always leave 4k available
if (uwsgi_buffer_ensure(cs->buffer, uwsgi.page_size)) return -1;
if (uwsgi_buffer_ensure(cs->buffer, uwsgi.page_size)) {
uwsgi_log("[uwsgi-https] cannot ensure the buffer (size: %d)\n", cs->buffer->len);
return -1;
}
struct http_session *hs = (struct http_session *) cs;
int ret = SSL_read(hs->ssl, cs->buffer->buf + cs->buffer_pos, cs->buffer->len - cs->buffer_pos);
if (ret > 0) {
// fix waiting
if (cs->event_hook_write) {
uwsgi_cr_hook_write(cs, NULL);
uwsgi_cr_hook_read(cs, hr_recv_http_ssl);
}
int ret2 = SSL_pending(hs->ssl);
if (ret2 > 0) {
if (uwsgi_buffer_fix(cs->buffer, cs->buffer->len + ret2 )) return -1;
if (uwsgi_buffer_fix(cs->buffer, cs->buffer->len + ret2 )) {
uwsgi_log("[uwsgi-https] cannot fix the buffer to %d\n", cs->buffer->len + ret2 );
return -1;
}
if (SSL_read(hs->ssl, cs->buffer->buf + cs->buffer_pos + ret, ret2) != ret2) {
uwsgi_log("[uwsgi-https] SSL_read() on %d bytes of pending data failed\n", ret2);
return -1;
}
ret += ret2;
@@ -1024,6 +1053,7 @@ ssize_t hs_http_manage(struct corerouter_session * cs, ssize_t len) {
// map the instance
cs->corerouter->cr_table[cs->instance_fd] = cs;
// wait for connection
cs->connecting = 1;
uwsgi_cr_hook_instance_write(cs, hr_instance_connected);
break;
}
@@ -1049,6 +1079,31 @@ void hr_session_close(struct corerouter_session *cs) {
}
}
static ssize_t hr_recv_stud4(struct corerouter_session * cs) {
struct http_session *hs = (struct http_session *) cs;
ssize_t len = read(cs->fd, hs->stud_prefix + hs->stud_prefix_pos, hs->stud_prefix_remains - hs->stud_prefix_pos);
if (len < 0) {
cr_try_again;
uwsgi_error("hr_recv_stud4()");
return -1;
}
hs->stud_prefix_pos += len;
if (hs->stud_prefix_pos == hs->stud_prefix_remains) {
if (hs->stud_prefix[0] != AF_INET) {
uwsgi_log("[uwsgi-http] invalid stud prefix\n");
return -1;
}
// set the passed ip address
memcpy(&hs->ip_addr, hs->stud_prefix + 1, 4);
uwsgi_cr_hook_read(cs, hr_recv_http);
}
return len;
}
#ifdef UWSGI_SSL
void hr_session_ssl_close(struct corerouter_session *cs) {
hr_session_close(cs);
@@ -1079,6 +1134,17 @@ void http_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socke
cs->modifier1 = uhttp.modifier1;
if (sa && sa->sa_family == AF_INET) {
hs->ip_addr = ((struct sockaddr_in *) sa)->sin_addr.s_addr;
struct uwsgi_string_list *usl = uhttp.stud_prefix;
while(usl) {
if (!memcmp(&hs->ip_addr, usl->value, 4)) {
hs->stud_prefix_remains = 5;
uwsgi_cr_hook_read(cs, hr_recv_stud4);
break;
}
usl = usl->next;
}
}
hs->rnrn = 0;
@@ -1100,7 +1166,9 @@ void http_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socke
}
else {
#endif
uwsgi_cr_hook_read(cs, hr_recv_http);
if (!cs->event_hook_read) {
uwsgi_cr_hook_read(cs, hr_recv_http);
}
cs->close = hr_session_close;
#ifdef UWSGI_SSL
}
@@ -1117,7 +1185,7 @@ int http_init() {
uhttp.cr.session_size = sizeof(struct http_session);
uhttp.cr.alloc_session = http_alloc_session;
if (uhttp.cr.has_sockets && !uwsgi.sockets && !uwsgi_courerouter_has_has_backends(&uhttp.cr)) {
if (uhttp.cr.has_sockets && !uwsgi.sockets && !uwsgi_corerouter_has_backends(&uhttp.cr)) {
uwsgi_new_socket(uwsgi_concat2("127.0.0.1:0", ""));
uhttp.cr.use_socket = 1;
uhttp.cr.socket_num = 0;
+3 -1
View File
@@ -308,6 +308,8 @@ static const luaL_reg uwsgi_api[] = {
static int uwsgi_lua_input(lua_State *L) {
struct wsgi_request *wsgi_req = current_wsgi_req();
int fd = wsgi_req->async_post ?
fileno(wsgi_req->async_post) : wsgi_req->poll.fd;
ssize_t sum, len, total;
char *buf, *ptr;
@@ -330,7 +332,7 @@ static int uwsgi_lua_input(lua_State *L) {
ptr = buf;
while(total) {
len = read(wsgi_req->poll.fd, ptr, total);
len = read(fd, ptr, total);
ptr += len;
total -= len;
}
+48 -1
View File
@@ -22,8 +22,10 @@ struct uwsgi_php {
struct uwsgi_string_list *index;
struct uwsgi_string_list *set;
struct uwsgi_string_list *append_config;
struct uwsgi_string_list *vars;
char *docroot;
char *app;
char *app_qs;
size_t ini_size;
int dump_config;
char *server_software;
@@ -49,6 +51,8 @@ struct uwsgi_option uwsgi_php_options[] = {
{"php-allowed-script", required_argument, 0, "list the allowed php scripts (require absolute path)", uwsgi_opt_add_string_list, &uphp.allowed_scripts, 0},
{"php-server-software", required_argument, 0, "force php SERVER_SOFTWARE", uwsgi_opt_set_str, &uphp.server_software, 0},
{"php-app", required_argument, 0, "force the php file to run at each request", uwsgi_opt_set_str, &uphp.app, 0},
{"php-var", required_argument, 0, "add/overwrite a CGI variable at each request", uwsgi_opt_add_string_list, &uphp.vars, 0},
{"php-app-qs", required_argument, 0, "when in app mode force QUERY_STRING to the specified value + PATH_INFO", uwsgi_opt_set_str, &uphp.app_qs, 0},
{"php-dump-config", no_argument, 0, "dump php config (if modified via --php-set or append options)", uwsgi_opt_true, &uphp.dump_config, 0},
{0, 0, 0, 0, 0, 0, 0},
@@ -285,6 +289,9 @@ static void sapi_uwsgi_register_variables(zval *track_vars_array TSRMLS_DC)
}
php_register_variable_safe("PATH_INFO", wsgi_req->path_info, wsgi_req->path_info_len, track_vars_array TSRMLS_CC);
if (wsgi_req->query_string_len > 0) {
php_register_variable_safe("QUERY_STRING", wsgi_req->query_string, wsgi_req->query_string_len, track_vars_array TSRMLS_CC);
}
php_register_variable_safe("SCRIPT_NAME", wsgi_req->script_name, wsgi_req->script_name_len, track_vars_array TSRMLS_CC);
php_register_variable_safe("SCRIPT_FILENAME", wsgi_req->file, wsgi_req->file_len, track_vars_array TSRMLS_CC);
@@ -304,6 +311,16 @@ static void sapi_uwsgi_register_variables(zval *track_vars_array TSRMLS_DC)
php_register_variable_safe("PHP_SELF", wsgi_req->script_name, wsgi_req->script_name_len, track_vars_array TSRMLS_CC);
struct uwsgi_string_list *usl = uphp.vars;
while(usl) {
char *equal = strchr(usl->value, '=');
if (equal) {
php_register_variable_safe( estrndup(usl->value, equal-usl->value),
equal+1, strlen(equal+1), track_vars_array TSRMLS_CC);
}
usl = usl->next;
}
}
@@ -658,6 +675,16 @@ int uwsgi_php_init(void) {
uwsgi_log("--- end of PHP custom config ---\n");
}
// fix docroot
if (uphp.docroot) {
char *orig_docroot = uphp.docroot;
uphp.docroot = uwsgi_expand_path(uphp.docroot, strlen(uphp.docroot), NULL);
if (!uphp.docroot) {
uwsgi_log("unable to set php docroot to %s\n", orig_docroot);
exit(1);
}
}
uwsgi_sapi_module.startup(&uwsgi_sapi_module);
// filling http status codes
@@ -766,6 +793,27 @@ int uwsgi_php_request(struct wsgi_request *wsgi_req) {
if (uphp.app) {
strcpy(real_filename, uphp.app);
if (wsgi_req->path_info_len == 1 && wsgi_req->path_info[0] == '/') {
goto appready;
}
if (uphp.app_qs) {
size_t app_qs_len = strlen(uphp.app_qs);
size_t qs_len = wsgi_req->path_info_len + app_qs_len;
if (wsgi_req->query_string_len > 0) {
qs_len += 1 + wsgi_req->query_string_len;
}
char *qs = ecalloc(1, qs_len+1);
memcpy(qs, uphp.app_qs, app_qs_len);
memcpy(qs+app_qs_len, wsgi_req->path_info, wsgi_req->path_info_len);
if (wsgi_req->query_string_len > 0) {
char *ptr = qs+app_qs_len+wsgi_req->path_info_len;
*ptr = '&';
memcpy(ptr+1, wsgi_req->query_string, wsgi_req->query_string_len);
}
wsgi_req->query_string = qs;
wsgi_req->query_string_len = qs_len;
}
appready:
wsgi_req->path_info = "";
wsgi_req->path_info_len = 0;
goto secure2;
@@ -891,7 +939,6 @@ secure2:
secure3:
if (wsgi_req->document_root[wsgi_req->document_root_len-1] == '/') {
wsgi_req->script_name = real_filename + (wsgi_req->document_root_len-1);
}
+1
View File
@@ -51,3 +51,4 @@ int uwsgi_perl_obj_can(SV *, char *, size_t);
int uwsgi_perl_obj_isa(SV *, char *);
int init_psgi_app(struct wsgi_request *, char *, uint16_t, PerlInterpreter **);
PerlInterpreter *uwsgi_perl_new_interpreter(void);
int uwsgi_perl_mule(char *);
+20
View File
@@ -233,6 +233,8 @@ xs_init(pTHX)
/* DynaLoader is a special case */
newXS("DynaLoader::boot_DynaLoader", boot_DynaLoader, file);
if (!uperl.tmp_input_stash) goto nonworker;
newXS("uwsgi::input::new", XS_input, "uwsgi::input");
newXS("uwsgi::input::read", XS_input_read, "uwsgi::input");
newXS("uwsgi::input::seek", XS_input_seek, "uwsgi::input");
@@ -253,6 +255,8 @@ xs_init(pTHX)
uperl.tmp_streaming_stash[uperl.tmp_current_i] = gv_stashpv("uwsgi::streaming", 0);
nonworker:
#ifdef UWSGI_EMBEDDED
init_perl_embedded_module();
#endif
@@ -496,3 +500,19 @@ void uwsgi_psgi_app() {
}
int uwsgi_perl_mule(char *opt) {
if (uwsgi_endswith(opt, ".pl")) {
PERL_SET_CONTEXT(uperl.main[0]);
uperl.embedding[1] = opt;
if (perl_parse(uperl.main[0], xs_init, 2, uperl.embedding, NULL)) {
return 0;
}
perl_run(uperl.main[0]);
return 1;
}
return 0;
}
+29
View File
@@ -595,6 +595,32 @@ void uwsgi_perl_enable_threads(void) {
#endif
}
int uwsgi_perl_signal_handler(uint8_t sig, void *handler) {
int ret = 0;
dSP;
ENTER;
SAVETMPS;
PUSHMARK(SP);
XPUSHs( sv_2mortal(newSViv(sig)));
PUTBACK;
call_sv( SvRV((SV*)handler), G_DISCARD);
if(SvTRUE(ERRSV)) {
uwsgi_log("[uwsgi-perl error] %s\n", SvPV_nolen(ERRSV));
ret = -1;
}
SPAGAIN;
PUTBACK;
FREETMPS;
LEAVE;
return ret;
}
struct uwsgi_plugin psgi_plugin = {
.name = "psgi",
@@ -606,6 +632,9 @@ struct uwsgi_plugin psgi_plugin = {
.mount_app = uwsgi_perl_mount_app,
.init_thread = uwsgi_perl_init_thread,
.signal_handler = uwsgi_perl_signal_handler,
.mule = uwsgi_perl_mule,
.post_fork = uwsgi_perl_post_fork,
.request = uwsgi_perl_request,
+5
View File
@@ -59,6 +59,11 @@ int psgi_response(struct wsgi_request *wsgi_req, AV *response) {
}
#endif
if (SvTYPE(response) != SVt_PVAV) {
uwsgi_log("invalid PSGI response type\n");
return UWSGI_OK;
}
status_code = av_fetch(response, 0, 0);
if (!status_code) { uwsgi_log("invalid PSGI status code\n"); return UWSGI_OK;}
+51
View File
@@ -1,6 +1,7 @@
#include "psgi.h"
extern struct uwsgi_server uwsgi;
extern struct uwsgi_plugin psgi_plugin;
#ifdef UWSGI_ASYNC
@@ -149,6 +150,27 @@ clear:
}
XS(XS_register_signal) {
dXSARGS;
if (!uwsgi.master_process) {
XSRETURN_NO;
}
psgi_check_args(3);
uint8_t signum = SvIV(ST(0));
STRLEN kindlen;
char *kind = SvPV(ST(1), kindlen);
if (uwsgi_register_signal(signum, kind, (void *) newRV_inc(ST(2)), psgi_plugin.modifier1)) {
XSRETURN_NO;
}
XSRETURN_YES;
}
XS(XS_log) {
dXSARGS;
@@ -219,6 +241,33 @@ XS(XS_suspend) {
XSRETURN_UNDEF;
}
XS(XS_signal_wait) {
dXSARGS;
psgi_check_args(0);
struct wsgi_request *wsgi_req = current_wsgi_req();
int received_signal = -1;
wsgi_req->signal_received = -1;
if (items > 0) {
received_signal = uwsgi_signal_wait(SvIV(ST(0)));
}
else {
received_signal = uwsgi_signal_wait(-1);
}
if (received_signal < 0) {
XSRETURN_NO;
}
wsgi_req->signal_received = received_signal;
XSRETURN_YES;
}
void init_perl_embedded_module() {
psgi_xs(reload);
psgi_xs(cache_set);
@@ -231,5 +280,7 @@ void init_perl_embedded_module() {
psgi_xs(async_connect);
psgi_xs(suspend);
psgi_xs(signal);
psgi_xs(register_signal);
psgi_xs(signal_wait);
}
+3 -3
View File
@@ -180,11 +180,11 @@ static PyObject *uwsgi_Input_read(uwsgi_Input *self, PyObject *args) {
ssize_t rlen = up.hook_wsgi_input_read(self->wsgi_req, tmp_buf, remains, &tmp_pos);
if (rlen < 0) {
free(tmp_buf);
return PyErr_Format(PyExc_IOError, "error reading 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);
return PyErr_Format(PyExc_IOError, "error reading for wsgi.input data: Content-Length %llu requested %llu received %llu pos %llu+%llu", (unsigned long long) self->wsgi_req->post_cl, (unsigned long long) remains, (unsigned long long) tmp_pos, (unsigned long long) self->pos, (unsigned long long) tmp_pos);
}
else if (tmp_pos == 0) {
else if (rlen == 0) {
free(tmp_buf);
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);
return PyErr_Format(PyExc_IOError, "error waiting for wsgi.input data: Content-Length %llu requested %llu received %llu pos %llu+%llu", (unsigned long long) self->wsgi_req->post_cl, (unsigned long long) remains, (unsigned long long) tmp_pos, (unsigned long long) self->pos, (unsigned long long) tmp_pos);
}
self->pos += tmp_pos;
+15 -4
View File
@@ -1024,10 +1024,21 @@ VALUE init_rack_app( VALUE script ) {
VALUE rack = rb_const_get(rb_cObject, rb_intern("Rack"));
#ifdef RUBY19
if (rb_eval_string("module Rack;class BodyProxy;def each(&block);@body.each(&block);end;end;end")) {
if (uwsgi.mywid <= 1) {
uwsgi_log("Rack::BodyProxy successfully patched for ruby 1.9.x\n");
}
if (rb_funcall(rack, rb_intern("const_defined?"), 1, ID2SYM(rb_intern("BodyProxy"))) == Qtrue) {
VALUE bodyproxy = rb_const_get(rack, rb_intern("BodyProxy"));
// get the list of available instance_methods
VALUE argv = Qfalse;
VALUE methods_list = rb_class_instance_methods(1, &argv, bodyproxy);
#ifdef UWSGI_DEBUG
uwsgi_log("%s\n", RSTRING_PTR(rb_inspect(methods_list)));
#endif
if (rb_ary_includes(methods_list, ID2SYM(rb_intern("each"))) == Qfalse) {
if (rb_eval_string("module Rack;class BodyProxy;def each(&block);@body.each(&block);end;end;end")) {
if (uwsgi.mywid <= 1) {
uwsgi_log("Rack::BodyProxy successfully patched for ruby 1.9.x\n");
}
}
}
}
#endif
+290 -20
View File
@@ -2,65 +2,335 @@
uWSGI rawrouter
requires:
- async
- caching
- pcre (optional)
*/
#include "../../uwsgi.h"
#include "../corerouter/cr.h"
struct uwsgi_rawrouter {
struct uwsgi_corerouter cr;
int xclient;
} urr;
extern struct uwsgi_server uwsgi;
#include "rr.h"
struct uwsgi_rawrouter urr;
struct rawrouter_session {
struct corerouter_session crs;
in_addr_t ip_addr;
// XCLIENT ADDR=xxx\r\n
char xclient[13+INET_ADDRSTRLEN+2];
size_t xclient_len;
off_t xclient_pos;
size_t xclient_remains;
// placeholder for \r\n
size_t xclient_rn;
};
struct uwsgi_option rawrouter_options[] = {
{"rawrouter", required_argument, 0, "run the rawrouter on the specified port", uwsgi_opt_corerouter, &urr, 0},
{"rawrouter", required_argument, 0, "run the rawrouter on the specified port", uwsgi_opt_undeferred_corerouter, &urr, 0},
{"rawrouter-processes", required_argument, 0, "prefork the specified number of rawrouter processes", uwsgi_opt_set_int, &urr.cr.processes, 0},
{"rawrouter-workers", required_argument, 0, "prefork the specified number of rawrouter processes", uwsgi_opt_set_int, &urr.cr.processes, 0},
{"rawrouter-zerg", required_argument, 0, "attach the rawrouter to a zerg server", uwsgi_opt_corerouter_zerg, &urr, 0 },
{"rawrouter-use-cache", no_argument, 0, "use uWSGI cache as address->server mapper for the rawrouter", uwsgi_opt_true, &urr.cr.use_cache, 0},
{"rawrouter-zerg", required_argument, 0, "attach the rawrouter to a zerg server", uwsgi_opt_corerouter_zerg, &urr, 0},
{"rawrouter-use-cache", no_argument, 0, "use uWSGI cache as hostname->server mapper for the rawrouter", uwsgi_opt_true, &urr.cr.use_cache, 0},
{"rawrouter-use-pattern", required_argument, 0, "use a pattern for rawrouter address->server mapping", uwsgi_opt_corerouter_use_pattern, &urr, 0},
{"rawrouter-use-base", required_argument, 0, "use a base dir for rawrouter address->server mapping", uwsgi_opt_corerouter_use_base, &urr, 0},
{"rawrouter-use-pattern", required_argument, 0, "use a pattern for rawrouter hostname->server mapping", uwsgi_opt_corerouter_use_pattern, &urr, 0},
{"rawrouter-use-base", required_argument, 0, "use a base dir for rawrouter hostname->server mapping", uwsgi_opt_corerouter_use_base, &urr, 0},
{"rawrouter-fallback", required_argument, 0, "fallback to the specified node in case of error", uwsgi_opt_add_string_list, &urr.cr.fallback, 0},
{"rawrouter-use-cluster", no_argument, 0, "load balance to nodes subscribed to the cluster", uwsgi_opt_true, &urr.cr.use_cluster, 0},
{"rawrouter-use-code-string", required_argument, 0, "use code string as address->server mapper for the rawrouter", uwsgi_opt_corerouter_cs, &urr, 0},
{"rawrouter-use-code-string", required_argument, 0, "use code string as hostname->server mapper for the rawrouter", uwsgi_opt_corerouter_cs, &urr, 0},
{"rawrouter-use-socket", optional_argument, 0, "forward request to the specified uwsgi socket", uwsgi_opt_corerouter_use_socket, &urr, 0},
{"rawrouter-to", required_argument, 0, "forward requests to the specified uwsgi server (you can specify it multiple times for load balancing)", uwsgi_opt_add_string_list, &urr.cr.static_nodes, 0},
{"rawrouter-gracetime", required_argument, 0, "retry connections to dead static nodes after the specified amount of seconds", uwsgi_opt_set_int, &urr.cr.static_node_gracetime, 0},
{"rawrouter-events", required_argument, 0, "set the maximum number of concurrent events", uwsgi_opt_set_int, &urr.cr.nevents, 0},
{"rawrouter-max-retries", required_argument, 0, "set the maximum number of retries/fallbacks to other nodes", uwsgi_opt_set_int, &urr.cr.max_retries, 0},
{"rawrouter-quiet", required_argument, 0, "do not report failed connections to instances", uwsgi_opt_true, &urr.cr.quiet, 0},
{"rawrouter-cheap", no_argument, 0, "run the rawrouter in cheap mode", uwsgi_opt_true, &urr.cr.cheap, 0},
{"rawrouter-subscription-server", required_argument, 0, "run the rawrouter subscription server on the spcified address", uwsgi_opt_corerouter_ss, &urr, 0},
{"rawrouter-subscription-slot", required_argument, 0, "*** deprecated ***", uwsgi_opt_deprecated, (void *) "useless thanks to the new implementation", 0},
{"rawrouter-timeout", required_argument, 0, "set rawrouter timeout", uwsgi_opt_set_int, &urr.cr.socket_timeout, 0},
{"rawrouter-stats", required_argument, 0, "run the rawrouter stats server", uwsgi_opt_set_str, &urr.cr.stats_server, 0},
{"rawrouter-stats-server", required_argument, 0, "run the rawrouter stats server", uwsgi_opt_set_str, &urr.cr.stats_server, 0},
{"rawrouter-ss", required_argument, 0, "run the rawrouter stats server", uwsgi_opt_set_str, &urr.cr.stats_server, 0},
{"rawrouter-harakiri", required_argument, 0, "enable rawrouter harakiri", uwsgi_opt_set_int, &urr.cr.harakiri, 0},
{"rawrouter-xclient", no_argument, 0, "use the xclient protocol to pass the client addres", uwsgi_opt_true, &urr.xclient, 0},
{"rawrouter-harakiri", required_argument, 0, "enable rawrouter harakiri", uwsgi_opt_set_int, &urr.cr.harakiri, 0 },
{0, 0, 0, 0, 0, 0, 0},
};
ssize_t rr_instance_read(struct corerouter_session *);
ssize_t rr_read(struct corerouter_session *);
// write to backend
ssize_t rr_instance_write(struct corerouter_session * cs) {
ssize_t len = write(cs->instance_fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos);
if (len < 0) {
cr_try_again;
uwsgi_error("fr_instance_write()");
return -1;
}
cs->buffer_pos += len;
// the chunk has been sent, start (again) reading from client and instance
if (cs->buffer_pos == (ssize_t) cs->buffer_len) {
uwsgi_cr_hook_instance_write(cs, NULL);
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
uwsgi_cr_hook_read(cs, rr_read);
}
return len;
}
// write to client
ssize_t rr_write(struct corerouter_session * cs) {
ssize_t len = write(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos);
if (len < 0) {
cr_try_again;
uwsgi_error("rr_write()");
return -1;
}
cs->buffer_pos += len;
// ok this response chunk is sent, let's wait for another one
if (cs->buffer_pos == (ssize_t) cs->buffer_len) {
uwsgi_cr_hook_write(cs, NULL);
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
}
return len;
}
ssize_t rr_instance_read(struct corerouter_session * cs) {
ssize_t len = read(cs->instance_fd, cs->buffer->buf, cs->buffer->len);
if (len < 0) {
cr_try_again;
uwsgi_error("rr_instance_read()");
return -1;
}
// end of the response
if (len == 0) {
return 0;
}
cs->buffer_pos = 0;
cs->buffer_len = len;
// ok stop reading from the instance, and start writing to the client
uwsgi_cr_hook_instance_read(cs, NULL);
uwsgi_cr_hook_write(cs, rr_write);
return len;
}
ssize_t rr_xclient_write(struct corerouter_session *);
ssize_t rr_xclient_read(struct corerouter_session * cs) {
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
cs->buffer_len = cs->buffer->len;
ssize_t len = read(cs->instance_fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos);
if (len < 0) {
cr_try_again;
uwsgi_error("rr_xclient_read()");
return -1;
}
if (len == 0) return 0;
char *ptr = cs->buffer->buf + cs->buffer_pos;
ssize_t i;
for(i=0;i<len;i++) {
if (rr->xclient_rn == 1) {
if (ptr[i] != '\n') {
return -1;
}
// banner received
cs->buffer_pos = len - (i+1);
uwsgi_cr_hook_instance_read(cs, NULL);
uwsgi_cr_hook_instance_write(cs, rr_xclient_write);
return len;
}
else if (ptr[i] == '\r') {
rr->xclient_rn = 1;
}
}
cs->buffer_pos += len;
return len;
}
ssize_t rr_xclient_write(struct corerouter_session * cs) {
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
ssize_t len = write(cs->instance_fd, rr->xclient + rr->xclient_pos, rr->xclient_len - rr->xclient_pos);
if (len < 0) {
cr_try_again;
uwsgi_error("rr_xclient_write()");
return -1;
}
rr->xclient_pos += len;
if (rr->xclient_pos == (ssize_t) rr->xclient_len) {
uwsgi_cr_hook_instance_write(cs, NULL);
if (cs->buffer_pos > 0) {
// send remaining data...
uwsgi_cr_hook_write(cs, rr_write);
}
else {
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
uwsgi_cr_hook_read(cs, rr_read);
}
}
return len;
}
ssize_t rr_instance_connected(struct corerouter_session * cs) {
cs->connecting = 0;
socklen_t solen = sizeof(int);
// first check for errors
if (getsockopt(cs->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&cs->soopt), &solen) < 0) {
uwsgi_error("rr_instance_connected()/getsockopt()");
cs->instance_failed = 1;
return -1;
}
if (cs->soopt) {
cs->instance_failed = 1;
return -1;
}
cs->buffer_pos = 0;
// ok instance is connected, begin...
if (cs->static_node) cs->static_node->custom2++;
if (cs->un) cs->un->requests++;
uwsgi_cr_hook_instance_write(cs, NULL);
if (urr.xclient) {
uwsgi_cr_hook_instance_read(cs, rr_xclient_read);
return 1;
}
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
uwsgi_cr_hook_read(cs, rr_read);
// return a value > 0
return 1;
}
ssize_t rr_read(struct corerouter_session * cs) {
ssize_t len = read(cs->fd, cs->buffer->buf, cs->buffer->len);
if (len < 0) {
cr_try_again;
uwsgi_error("rr_recv()");
return -1;
}
if (len == 0) return 0;
cs->buffer_pos = 0;
cs->buffer_len = len;
uwsgi_cr_hook_read(cs, NULL);
uwsgi_cr_hook_instance_read(cs, NULL);
uwsgi_cr_hook_instance_write(cs, rr_instance_write);
return len;
}
int rr_retry(struct uwsgi_corerouter *ucr, struct corerouter_session *cs) {
if (cs->instance_address_len > 0) goto retry;
if (ucr->mapper(ucr, cs)) {
cs->instance_failed = 1;
return -1;
}
if (cs->instance_address_len == 0) {
cs->instance_failed = 1;
return -1;
}
retry:
// start async connect
cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1);
if (cs->instance_fd < 0) {
cs->instance_failed = 1;
cs->soopt = errno;
return -1;
}
// map the instance
cs->corerouter->cr_table[cs->instance_fd] = cs;
// wait for connection
cs->connecting = 1;
// wait for connection
uwsgi_cr_hook_instance_write(cs, rr_instance_connected);
return 0;
}
void rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) {
// use the address as hostname
cs->hostname = cs->ugs->name;
cs->hostname_len = cs->ugs->name_len;
if (sa && sa->sa_family == AF_INET) {
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
rr->ip_addr = ((struct sockaddr_in *) sa)->sin_addr.s_addr;
if (urr.xclient) {
if (!inet_ntop(AF_INET, &rr->ip_addr, rr->xclient+13, INET_ADDRSTRLEN)) {
uwsgi_error("rawrouter_alloc_session() -> inet_ntop()");
cs->instance_failed = 1;
return;
}
// fix string
size_t ip_addr_len = strlen(rr->xclient+13);
memcpy(rr->xclient,"XCLIENT ADDR=", 13);
rr->xclient[13+ip_addr_len] = '\r';
rr->xclient[13+ip_addr_len+1] = '\n';
rr->xclient_len = 13 + ip_addr_len + 2;
}
}
// the mapper hook
if (ucr->mapper(ucr, cs)) {
cs->instance_failed = 1;
return;
}
if (cs->instance_address_len == 0) {
cs->instance_failed = 1;
return;
}
// ok, now we could retry
cs->retry = rr_retry;
// start async connect
cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1);
if (cs->instance_fd < 0) {
cs->instance_failed = 1;
cs->soopt = errno;
return;
}
// map the instance
cs->corerouter->cr_table[cs->instance_fd] = cs;
// wait for connection
cs->connecting = 1;
uwsgi_cr_hook_instance_write(cs, rr_instance_connected);
}
int rawrouter_init() {
urr.cr.session_size = sizeof(struct rawrouter_session);
urr.cr.switch_events = uwsgi_rawrouter_switch_events;
urr.cr.alloc_session = rawrouter_alloc_session;
uwsgi_corerouter_init((struct uwsgi_corerouter *) &urr);
uwsgi_corerouter_init((struct uwsgi_corerouter *) &urr);
return 0;
}
-16
View File
@@ -1,16 +0,0 @@
#include "../corerouter/cr.h"
struct uwsgi_rawrouter {
struct uwsgi_corerouter cr;
};
struct rawrouter_session {
struct corerouter_session crs;
};
void uwsgi_rawrouter_switch_events(struct uwsgi_corerouter *, struct corerouter_session *, int interesting_fd);
-145
View File
@@ -1,145 +0,0 @@
#include "../../uwsgi.h"
#include "rr.h"
extern struct uwsgi_server uwsgi;
extern struct uwsgi_rawrouter urr;
void uwsgi_rawrouter_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_session *cs, int interesting_fd) {
socklen_t solen = sizeof(int);
ssize_t len;
char buf[8192];
switch (cs->status) {
case COREROUTER_STATUS_RECV_HDR:
#ifdef UWSGI_EVENT_USE_PORT
event_queue_add_fd_read(ucr->queue, cs->fd);
#endif
// use the address as hostname
cs->hostname = cs->ugs->name;
cs->hostname_len = cs->ugs->name_len;
// the mapper hook
if (ucr->mapper(ucr, cs))
break;
// no address found
if (!cs->instance_address_len) {
// if fallback nodes are configured, trigger them
if (ucr->fallback) {
cs->instance_failed = 1;
}
corerouter_close_session(ucr, cs);
break;
}
cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1);
if (cs->instance_fd < 0) {
cs->instance_failed = 1;
cs->soopt = errno;
corerouter_close_session(ucr, cs);
break;
}
cs->status = COREROUTER_STATUS_CONNECTING;
ucr->cr_table[cs->instance_fd] = cs;
event_queue_add_fd_write(ucr->queue, cs->instance_fd);
break;
case COREROUTER_STATUS_CONNECTING:
if (interesting_fd == cs->instance_fd) {
if (getsockopt(cs->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&cs->soopt), &solen) < 0) {
uwsgi_error("getsockopt()");
cs->instance_failed = 1;
corerouter_close_session(ucr, cs);
break;
}
if (cs->soopt) {
cs->instance_failed = 1;
corerouter_close_session(ucr, cs);
break;
}
// increment node requests counter
if (cs->un)
cs->un->requests++;
event_queue_fd_write_to_read(ucr->queue, cs->instance_fd);
cs->status = COREROUTER_STATUS_RESPONSE;
}
break;
case COREROUTER_STATUS_RESPONSE:
// data from instance
if (interesting_fd == cs->instance_fd) {
len = recv(cs->instance_fd, buf, 8192, 0);
#ifdef UWSGI_EVENT_USE_PORT
event_queue_add_fd_read(ucr->queue, cs->instance_fd);
#endif
if (len <= 0) {
if (len < 0)
uwsgi_error("recv()");
corerouter_close_session(ucr, cs);
break;
}
len = send(cs->fd, buf, len, 0);
if (len <= 0) {
if (len < 0)
uwsgi_error("send()");
corerouter_close_session(ucr, cs);
break;
}
// update transfer statistics
if (cs->un)
cs->un->transferred += len;
}
// body from client
else if (interesting_fd == cs->fd) {
//uwsgi_log("receiving body...\n");
len = recv(cs->fd, buf, 8192, 0);
#ifdef UWSGI_EVENT_USE_PORT
event_queue_add_fd_read(ucr->queue, cs->fd);
#endif
if (len <= 0) {
if (len < 0)
uwsgi_error("recv()");
corerouter_close_session(ucr, cs);
break;
}
len = send(cs->instance_fd, buf, len, 0);
if (len <= 0) {
if (len < 0)
uwsgi_error("send()");
corerouter_close_session(ucr, cs);
break;
}
}
break;
// fallback to destroy !!!
default:
uwsgi_log("unknown event: closing session\n");
corerouter_close_session(ucr, cs);
break;
}
}
+1 -1
View File
@@ -6,4 +6,4 @@ LIBS = []
REQUIRES = ['corerouter']
GCC_LIST = ['rawrouter', 'rr_events']
GCC_LIST = ['rawrouter']
+18 -1
View File
@@ -12,13 +12,24 @@ int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_route *u
// get the http address from the route
char *addr = ur->data;
char *uri = NULL;
uint16_t uri_len = 0;
if (ur->data3_len) {
uri = uwsgi_regexp_apply_ovec(wsgi_req->uri, wsgi_req->uri_len, ur->data3, ur->data3_len, ur->ovector, ur->ovn);
uri_len = strlen(uri);
}
// convert the wsgi_request to an http proxy request
struct uwsgi_buffer *ub = uwsgi_to_http(wsgi_req, ur->data2, ur->data2_len);
struct uwsgi_buffer *ub = uwsgi_to_http(wsgi_req, ur->data2, ur->data2_len, uri, uri_len);
if (!ub) {
if (uri) free(uri);
uwsgi_log("unable to generate http request for %s\n", addr);
return UWSGI_ROUTE_NEXT;
}
if (uri) free(uri);
// ok now if have offload threads, directly use them
if (wsgi_req->socket->can_offload) {
if (!uwsgi_offload_request_net_do(wsgi_req, addr, ub)) {
@@ -86,6 +97,12 @@ int uwsgi_router_http(struct uwsgi_route *ur, char *args) {
*comma = 0;
ur->data_len = strlen(ur->data);
ur->data2 = comma+1;
comma = strchr(ur->data2, ',');
if (comma) {
*comma = 0;
ur->data3 = comma+1;
ur->data3_len = strlen(ur->data3);
}
ur->data2_len = strlen(ur->data2);
}
return 0;
+17 -29
View File
@@ -5,6 +5,8 @@ extern struct uwsgi_server uwsgi;
int uwsgi_routing_func_rewrite(struct wsgi_request *wsgi_req, struct uwsgi_route *ur) {
char *tmp_qs = NULL;
char **subject = (char **) (((char *)(wsgi_req))+ur->subject);
uint16_t *subject_len = (uint16_t *) (((char *)(wsgi_req))+ur->subject_len);
@@ -18,28 +20,26 @@ int uwsgi_routing_func_rewrite(struct wsgi_request *wsgi_req, struct uwsgi_route
path_info_len = query_string - path_info;
query_string++;
query_string_len = strlen(query_string);
if (wsgi_req->query_string_len > 0) {
tmp_qs = uwsgi_concat4n(query_string, query_string_len, "&", 1, wsgi_req->query_string, wsgi_req->query_string_len, "", 0);
query_string = tmp_qs;
query_string_len = strlen(query_string);
}
}
// over engineering, could be requiredin the future...
else {
query_string = "";
if (wsgi_req->query_string_len > 0) {
query_string = wsgi_req->query_string;
query_string_len = wsgi_req->query_string_len;
}
else {
query_string = "";
}
}
char *ptr = uwsgi_req_append(wsgi_req, "PATH_INFO", 9, path_info, path_info_len);
if (!ptr) goto clear;
// fill iovec
if (wsgi_req->var_cnt + 2 >= uwsgi.vec_size - (4 + 1)) {
uwsgi_log("not enough io vectors for rewriting url\n");
goto clear;
}
wsgi_req->hvec[wsgi_req->var_cnt].iov_base = ptr - (2 + 9);
wsgi_req->hvec[wsgi_req->var_cnt].iov_len = 9;
wsgi_req->var_cnt++;
wsgi_req->hvec[wsgi_req->var_cnt].iov_base = ptr;
wsgi_req->hvec[wsgi_req->var_cnt].iov_len = path_info_len;
wsgi_req->var_cnt++;
// set new path_info
wsgi_req->path_info = ptr;
wsgi_req->path_info_len = path_info_len;
@@ -47,31 +47,19 @@ int uwsgi_routing_func_rewrite(struct wsgi_request *wsgi_req, struct uwsgi_route
ptr = uwsgi_req_append(wsgi_req, "QUERY_STRING", 12, query_string, query_string_len);
if (!ptr) goto clear;
// fill iovec
if (wsgi_req->var_cnt + 2 >= uwsgi.vec_size - (4 + 1)) {
uwsgi_log("not enough io vectors for rewriting url\n");
goto clear;
}
wsgi_req->hvec[wsgi_req->var_cnt].iov_base = ptr - (2 + 12);
wsgi_req->hvec[wsgi_req->var_cnt].iov_len = 12;
wsgi_req->var_cnt++;
wsgi_req->hvec[wsgi_req->var_cnt].iov_base = ptr;
wsgi_req->hvec[wsgi_req->var_cnt].iov_len = query_string_len;
wsgi_req->var_cnt++;
// set new query_string
wsgi_req->query_string = ptr;
wsgi_req->query_string_len = query_string_len;
free(path_info);
if (tmp_qs) free(tmp_qs);
if (ur->custom)
return UWSGI_ROUTE_CONTINUE;
return UWSGI_ROUTE_NEXT;
clear:
free(path_info);
if (tmp_qs) free(tmp_qs);
return UWSGI_ROUTE_BREAK;
}
+19 -3
View File
@@ -393,14 +393,19 @@ void uwsgi_httpize_var(char *buf, size_t len) {
}
}
struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *wsgi_req, char *host, uint16_t host_len) {
struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *wsgi_req, char *host, uint16_t host_len, char *uri, uint16_t uri_len) {
struct uwsgi_buffer *ub = uwsgi_buffer_new(4096);
if (uwsgi_buffer_append(ub, wsgi_req->method, wsgi_req->method_len)) goto clear;
if (uwsgi_buffer_append(ub, " ", 1)) goto clear;
if (uwsgi_buffer_append(ub, wsgi_req->uri, wsgi_req->uri_len)) goto clear;
if (uri_len && uri) {
if (uwsgi_buffer_append(ub, uri, uri_len)) goto clear;
}
else {
if (uwsgi_buffer_append(ub, wsgi_req->uri, wsgi_req->uri_len)) goto clear;
}
if (uwsgi_buffer_append(ub, " HTTP/1.0\r\n", 11)) goto clear;
@@ -447,11 +452,22 @@ next:
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
}
if (wsgi_req->content_type_len > 0) {
if (uwsgi_buffer_append(ub, "Content-Type: ", 14)) goto clear;
if (uwsgi_buffer_append(ub, wsgi_req->content_type, wsgi_req->content_type_len)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
}
if (wsgi_req->post_cl > 0) {
if (uwsgi_buffer_append(ub, "Content-Length: ", 16)) goto clear;
if (uwsgi_buffer_num64(ub, (int64_t) wsgi_req->post_cl)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
}
// append required headers
if (uwsgi_buffer_append(ub, "Connection: close\r\n", 19)) goto clear;
if (uwsgi_buffer_append(ub, "X-Forwarded-For: ", 17)) goto clear;
if (x_forwarded_for_len > 0) {
if (uwsgi_buffer_append(ub, x_forwarded_for, x_forwarded_for_len)) goto clear;
if (uwsgi_buffer_append(ub, ", ", 2)) goto clear;
+11
View File
@@ -19,6 +19,14 @@ my $two = sub {
print "two\n";
};
my $four = sub {
my $signum = shift;
print "i am signal ".$signum."\n" ;
};
uwsgi::register_signal(17, '', $four);
uwsgi::register_signal(30, '', $two);
my $three = sub {
my $env = shift;
sleep(1);
@@ -27,6 +35,9 @@ my $three = sub {
my $app = sub {
my $env = shift;
uwsgi::signal(17);
uwsgi::signal(30);
if ($env->{'psgix.cleanup'}) {
print "cleanup supported\n";
push @{$env->{'psgix.cleanup.handlers'}}, $one;
+19 -7
View File
@@ -37,7 +37,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) || !uwsgi_startswith(file, "section://", 10))
#define uwsgi_check_scheme(file) (!uwsgi_startswith(file, "emperor://", 10) || !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
@@ -308,6 +308,8 @@ struct uwsgi_string_list {
char *value;
size_t len;
uint64_t custom;
uint64_t custom2;
void *custom_ptr;
struct uwsgi_string_list *next;
};
@@ -441,6 +443,8 @@ struct uwsgi_gateway_socket {
char *port;
int port_len;
int no_defer;
void *data;
// this requires UDP
int subscription;
@@ -638,6 +642,7 @@ struct uwsgi_socket {
void *ctx;
int queue;
int no_defer;
int auto_port;
// true if connection must be initialized for each core
@@ -846,6 +851,9 @@ struct uwsgi_route {
void *data2;
size_t data2_len;
void *data3;
size_t data3_len;
// 64bit value for custom usage
uint64_t custom;
@@ -1214,6 +1222,7 @@ struct uwsgi_server {
// enable threads
int has_threads;
int no_threads_wait;
// default app id
int default_app;
@@ -1524,10 +1533,6 @@ struct uwsgi_server {
int check_static_docroot;
// linked list for offloaded requests
struct uwsgi_offload_request *offload_requests_head;
struct uwsgi_offload_request *offload_requests_tail;
char *daemonize;
char *daemonize2;
int do_not_change_umask;
@@ -2807,7 +2812,7 @@ int uwsgi_get_shared_socket_num(struct uwsgi_socket *);
#ifdef __linux__
void uwsgi_set_cgroup(void);
long uwsgi_num_from_file(char *);
long uwsgi_num_from_file(char *, int);
#endif
void uwsgi_add_sockets_to_queue(int, int);
@@ -2992,6 +2997,7 @@ void uwsgi_opt_set_str(char *, char *, void *);
void uwsgi_opt_set_logger(char *, char *, void *);
void uwsgi_opt_set_str_spaced(char *, char *, void *);
void uwsgi_opt_add_string_list(char *, char *, void *);
void uwsgi_opt_add_addr_list(char *, char *, void *);
void uwsgi_opt_add_dyn_dict(char *, char *, void *);
#ifdef UWSGI_PCRE
void uwsgi_opt_pcre_jit(char *, char *, void *);
@@ -3362,9 +3368,10 @@ int uwsgi_buffer_append(struct uwsgi_buffer *, char *, size_t);
int uwsgi_buffer_fix(struct uwsgi_buffer *, size_t);
int uwsgi_buffer_ensure(struct uwsgi_buffer *, size_t);
void uwsgi_buffer_destroy(struct uwsgi_buffer *);
int uwsgi_buffer_num64(struct uwsgi_buffer *, int64_t);
void uwsgi_httpize_var(char *, size_t);
struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *, char *, uint16_t);
struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *, char *, uint16_t, char *, uint16_t);
ssize_t uwsgi_pipe(int, int, int);
ssize_t uwsgi_pipe_sized(int, int, size_t, int);
@@ -3404,6 +3411,9 @@ struct uwsgi_thread {
uint64_t custom1;
uint64_t custom2;
uint64_t custom3;
// linked list for offloaded requests
struct uwsgi_offload_request *offload_requests_head;
struct uwsgi_offload_request *offload_requests_tail;
void (*func)(struct uwsgi_thread *);
};
struct uwsgi_thread *uwsgi_thread_new(void (*)(struct uwsgi_thread *));
@@ -3475,6 +3485,8 @@ void uwsgi_user_lock(int);
void uwsgi_user_unlock(int);
void simple_loop_run_int(int);
int uwsgi_valid_fd(int);
void uwsgi_close_all_fds(void);
void uwsgi_check_emperor(void);
#ifdef UWSGI_AS_SHARED_LIBRARY
+17 -8
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '1.4-rc2'
uwsgi_version = '1.4.4'
import os
import re
@@ -33,15 +33,20 @@ if not GCC:
CPP = os.environ.get('CPP', 'cpp')
CPUCOUNT = 1
try:
import multiprocessing
CPUCOUNT = multiprocessing.cpu_count()
CPUCOUNT = int(os.environ.get('CPUCOUNT', -1))
except:
CPUCOUNT = -1
if CPUCOUNT < 1:
try:
CPUCOUNT = int(os.sysconf('SC_NPROCESSORS_ONLN'))
import multiprocessing
CPUCOUNT = multiprocessing.cpu_count()
except:
pass
try:
CPUCOUNT = int(os.sysconf('SC_NPROCESSORS_ONLN'))
except:
CPUCOUNT = 1
binary_list = []
@@ -340,12 +345,16 @@ def build_uwsgi(uc, print_only=False):
pass
for cfile in up.GCC_LIST:
if not cfile.endswith('.a'):
if cfile.endswith('.a'):
gcc_list.append(cfile)
elif not cfile.endswith('.c') and not cfile.endswith('.cc') and not cfile.endswith('.m'):
compile(' '.join(uniq_warnings(p_cflags)), last_cflags_ts,
path + '/' + cfile + '.o', path + '/' + cfile + '.c')
gcc_list.append('%s/%s' % (path, cfile))
else:
gcc_list.append(cfile)
compile(' '.join(uniq_warnings(p_cflags)), last_cflags_ts,
path + '/' + cfile + '.o', path + '/' + cfile)
gcc_list.append('%s/%s' % (path, cfile))
libs += up.LIBS