mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-10-04 16:08:31 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
87ef444ea6 | ||
|
|
714be76f9b | ||
|
|
af0f0da8c2 | ||
|
|
8407372801 | ||
|
|
d19c4e84fc | ||
|
|
3342852a9a | ||
|
|
586a02d849 | ||
|
|
188dad744a | ||
|
|
6cf3bbea69 | ||
|
|
04736d371e | ||
|
|
cc9151f89e | ||
|
|
e13734ae6c | ||
|
|
cf62b03e33 | ||
|
|
b2f292c2da | ||
|
|
b96b6a24a6 | ||
|
|
b2a1237781 | ||
|
|
a48af2e226 | ||
|
|
33446c3179 | ||
|
|
242e01e37a | ||
|
|
41eba2b486 | ||
|
|
24763fdbd2 | ||
|
|
951cbf3eba | ||
|
|
8759b66206 |
@@ -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 ***
|
||||
|
||||
|
||||
+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, 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
|
||||
|
||||
+2
-2
@@ -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;
|
||||
|
||||
@@ -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, ':');
|
||||
|
||||
+19
-18
@@ -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
|
||||
|
||||
+16
-1
@@ -1082,7 +1082,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 +1092,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 +1129,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 +1664,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);
|
||||
|
||||
|
||||
+13
-3
@@ -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) {
|
||||
@@ -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++) {
|
||||
|
||||
+11
-8
@@ -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
|
||||
@@ -2684,15 +2691,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;
|
||||
}
|
||||
@@ -3101,6 +3103,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;
|
||||
}
|
||||
@@ -4798,7 +4801,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)) {
|
||||
|
||||
+32
-11
@@ -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},
|
||||
@@ -675,6 +677,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 +697,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 +740,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 +774,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);
|
||||
@@ -3294,7 +3312,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) {
|
||||
|
||||
@@ -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()");
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
@@ -860,7 +859,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 +895,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 +923,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 +970,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;
|
||||
|
||||
@@ -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 *));
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
+25
-3
@@ -399,9 +399,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') {
|
||||
@@ -518,6 +518,14 @@ ssize_t hr_read_ssl_body(struct corerouter_session * cs) {
|
||||
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 )) return -1;
|
||||
if (SSL_read(hs->ssl, hs->post_buf->buf + ret, ret2) != ret2) {
|
||||
return -1;
|
||||
}
|
||||
ret += ret2;
|
||||
}
|
||||
hs->post_buf_len = ret;
|
||||
hs->post_buf_pos = 0;
|
||||
uwsgi_cr_hook_read(cs, NULL);
|
||||
@@ -762,6 +770,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 +789,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);
|
||||
}
|
||||
|
||||
@@ -856,6 +874,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 +902,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
|
||||
@@ -1024,6 +1045,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;
|
||||
}
|
||||
@@ -1117,7 +1139,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;
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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 *);
|
||||
|
||||
@@ -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;
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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;}
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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;
|
||||
|
||||
}
|
||||
}
|
||||
@@ -6,4 +6,4 @@ LIBS = []
|
||||
|
||||
REQUIRES = ['corerouter']
|
||||
|
||||
GCC_LIST = ['rawrouter', 'rr_events']
|
||||
GCC_LIST = ['rawrouter']
|
||||
|
||||
@@ -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,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;
|
||||
|
||||
@@ -308,6 +308,7 @@ struct uwsgi_string_list {
|
||||
char *value;
|
||||
size_t len;
|
||||
uint64_t custom;
|
||||
uint64_t custom2;
|
||||
struct uwsgi_string_list *next;
|
||||
};
|
||||
|
||||
@@ -441,6 +442,8 @@ struct uwsgi_gateway_socket {
|
||||
char *port;
|
||||
int port_len;
|
||||
|
||||
int no_defer;
|
||||
|
||||
void *data;
|
||||
// this requires UDP
|
||||
int subscription;
|
||||
@@ -638,6 +641,7 @@ struct uwsgi_socket {
|
||||
void *ctx;
|
||||
|
||||
int queue;
|
||||
int no_defer;
|
||||
|
||||
int auto_port;
|
||||
// true if connection must be initialized for each core
|
||||
@@ -1214,6 +1218,7 @@ struct uwsgi_server {
|
||||
|
||||
// enable threads
|
||||
int has_threads;
|
||||
int no_threads_wait;
|
||||
|
||||
// default app id
|
||||
int default_app;
|
||||
@@ -1524,10 +1529,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;
|
||||
@@ -3404,6 +3405,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 *));
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
# uWSGI build system
|
||||
|
||||
uwsgi_version = '1.4-rc2'
|
||||
uwsgi_version = '1.4.2'
|
||||
|
||||
import os
|
||||
import re
|
||||
|
||||
Reference in New Issue
Block a user