Compare commits

..
13 Commits
18 changed files with 454 additions and 247 deletions
+3 -1
View File
@@ -9,11 +9,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
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
+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;
+12 -2
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) {
@@ -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++) {
+8
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
@@ -3101,6 +3108,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;
}
+9 -1
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},
@@ -759,6 +761,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 +3299,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) {
+76 -51
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;
}
@@ -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);
@@ -896,6 +895,9 @@ 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);
@@ -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;
+6 -2
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 *);
+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);
}
+14
View File
@@ -762,6 +762,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 +781,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 +866,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 +894,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 +1037,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;
}
+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']
+5
View File
@@ -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;
+1 -1
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '1.4-rc2'
uwsgi_version = '1.4'
import os
import re