From 990e7241e601d0e1b32ae0ee81e3d064ed306161 Mon Sep 17 00:00:00 2001 From: Unbit Date: Tue, 29 Jan 2013 09:08:41 +0100 Subject: [PATCH] added the sslrouter and various corerouters micro-optimizations --- buildconf/base.ini | 2 +- plugins/corerouter/corerouter.c | 26 ++- plugins/corerouter/cr.h | 3 +- plugins/fastrouter/fastrouter.c | 30 +-- plugins/http/https.c | 1 - plugins/rawrouter/rawrouter.c | 28 +-- plugins/sslrouter/sslrouter.c | 330 +++++++++++++++++++++++++++++++ plugins/sslrouter/uwsgiplugin.py | 9 + 8 files changed, 389 insertions(+), 40 deletions(-) create mode 100644 plugins/sslrouter/sslrouter.c create mode 100644 plugins/sslrouter/uwsgiplugin.py diff --git a/buildconf/base.ini b/buildconf/base.ini index a247759b..a8cd89f0 100644 --- a/buildconf/base.ini +++ b/buildconf/base.ini @@ -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, rawrouter, router_static +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, router_static, sslrouter as_shared_library = false locking = auto diff --git a/plugins/corerouter/corerouter.c b/plugins/corerouter/corerouter.c index edd707de..eb2ee37d 100644 --- a/plugins/corerouter/corerouter.c +++ b/plugins/corerouter/corerouter.c @@ -416,9 +416,15 @@ struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *ucr, st return cr_add_timeout(ucr, peer); } -static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) { +struct uwsgi_rb_timer *corerouter_reset_timeout_fast(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer, time_t now) { + cr_del_timeout(ucr, peer); + return cr_add_timeout_fast(ucr, peer, now); +} - uint64_t current = (uint64_t) uwsgi_now(); + +static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr, time_t now) { + + uint64_t current = (uint64_t) now; struct uwsgi_rb_timer *urbt; struct corerouter_peer *peer; @@ -695,15 +701,17 @@ void uwsgi_corerouter_loop(int id, void *data) { for (;;) { + time_t now = uwsgi_now(); + // set timeouts and harakiri min_timeout = uwsgi_min_rb_timer(ucr->timeouts, NULL); if (min_timeout == NULL) { delta = -1; } else { - delta = min_timeout->value - uwsgi_now(); + delta = min_timeout->value - now; if (delta <= 0) { - corerouter_expire_timeouts(ucr); + corerouter_expire_timeouts(ucr, now); delta = 0; } } @@ -715,12 +723,14 @@ void uwsgi_corerouter_loop(int id, void *data) { // wait for events nevents = event_queue_wait_multi(ucr->queue, delta, events, ucr->nevents); + now = uwsgi_now(); + if (uwsgi.master_process && ucr->harakiri > 0) { - ushared->gateways_harakiri[id] = uwsgi_now() + ucr->harakiri; + ushared->gateways_harakiri[id] = now + ucr->harakiri; } if (nevents == 0) { - corerouter_expire_timeouts(ucr); + corerouter_expire_timeouts(ucr, now); } for (i = 0; i < nevents; i++) { @@ -796,8 +806,8 @@ void uwsgi_corerouter_loop(int id, void *data) { } // set timeout (in main_peer too) - peer->timeout = corerouter_reset_timeout(ucr, peer); - peer->session->main_peer->timeout = corerouter_reset_timeout(ucr, peer->session->main_peer); + peer->timeout = corerouter_reset_timeout_fast(ucr, peer, now); + peer->session->main_peer->timeout = corerouter_reset_timeout_fast(ucr, peer->session->main_peer, now); ssize_t (*hook)(struct corerouter_peer *) = NULL; diff --git a/plugins/corerouter/cr.h b/plugins/corerouter/cr.h index 37523d4a..39d3147d 100644 --- a/plugins/corerouter/cr.h +++ b/plugins/corerouter/cr.h @@ -3,7 +3,8 @@ #define COREROUTER_STATUS_RECV_HDR 2 #define COREROUTER_STATUS_RESPONSE 3 -#define cr_add_timeout(u, x) uwsgi_add_rb_timer(u->timeouts, time(NULL)+u->socket_timeout, x) +#define cr_add_timeout(u, x) uwsgi_add_rb_timer(u->timeouts, uwsgi_now()+u->socket_timeout, x) +#define cr_add_timeout_fast(u, x, t) uwsgi_add_rb_timer(u->timeouts, t+u->socket_timeout, x) #define cr_del_timeout(u, x) uwsgi_del_rb_timer(u->timeouts, x->timeout); free(x->timeout); #define cr_try_again if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) {\ diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c index 04d94328..2664d43d 100644 --- a/plugins/fastrouter/fastrouter.c +++ b/plugins/fastrouter/fastrouter.c @@ -7,7 +7,7 @@ #include "../../uwsgi.h" #include "../corerouter/cr.h" -struct uwsgi_fastrouter { +static struct uwsgi_fastrouter { struct uwsgi_corerouter cr; } ufr; @@ -18,7 +18,7 @@ struct fastrouter_session { int has_key; }; -struct uwsgi_option fastrouter_options[] = { +static struct uwsgi_option fastrouter_options[] = { {"fastrouter", required_argument, 0, "run the fastrouter on the specified port", uwsgi_opt_corerouter, &ufr, 0}, {"fastrouter-processes", required_argument, 0, "prefork the specified number of fastrouter processes", uwsgi_opt_set_int, &ufr.cr.processes, 0}, {"fastrouter-workers", required_argument, 0, "prefork the specified number of fastrouter processes", uwsgi_opt_set_int, &ufr.cr.processes, 0}, @@ -53,7 +53,7 @@ struct uwsgi_option fastrouter_options[] = { {0, 0, 0, 0, 0, 0, 0}, }; -void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) { +static void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) { struct corerouter_peer *peer = (struct corerouter_peer *) data; struct fastrouter_session *fr = (struct fastrouter_session *) peer->session; @@ -80,7 +80,7 @@ void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, voi } // writing client body to the instance -ssize_t fr_instance_write_body(struct corerouter_peer *peer) { +static ssize_t fr_instance_write_body(struct corerouter_peer *peer) { ssize_t len = cr_write(peer, "fr_instance_write_body()"); // end on empty write if (!len) return 0; @@ -97,7 +97,7 @@ ssize_t fr_instance_write_body(struct corerouter_peer *peer) { // read client body -ssize_t fr_read_body(struct corerouter_peer *main_peer) { +static ssize_t fr_read_body(struct corerouter_peer *main_peer) { ssize_t len = cr_read(main_peer, "fr_read_body()"); if (!len) return 0; @@ -109,7 +109,7 @@ ssize_t fr_read_body(struct corerouter_peer *main_peer) { } // write to the client -ssize_t fr_write(struct corerouter_peer *main_peer) { +static ssize_t fr_write(struct corerouter_peer *main_peer) { ssize_t len = cr_write(main_peer, "fr_write()"); // end on empty write if (!len) return 0; @@ -125,7 +125,7 @@ ssize_t fr_write(struct corerouter_peer *main_peer) { } // data from instance -ssize_t fr_instance_read(struct corerouter_peer *peer) { +static ssize_t fr_instance_read(struct corerouter_peer *peer) { ssize_t len = cr_read(peer, "fr_instance_read()"); if (!len) return 0; @@ -138,7 +138,7 @@ ssize_t fr_instance_read(struct corerouter_peer *peer) { } // send the uwsgi request header and vars -ssize_t fr_instance_send_request(struct corerouter_peer *peer) { +static ssize_t fr_instance_send_request(struct corerouter_peer *peer) { ssize_t len = cr_write(peer, "fr_instance_send_request()"); // end on empty write if (!len) return 0; @@ -156,7 +156,7 @@ ssize_t fr_instance_send_request(struct corerouter_peer *peer) { } // instance is connected -ssize_t fr_instance_connected(struct corerouter_peer *peer) { +static ssize_t fr_instance_connected(struct corerouter_peer *peer) { cr_peer_connected(peer, "fr_instance_connected()"); @@ -173,7 +173,7 @@ ssize_t fr_instance_connected(struct corerouter_peer *peer) { } // called after receaving the uwsgi header (read vars) -ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) { +static ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) { struct uwsgi_header *uh = (struct uwsgi_header *) main_peer->in->buf; // increase buffer if needed if (uwsgi_buffer_fix(main_peer->in, uh->pktsize+4)) @@ -212,7 +212,7 @@ ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) { } // called soon after accept() -ssize_t fr_recv_uwsgi_header(struct corerouter_peer *main_peer) { +static ssize_t fr_recv_uwsgi_header(struct corerouter_peer *main_peer) { ssize_t len = cr_read_exact(main_peer, 4, "fr_recv_uwsgi_header()"); if (!len) return 0; @@ -227,7 +227,7 @@ ssize_t fr_recv_uwsgi_header(struct corerouter_peer *main_peer) { } // retry connection to the backend -int fr_retry(struct corerouter_peer *peer) { +static int fr_retry(struct corerouter_peer *peer) { struct uwsgi_corerouter *ucr = peer->session->corerouter; @@ -249,7 +249,7 @@ retry: // called when a new session is created -int fastrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) { +static int fastrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) { // set the retry hook cs->retry = fr_retry; // wait for requests... @@ -257,7 +257,7 @@ int fastrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_ return 0; } -int fastrouter_init() { +static int fastrouter_init() { ufr.cr.session_size = sizeof(struct fastrouter_session); ufr.cr.alloc_session = fastrouter_alloc_session; @@ -266,7 +266,7 @@ int fastrouter_init() { return 0; } -void fastrouter_setup() { +static void fastrouter_setup() { ufr.cr.name = uwsgi_str("uWSGI fastrouter"); ufr.cr.short_name = uwsgi_str("fastrouter"); } diff --git a/plugins/http/https.c b/plugins/http/https.c index 9c54add6..77e8acd2 100644 --- a/plugins/http/https.c +++ b/plugins/http/https.c @@ -200,7 +200,6 @@ int hr_https_add_vars(struct http_session *hr, struct uwsgi_buffer *out) { void hr_session_ssl_close(struct corerouter_session *cs) { hr_session_close(cs); struct http_session *hr = (struct http_session *) cs; - SSL_shutdown(hr->ssl); if (hr->ssl_client_dn) { OPENSSL_free(hr->ssl_client_dn); } diff --git a/plugins/rawrouter/rawrouter.c b/plugins/rawrouter/rawrouter.c index 8f113036..d72da5d3 100644 --- a/plugins/rawrouter/rawrouter.c +++ b/plugins/rawrouter/rawrouter.c @@ -7,7 +7,7 @@ #include "../../uwsgi.h" #include "../corerouter/cr.h" -struct uwsgi_rawrouter { +static struct uwsgi_rawrouter { struct uwsgi_corerouter cr; int xclient; } urr; @@ -26,7 +26,7 @@ struct rawrouter_session { size_t xclient_rn; }; -struct uwsgi_option rawrouter_options[] = { +static struct uwsgi_option rawrouter_options[] = { {"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}, @@ -64,7 +64,7 @@ struct uwsgi_option rawrouter_options[] = { }; // write to backend -ssize_t rr_instance_write(struct corerouter_peer *peer) { +static ssize_t rr_instance_write(struct corerouter_peer *peer) { ssize_t len = cr_write(peer, "rr_instance_write()"); // end on empty write if (!len) return 0; @@ -80,7 +80,7 @@ ssize_t rr_instance_write(struct corerouter_peer *peer) { } // write to client -ssize_t rr_write(struct corerouter_peer *main_peer) { +static ssize_t rr_write(struct corerouter_peer *main_peer) { ssize_t len = cr_write(main_peer, "rr_write()"); // end on empty write if (!len) return 0; @@ -96,7 +96,7 @@ ssize_t rr_write(struct corerouter_peer *main_peer) { } // read from backend -ssize_t rr_instance_read(struct corerouter_peer *peer) { +static ssize_t rr_instance_read(struct corerouter_peer *peer) { ssize_t len = cr_read(peer, "rr_instance_read()"); if (!len) return 0; @@ -109,7 +109,7 @@ ssize_t rr_instance_read(struct corerouter_peer *peer) { } // write the xclient banner -ssize_t rr_xclient_write(struct corerouter_peer *peer) { +static ssize_t rr_xclient_write(struct corerouter_peer *peer) { struct corerouter_session *cs = peer->session; struct rawrouter_session *rr = (struct rawrouter_session *) cs; ssize_t len = cr_write_buf(peer, rr->xclient, "rr_xclient_write()"); @@ -130,7 +130,7 @@ ssize_t rr_xclient_write(struct corerouter_peer *peer) { } // read the first line from the backend and skip it -ssize_t rr_xclient_read(struct corerouter_peer *peer) { +static ssize_t rr_xclient_read(struct corerouter_peer *peer) { struct corerouter_session *cs = peer->session; struct rawrouter_session *rr = (struct rawrouter_session *) cs; ssize_t len = cr_read(peer, "rr_xclient_read()"); @@ -161,7 +161,7 @@ ssize_t rr_xclient_read(struct corerouter_peer *peer) { } // the instance is connected now we cannot retry connections -ssize_t rr_instance_connected(struct corerouter_peer *peer) { +static ssize_t rr_instance_connected(struct corerouter_peer *peer) { struct corerouter_session *cs = peer->session; struct rawrouter_session *rr = (struct rawrouter_session *) cs; @@ -176,7 +176,7 @@ ssize_t rr_instance_connected(struct corerouter_peer *peer) { } // read from client -ssize_t rr_read(struct corerouter_peer *main_peer) { +static ssize_t rr_read(struct corerouter_peer *main_peer) { ssize_t len = cr_read(main_peer, "rr_read()"); if (!len) return 0; @@ -188,7 +188,7 @@ ssize_t rr_read(struct corerouter_peer *main_peer) { } // retry the connection -int rr_retry(struct corerouter_peer *peer) { +static int rr_retry(struct corerouter_peer *peer) { struct corerouter_session *cs = peer->session; struct uwsgi_corerouter *ucr = cs->corerouter; @@ -209,7 +209,7 @@ retry: return 0; } -void rr_session_close(struct corerouter_session *cs) { +static void rr_session_close(struct corerouter_session *cs) { struct rawrouter_session *rr = (struct rawrouter_session *) cs; if (rr->xclient) { uwsgi_buffer_destroy(rr->xclient); @@ -217,7 +217,7 @@ void rr_session_close(struct corerouter_session *cs) { } // allocate a new session -int rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) { +static int rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) { // set default read hook cs->main_peer->last_hook_read = rr_read; @@ -260,7 +260,7 @@ int rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_s return 0; } -int rawrouter_init() { +static int rawrouter_init() { urr.cr.session_size = sizeof(struct rawrouter_session); urr.cr.alloc_session = rawrouter_alloc_session; @@ -269,7 +269,7 @@ int rawrouter_init() { return 0; } -void rawrouter_setup() { +static void rawrouter_setup() { urr.cr.name = uwsgi_str("uWSGI rawrouter"); urr.cr.short_name = uwsgi_str("rawrouter"); } diff --git a/plugins/sslrouter/sslrouter.c b/plugins/sslrouter/sslrouter.c new file mode 100644 index 00000000..a9fb6e9a --- /dev/null +++ b/plugins/sslrouter/sslrouter.c @@ -0,0 +1,330 @@ +#ifdef UWSGI_SSL + +/* + + uWSGI sslrouter + +*/ + +#include +#include "../corerouter/cr.h" + +extern struct uwsgi_server uwsgi; + +static struct uwsgi_sslrouter { + struct uwsgi_corerouter cr; + char *ssl_session_context; + int sni; +} usr; + +struct sslrouter_session { + struct corerouter_session session; + in_addr_t ip_addr; + SSL *ssl; +}; + +static void uwsgi_opt_sslrouter(char *opt, char *value, void *cr) { + struct uwsgi_corerouter *ucr = (struct uwsgi_corerouter *) cr; + + char *s_addr = NULL; + char *s_cert = NULL; + char *s_key = NULL; + char *s_ciphers = NULL; + char *s_clientca = NULL; + + if (uwsgi_kvlist_parse(value, strlen(value), ',', '=', + "addr", &s_addr, + "cert", &s_cert, + "crt", &s_cert, + "key", &s_key, + "ciphers", &s_ciphers, + "clientca", &s_clientca, + "client_ca", &s_clientca, + NULL)) { + uwsgi_log("error parsing --sslrouter option\n"); + exit(1); + } + + if (!s_addr || !s_cert || !s_key) { + uwsgi_log("--sslrouter option needs addr, cert and key items\n"); + exit(1); + } + + struct uwsgi_gateway_socket *ugs = uwsgi_new_gateway_socket(s_addr, ucr->name); + // ok we have the socket, initialize ssl if required + if (!uwsgi.ssl_initialized) { + uwsgi_ssl_init(); + } + + // initialize ssl context + char *name = usr.ssl_session_context; + if (!name) { + name = uwsgi_concat3(ucr->short_name, "-", ugs->name); + } + + ugs->ctx = uwsgi_ssl_new_server_context(name, s_cert, s_key, s_ciphers, s_clientca); + if (!ugs->ctx) { + exit(1); + } + ucr->has_sockets++; +} + + +static struct uwsgi_option sslrouter_options[] = { + {"sslrouter", required_argument, 0, "run the sslrouter on the specified port", uwsgi_opt_sslrouter, &usr, 0}, + {"sslrouter-session-context", required_argument, 0, "set the session id context to the specified value", uwsgi_opt_set_str, &usr.ssl_session_context, 0}, + {"sslrouter-processes", required_argument, 0, "prefork the specified number of sslrouter processes", uwsgi_opt_set_int, &usr.cr.processes, 0}, + {"sslrouter-workers", required_argument, 0, "prefork the specified number of sslrouter processes", uwsgi_opt_set_int, &usr.cr.processes, 0}, + {"sslrouter-zerg", required_argument, 0, "attach the sslrouter to a zerg server", uwsgi_opt_corerouter_zerg, &usr, 0}, + {"sslrouter-use-cache", optional_argument, 0, "use uWSGI cache as hostname->server mapper for the sslrouter", uwsgi_opt_set_str, &usr.cr.use_cache, 0}, + + {"sslrouter-use-pattern", required_argument, 0, "use a pattern for sslrouter hostname->server mapping", uwsgi_opt_corerouter_use_pattern, &usr, 0}, + {"sslrouter-use-base", required_argument, 0, "use a base dir for sslrouter hostname->server mapping", uwsgi_opt_corerouter_use_base, &usr, 0}, + + {"sslrouter-fallback", required_argument, 0, "fallback to the specified node in case of error", uwsgi_opt_add_string_list, &usr.cr.fallback, 0}, + + {"sslrouter-use-cluster", no_argument, 0, "load balance to nodes subscribed to the cluster", uwsgi_opt_true, &usr.cr.use_cluster, 0}, + + {"sslrouter-use-code-string", required_argument, 0, "use code string as hostname->server mapper for the sslrouter", uwsgi_opt_corerouter_cs, &usr, 0}, + {"sslrouter-use-socket", optional_argument, 0, "forward request to the specified uwsgi socket", uwsgi_opt_corerouter_use_socket, &usr, 0}, + {"sslrouter-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, &usr.cr.static_nodes, 0}, + {"sslrouter-gracetime", required_argument, 0, "retry connections to dead static nodes after the specified amount of seconds", uwsgi_opt_set_int, &usr.cr.static_node_gracetime, 0}, + {"sslrouter-events", required_argument, 0, "set the maximum number of concusrent events", uwsgi_opt_set_int, &usr.cr.nevents, 0}, + {"sslrouter-max-retries", required_argument, 0, "set the maximum number of retries/fallbacks to other nodes", uwsgi_opt_set_int, &usr.cr.max_retries, 0}, + {"sslrouter-quiet", required_argument, 0, "do not report failed connections to instances", uwsgi_opt_true, &usr.cr.quiet, 0}, + {"sslrouter-cheap", no_argument, 0, "run the sslrouter in cheap mode", uwsgi_opt_true, &usr.cr.cheap, 0}, + {"sslrouter-subscription-server", required_argument, 0, "run the sslrouter subscription server on the spcified address", uwsgi_opt_corerouter_ss, &usr, 0}, + + {"sslrouter-timeout", required_argument, 0, "set sslrouter timeout", uwsgi_opt_set_int, &usr.cr.socket_timeout, 0}, + + {"sslrouter-stats", required_argument, 0, "run the sslrouter stats server", uwsgi_opt_set_str, &usr.cr.stats_server, 0}, + {"sslrouter-stats-server", required_argument, 0, "run the sslrouter stats server", uwsgi_opt_set_str, &usr.cr.stats_server, 0}, + {"sslrouter-ss", required_argument, 0, "run the sslrouter stats server", uwsgi_opt_set_str, &usr.cr.stats_server, 0}, + {"sslrouter-harakiri", required_argument, 0, "enable sslrouter harakiri", uwsgi_opt_set_int, &usr.cr.harakiri, 0}, + + {"sslrouter-sni", no_argument, 0, "use SNI to route requests", uwsgi_opt_true, &usr.sni, 0}, + + {0, 0, 0, 0, 0, 0, 0}, +}; + +static ssize_t sr_write(struct corerouter_peer *); + +// write to backend +static ssize_t sr_instance_write(struct corerouter_peer *peer) { + ssize_t len = cr_write(peer, "sr_instance_write()"); + // end on empty write + if (!len) return 0; + + // the chunk has been sent, start (again) reading from client and instances + if (cr_write_complete(peer)) { + // reset the buffer + peer->out->pos = 0; + cr_reset_hooks(peer); + } + + return len; +} + +// read from backend +static ssize_t sr_instance_read(struct corerouter_peer *peer) { + ssize_t len = cr_read(peer, "sr_instance_read()"); + if (!len) return 0; + + // set the input buffer as the main output one + peer->session->main_peer->out = peer->in; + peer->session->main_peer->out_pos = 0; + + cr_write_to_main(peer, sr_write); + return len; +} + +// the instance is connected now we cannot retry connections +static ssize_t sr_instance_connected(struct corerouter_peer *peer) { + cr_peer_connected(peer, "sr_instance_connected()"); + // set the output buffer as the main_peer output one + peer->out = peer->session->main_peer->in; + peer->out_pos = 0; + return sr_instance_write(peer); +} + +// retry the connection +static int sr_retry(struct corerouter_peer *peer) { + + struct corerouter_session *cs = peer->session; + struct uwsgi_corerouter *ucr = cs->corerouter; + + if (peer->instance_address_len > 0) goto retry; + + if (ucr->mapper(ucr, peer)) { + return -1; + } + + if (peer->instance_address_len == 0) { + return -1; + } + +retry: + // start async connect (again) + cr_connect(peer, sr_instance_connected); + return 0; +} + +static ssize_t sr_write(struct corerouter_peer *main_peer) { + struct corerouter_session *cs = main_peer->session; + struct sslrouter_session *sr = (struct sslrouter_session *) cs; + + int ret = SSL_write(sr->ssl, main_peer->out->buf + main_peer->out_pos, main_peer->out->pos - main_peer->out_pos); + if (ret > 0) { + main_peer->out_pos += ret; + if (main_peer->out->pos == main_peer->out_pos) { + main_peer->out->pos = 0; + cr_reset_hooks(main_peer); + } + return ret; + } + if (ret == 0) return 0; + int err = SSL_get_error(sr->ssl, ret); + + if (err == SSL_ERROR_WANT_READ) { + cr_reset_hooks_and_read(main_peer, sr_write); + return 1; + } + + else if (err == SSL_ERROR_WANT_WRITE) { + cr_write_to_main(main_peer, sr_write); + return 1; + } + + else if (err == SSL_ERROR_SYSCALL) { + uwsgi_error("sr_write()"); + } + + else if (err == SSL_ERROR_SSL && uwsgi.ssl_verbose) { + ERR_print_errors_fp(stderr); + } + + return -1; +} + +static ssize_t sr_read(struct corerouter_peer *main_peer) { + struct corerouter_session *cs = main_peer->session; + struct sslrouter_session *sr = (struct sslrouter_session *) cs; + + int ret = SSL_read(sr->ssl, main_peer->in->buf + main_peer->in->pos, main_peer->in->len - main_peer->in->pos); + if (ret > 0) { + // fix the buffer + main_peer->in->pos += ret; + // check for pending data + int ret2 = SSL_pending(sr->ssl); + if (ret2 > 0) { + if (uwsgi_buffer_fix(main_peer->in, main_peer->in->len + ret2 )) { + uwsgi_log("[uwsgi-sslrouter] cannot fix the buffer to %d\n", main_peer->in->len + ret2); + return -1; + } + if (SSL_read(sr->ssl, main_peer->in->buf + main_peer->in->pos, ret2) != ret2) { + uwsgi_log("[uwsgi-sslrouter] SSL_read() on %d bytes of pending data failed\n", ret2); + return -1; + } + // fix the buffer + main_peer->in->pos += ret2; + } + if (!main_peer->session->peers) { + // add a new peer + struct corerouter_peer *peer = uwsgi_cr_peer_add(cs); + // set default peer hook + peer->last_hook_read = sr_instance_read; + // use the address as hostname + peer->key = cs->ugs->name; + peer->key_len = cs->ugs->name_len; + // the mapper hook + if (cs->corerouter->mapper(cs->corerouter, peer)) { + return -1; + } + + if (peer->instance_address_len == 0) { + return -1; + } + + cr_connect(peer, sr_instance_connected); + return 1; + } + main_peer->session->peers->out = main_peer->in; + main_peer->session->peers->out_pos = 0; + cr_write_to_backend(main_peer, sr_instance_write); + return ret; + } + if (ret == 0) return 0; + int err = SSL_get_error(sr->ssl, ret); + + if (err == SSL_ERROR_WANT_READ) { + cr_reset_hooks_and_read(main_peer, sr_read); + return 1; + } + + else if (err == SSL_ERROR_WANT_WRITE) { + cr_write_to_main(main_peer, sr_read); + return 1; + } + + else if (err == SSL_ERROR_SYSCALL) { + uwsgi_error("sr_ssl_read()"); + } + + else if (err == SSL_ERROR_SSL && uwsgi.ssl_verbose) { + ERR_print_errors_fp(stderr); + } + + return -1; +} + +static void sr_session_close(struct corerouter_session *cs) { + struct sslrouter_session *sr = (struct sslrouter_session *) cs; + SSL_free(sr->ssl); +} + +// allocate a new session +static int sslrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) { + + // set close hook + cs->close = sr_session_close; + // set retry hook + cs->retry = sr_retry; + + struct sslrouter_session *sr = (struct sslrouter_session *) cs; + + if (sa && sa->sa_family == AF_INET) { + sr->ip_addr = ((struct sockaddr_in *) sa)->sin_addr.s_addr; + } + + sr->ssl = SSL_new(ugs->ctx); + SSL_set_fd(sr->ssl, cs->main_peer->fd); + SSL_set_accept_state(sr->ssl); + + uwsgi_cr_set_hooks(cs->main_peer, sr_read, NULL); + + return 0; +} + +static int sslrouter_init() { + + usr.cr.session_size = sizeof(struct sslrouter_session); + usr.cr.alloc_session = sslrouter_alloc_session; + uwsgi_corerouter_init((struct uwsgi_corerouter *) &usr); + + return 0; +} + +static void sslrouter_setup() { + usr.cr.name = uwsgi_str("uWSGI sslrouter"); + usr.cr.short_name = uwsgi_str("sslrouter"); +} + +struct uwsgi_plugin sslrouter_plugin = { + + .name = "sslrouter", + .options = sslrouter_options, + .init = sslrouter_init, + .on_load = sslrouter_setup +}; + +#endif diff --git a/plugins/sslrouter/uwsgiplugin.py b/plugins/sslrouter/uwsgiplugin.py new file mode 100644 index 00000000..171fb5a3 --- /dev/null +++ b/plugins/sslrouter/uwsgiplugin.py @@ -0,0 +1,9 @@ + +NAME='sslrouter' +CFLAGS = [] +LDFLAGS = [] +LIBS = [] + +REQUIRES = ['corerouter'] + +GCC_LIST = ['sslrouter']