Compare commits

...
47 Commits
Author SHA1 Message Date
Unbit 66099275ca updated gemspec file 2014-12-29 10:07:07 +01:00
Unbit e7a40eb173 completed backporting of fastrouter post buffering 2014-12-27 13:30:16 +01:00
Riccardo Magliocchetti 8342150b78 Fix fastrouter compilation
After 6a8f18a9bb

Fix #801
2014-12-27 13:28:41 +01:00
Roberto De Ioris e2253f3693 backport of fastrouter buffering 2014-12-27 13:28:08 +01:00
Roberto De Ioris d2f07c7583 implemented uwsgi::opt in perl 2014-12-12 07:01:28 +01:00
Roberto De Ioris 1c87053a13 implemented --pull-header 2014-12-11 15:54:21 +01:00
Unbit d69e23414f another check for mod_proxy_uwsgi 2014-12-10 10:45:06 +01:00
Roberto De Ioris 3d37d272ee fixed master-fifo + cheaper 2014-12-09 09:53:45 +01:00
Riccardo Magliocchetti 9094a0f811 Fix leak on error in bind_to_unix
Spotted by cppcheck.
2014-12-09 09:44:36 +01:00
Unbit 97efcca2b7 improved mono plugin 2014-12-04 20:55:58 +01:00
Roberto De Ioris 2445fd1f61 fixed #776 2014-12-06 10:39:27 +01:00
Unbit 3234e0aae8 fixed #680 2014-12-06 09:14:17 +01:00
Unbit c5e7eb5d95 fix #785 2014-12-05 15:22:53 +01:00
Unbit e051e4fe1c another attempt at fixing mod_proxy_uwsgi 2014-12-05 15:08:47 +01:00
Unbit 5f691ec987 attempt to fix #774 2014-12-05 09:23:55 +01:00
Unbit a8e0aac9b0 fixed #787 2014-12-05 09:11:11 +01:00
Unbit bcdbab247c fix peer name in corerouters [2] 2014-11-22 13:22:21 +01:00
Roberto De Ioris 2e378dfc29 Merge branch 'uwsgi-2.0' of github.com:unbit/uwsgi into uwsgi-2.0 2014-11-22 13:19:28 +01:00
Roberto De Ioris acc270ad81 try to fix peer name in corerouter 2014-11-22 13:12:11 +01:00
Unbit b0c3ee02f9 detect modern pypy 2014-11-21 08:35:41 +01:00
Unbit 06ab5eaee5 fix #778 2014-11-21 08:34:44 +01:00
Roberto De Ioris 46de2ae0cb attempt to fix #340 2014-11-21 06:47:04 +01:00
Unbit 65fe93ac14 fixed hr_ssl_clear_errors 2014-11-20 07:35:18 +01:00
Roberto De Ioris f7cdd22f9a attempt to fix error propagation in https 2014-11-20 07:25:19 +01:00
Andjelko Horvat 3d553d2ae0 Daemon touch kill failback. 2014-11-19 10:45:59 +01:00
Unbit 582d5828aa another round of ssl fixes 2014-11-18 19:15:29 +01:00
Unbit 245a66b0e5 try to improve ssl management 2014-11-18 18:57:59 +01:00
Roberto De Ioris 596c81124e fixed destruction 2014-11-16 20:59:45 +01:00
Unbit a86a1fcc40 added active-workers signal target, aimed at improving #58 2014-11-16 18:47:33 +01:00
Unbit 34036ac8b4 Merge branch 'uwsgi-2.0' of https://github.com/unbit/uwsgi into uwsgi-2.0 2014-11-15 13:17:56 +01:00
Unbit 11cd6c34a5 fixed python3 --py-auto-reload-ignore 2014-11-15 13:17:52 +01:00
Unbit aaea6ff2e2 fixed modifiers usage in corerouters 2014-11-15 06:32:25 +01:00
Unbit 38c30bc257 bump version to 2.0.9 2014-11-15 05:26:20 +01:00
Unbit 2f3e83b78e added httpdumb router and subscribe-with-modifier 2014-11-14 15:41:42 +01:00
Roberto De Ioris 43e8c65841 support yajl brew 2014-11-14 06:03:08 +01:00
Unbit 555902a0ee less strict http generator 2014-11-13 11:56:44 +01:00
Unbit 9d4f4eb3c3 backported #772 2014-11-12 07:43:34 +01:00
Ævar Arnfjörð Bjarmason 5d338cfb3a psgi: Ensure that we call any DESTROY hooks on psgix.harakiri.commit
Before this we'd just exit(0) and let the OS clean up after us, but
e.g. with post-buffering=1 we'll end up with a temporary file in /tmp
that we won't clean up when we exit unless DESTROY is called.

This resulted in us leaking files in /tmp if we ever had a request where
the last request before a harakiri was a POST request with a body we'd
buffer to /tmp.

We'd have similar leaks in any user-defined code that required DESTROY
to run.

Aside from this I'm still not very comfortable with what this whole code
here in psgi_plugin.c and psgi_loader.c is doing when managing the
interpreter(s). It:

 * Doesn't consistently call PERL_SET_CONTEXT() as described in "perldoc
   perlembed".

 * Nothing calls PERL_SYS_TERM() either.

 * Should we be calling uwsgi_perl_free_stashes() here too?

To test this:

    UWSGI_PROFILE=psgi python uwsgiconfig.py --build
    ./uwsgi --master --http-socket localhost:1234 --psgi t/perl/test_harakiri.psgi

Then elsewhere:

    curl 'localhost:1234?0'
    curl 'localhost:1234?1'

Both of those should emit "Calling DESTROY".
2014-11-12 07:38:18 +01:00
Riccardo Magliocchetti 3406441d95 systemdlogger: fix compilation with -Werror=format-security 2014-11-06 18:53:30 +01:00
Unbit df08c7bfb5 fixed unmasked websocket 2014-11-06 09:02:21 +01:00
Unbit 75b4a3cb28 clear if sweep_on_full fails 2014-11-05 13:06:24 +01:00
Unbit c356c01f12 Merge branch 'uwsgi-2.0' of https://github.com/unbit/uwsgi into uwsgi-2.0 2014-11-02 10:51:50 +01:00
Unbit ba8aa44baa attempt to fix #766 2014-11-02 10:51:39 +01:00
Roberto De Ioris a1dc7a660c use uwsgi_buffer for send_ack in #760 2014-11-02 09:33:07 +01:00
Roberto De Ioris 3a4663844b fixed #760 2014-11-02 06:42:53 +01:00
Mattia Barbon 309f449a43 Fix latent refcounting bug
It can be reproduced by enabling the Perl debugger inside a PSGI application:

    {
        package DB;

        sub DB { }
        sub sub { &$sub }
    }

    $^P = 0x73f;

    sub { [200, ['Content-Type' => 'text/plain'], ['Hello World']] }

For every request the following warnings are emitted:

    Attempt to free unreferenced scalar: SV 0xfea6e8, Perl interpreter: 0xd534a0.
    Attempt to free unreferenced scalar: SV 0xfea718, Perl interpreter: 0xd534a0.

where the unreferenced scalars are the uwsgi::input/uwsgi::error instances
created in build_psgi_env.

The calling convention for Perl subroutines is that the values pushed on the
stack must be mortalized in the callee, and if the caller wants to retain them,
it must do a SvREFCNT_inc to undo the effect of the mortalization.

Before this patch XS_input/XS_error were not mortalizing the value, and
uwsgi_perl_obj_new was not incrementing the reference count, so the two bugs
balanced each other.

When running under debugger, Perl forwards all function/method calls to
DB::sub, which causes a mortal copy of the return value of
uwsgi::input/error::new to be pushed on the stack. The value is cleared by the
FREETMPS at the end of uwsgi_perl_obj_new, and the freed value is added to the
environment hash. The warning is emitted at the end of the request when the
environment hash is freed and Perl notices that some of the values has been
already freed.
2014-11-02 06:13:20 +01:00
Mattia Barbon f059a69312 Remove unnecessary mortalization
newRV(sv_newmortal()) is equivalent to newRV_noinc(newSV(0)): the former
creates a new SV with refcount 1, schedules a decrement "soon" (the
mortalization) and increments the refcount, the net result is a refcount of 1,
which is what the latter does.
2014-11-02 06:13:04 +01:00
41 changed files with 647 additions and 125 deletions
+14 -4
View File
@@ -365,12 +365,13 @@ static int uwsgi_response(request_rec *r, proxy_conn_rec *backend, proxy_server_
ap_set_content_type(r, apr_pstrdup(r->pool, buf));
}
for(;;) {
int finish = 0;
while(!finish) {
rv = ap_get_brigade(rp->input_filters, bb,
AP_MODE_READBYTES, mode,
conf->io_buffer_size);
if (mode == APR_NONBLOCK_READ && (APR_STATUS_IS_EAGAIN(rv)
|| (rv == APR_SUCCESS && APR_BRIGADE_EMPTY(bb)))) {
if (APR_STATUS_IS_EAGAIN(rv)
|| (rv == APR_SUCCESS && APR_BRIGADE_EMPTY(bb)) ) {
e = apr_bucket_flush_create(c->bucket_alloc);
APR_BRIGADE_INSERT_TAIL(bb, e);
if (ap_pass_brigade(r->output_filters, bb) || c->aborted) {
@@ -401,7 +402,16 @@ static int uwsgi_response(request_rec *r, proxy_conn_rec *backend, proxy_server_
ap_proxy_buckets_lifetime_transform(r, bb, pass_bb);
ap_pass_brigade(r->output_filters, pass_bb);
// found the last brigade?
if (APR_BUCKET_IS_EOS(APR_BRIGADE_LAST(bb))) finish = 1;
// do not pass chunk if it is zero_sized
apr_brigade_length(pass_bb, 0, &readbytes);
if ((readbytes > 0 && ap_pass_brigade(r->output_filters, pass_bb) != APR_SUCCESS) || c->aborted) {
finish = 1;
}
apr_brigade_cleanup(bb);
apr_brigade_cleanup(pass_bb);
}
+9 -2
View File
@@ -26,6 +26,7 @@ extern struct uwsgi_server uwsgi;
static void cache_full(struct uwsgi_cache *uc) {
uint64_t i;
int force_clear = 0;
if (!uc->ignore_full) {
if (uc->purge_lru)
@@ -41,19 +42,25 @@ static void cache_full(struct uwsgi_cache *uc) {
// we do not need locking here !
if (uc->sweep_on_full) {
uint64_t removed = 0;
uint64_t now = (uint64_t) uwsgi_now();
if (uc->next_scan <= now) {
uc->next_scan = now + uc->sweep_on_full;
for (i = 1; i < uc->max_items; i++) {
struct uwsgi_cache_item *uci = cache_item(i);
if (uci->expires > 0 && uci->expires <= now) {
uwsgi_cache_del2(uc, NULL, 0, i, 0);
if (!uwsgi_cache_del2(uc, NULL, 0, i, 0)) {
removed++;
}
}
}
if (removed == 0) {
force_clear = 1;
}
}
}
if (uc->clear_on_full) {
if (uc->clear_on_full || force_clear) {
for (i = 1; i < uc->max_items; i++) {
uwsgi_cache_del2(uc, NULL, 0, i, 0);
}
+12 -1
View File
@@ -13,6 +13,7 @@ void emperor_send_stats(int);
time_t emperor_throttle;
int emperor_throttle_level;
int emperor_warming_up = 1;
struct uwsgi_instance *ui;
@@ -867,7 +868,16 @@ void emperor_add(struct uwsgi_emperor_scanner *ues, char *name, time_t born, cha
#ifdef UWSGI_DEBUG
uwsgi_log("emperor throttle = %d\n", emperor_throttle_level);
#endif
usleep(emperor_throttle_level);
if (emperor_warming_up) {
if (emperor_throttle_level > 0) {
// wait 10 milliseconds in case of fork-bombing
// pretty random value, but should avoid the load average to increase
usleep(10);
}
}
else {
usleep(emperor_throttle_level);
}
if (uwsgi.emperor_tyrant) {
if (uid == 0 || gid == 0) {
@@ -1554,6 +1564,7 @@ void uwsgi_emperor_run_scanners(void) {
ues->monitor->func(ues);
ues = ues->next;
}
emperor_warming_up = 0;
}
void emperor_build_scanners() {
+4 -1
View File
@@ -861,7 +861,10 @@ int master_loop(char **argv, char **environ) {
if (touched) {
uwsgi_log_verbose("*** %s has been touched... reloading daemon \"%s\" (pid: %d) !!! ***\n", touched, ud->command, (int) ud->pid);
if (kill(-ud->pid, ud->stop_signal)) {
uwsgi_error("[uwsgi-daemon/touch] kill()");
// killing process group failed, try to kill by process id
if (kill(ud->pid, ud->stop_signal)) {
uwsgi_error("[uwsgi-daemon/touch] kill()");
}
}
}
}
+14
View File
@@ -138,12 +138,25 @@ void uwsgi_master_check_idle() {
uwsgi.workers[i].cheaped = 1;
if (uwsgi.workers[i].pid == 0)
continue;
// first send SIGINT
kill(uwsgi.workers[i].pid, SIGINT);
// and start waiting upto 3 seconds
int j;
for(j=0;j<3;j++) {
sleep(1);
int ret = waitpid(uwsgi.workers[i].pid, &waitpid_status, WNOHANG);
if (ret == 0) continue;
if (ret == (int) uwsgi.workers[i].pid) goto done;
// on error, directly send SIGKILL
break;
}
kill(uwsgi.workers[i].pid, SIGKILL);
if (waitpid(uwsgi.workers[i].pid, &waitpid_status, 0) < 0) {
if (errno != ECHILD)
uwsgi_error("uwsgi_master_check_idle()/waitpid()");
}
else {
done:
uwsgi.workers[i].pid = 0;
uwsgi.workers[i].rss_size = 0;
uwsgi.workers[i].vsz_size = 0;
@@ -335,6 +348,7 @@ int uwsgi_master_check_daemons_death(int diedpid) {
int uwsgi_worker_is_busy(int wid) {
int i;
if (uwsgi.workers[uwsgi.mywid].sig) return 1;
for(i=0;i<uwsgi.cores;i++) {
if (uwsgi.workers[wid].cores[i].in_request) {
return 1;
+2
View File
@@ -162,6 +162,7 @@ int uwsgi_calc_cheaper(void) {
ignore_algo = 1;
}
uwsgi.cheaper_fifo_delta = 0;
goto safe;
}
// if cheaper limits wants to change worker count, then skip cheaper algo
@@ -172,6 +173,7 @@ int uwsgi_calc_cheaper(void) {
needed_workers = 0;
}
safe:
if (needed_workers > 0) {
for (i = 1; i <= uwsgi.numproc; i++) {
if (uwsgi.workers[i].cheaped == 1 && uwsgi.workers[i].pid == 0) {
+14
View File
@@ -2045,5 +2045,19 @@ next4:
usl->custom2 = strlen(space+1);
uwsgi_log("collecting header %.*s to var %s\n", usl->custom, usl->value, usl->custom_ptr);
}
uwsgi_foreach(usl, uwsgi.pull_headers) {
char *space = strchr(usl->value, ' ');
if (!space) {
uwsgi_log("invalid pull header syntax, must be <header> <var>\n");
exit(1);
}
*space = 0;
usl->custom = strlen(usl->value);
*space = ' ';
usl->custom_ptr = space+1;
usl->custom2 = strlen(space+1);
uwsgi_log("pulling header %.*s to var %s\n", usl->custom, usl->value, usl->custom_ptr);
}
}
#endif
+10
View File
@@ -331,6 +331,16 @@ void uwsgi_route_signal(uint8_t sig) {
}
}
}
// send to al lactive workers
else if (!strcmp(use->receiver, "active-workers")) {
for (i = 1; i <= uwsgi.numproc; i++) {
if (uwsgi.workers[i].pid > 0 && !uwsgi.workers[i].cheaped && !uwsgi.workers[i].suspended) {
if (uwsgi_signal_send(uwsgi.workers[i].signal_pipe[0], sig)) {
uwsgi_log("could not deliver signal %d to worker %d\n", sig, i);
}
}
}
}
// route to specific worker
else if (!strncmp(use->receiver, "worker", 6)) {
i = atoi(use->receiver + 6);
+4 -1
View File
@@ -189,7 +189,10 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst
memset(uws_addr, 0, sizeof(struct sockaddr_un));
serverfd = create_server_socket(AF_UNIX, SOCK_STREAM);
if (serverfd < 0) return -1;
if (serverfd < 0) {
free(uws_addr);
return -1;
}
if (abstract_socket == 0) {
if (unlink(socket_name) != 0 && errno != ENOENT) {
uwsgi_error("error removing unix socket, unlink()");
+4
View File
@@ -679,6 +679,10 @@ static struct uwsgi_buffer *uwsgi_subscription_ub(char *key, size_t keysize, uin
goto end;
if (uwsgi_buffer_append_keyval(ub, "address", 7, socket_name, strlen(socket_name)))
goto end;
if (uwsgi.subscribe_with_modifier1) {
modifier1 = atoi(uwsgi.subscribe_with_modifier1);
}
if (uwsgi_buffer_append_keynum(ub, "modifier1", 9, modifier1))
goto end;
if (uwsgi_buffer_append_keynum(ub, "modifier2", 9, modifier2))
+5
View File
@@ -632,6 +632,9 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"subscription-tolerance", required_argument, 0, "set tolerance for subscription servers", uwsgi_opt_set_int, &uwsgi.subscription_tolerance, 0},
{"unsubscribe-on-graceful-reload", no_argument, 0, "force unsubscribe request even during graceful reload", uwsgi_opt_true, &uwsgi.unsubscribe_on_graceful_reload, 0},
{"start-unsubscribed", no_argument, 0, "configure subscriptions but do not send them (useful with master fifo)", uwsgi_opt_true, &uwsgi.subscriptions_blocked, 0},
{"subscribe-with-modifier1", required_argument, 0, "force the specififed modifier1 when subscribing", uwsgi_opt_set_str, &uwsgi.subscribe_with_modifier1, UWSGI_OPT_MASTER},
{"snmp", optional_argument, 0, "enable the embedded snmp server", uwsgi_opt_snmp, NULL, 0},
{"snmp-community", required_argument, 0, "set the snmp community string", uwsgi_opt_snmp_community, NULL, 0},
#ifdef UWSGI_SSL
@@ -859,6 +862,8 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"collect-header", required_argument, 0, "store the specified response header in a request var (syntax: header var)", uwsgi_opt_add_string_list, &uwsgi.collect_headers, 0},
{"response-header-collect", required_argument, 0, "store the specified response header in a request var (syntax: header var)", uwsgi_opt_add_string_list, &uwsgi.collect_headers, 0},
{"pull-header", required_argument, 0, "store the specified response header in a request var and remove it from the response (syntax: header var)", uwsgi_opt_add_string_list, &uwsgi.pull_headers, 0},
{"check-static", required_argument, 0, "check for static files in the specified directory", uwsgi_opt_check_static, NULL, UWSGI_OPT_MIME},
{"check-static-docroot", no_argument, 0, "check for static files in the requested DOCUMENT_ROOT", uwsgi_opt_true, &uwsgi.check_static_docroot, UWSGI_OPT_MIME},
{"static-check", required_argument, 0, "check for static files in the specified directory", uwsgi_opt_check_static, NULL, UWSGI_OPT_MIME},
+1
View File
@@ -273,6 +273,7 @@ static struct uwsgi_buffer *uwsgi_websocket_recv_do(struct wsgi_request *wsgi_re
}
else {
wsgi_req->websocket_need += wsgi_req->websocket_size;
wsgi_req->websocket_pktsize += wsgi_req->websocket_size;
wsgi_req->websocket_phase = 4;
}
break;
+10 -1
View File
@@ -124,9 +124,18 @@ error:
static int uwsgi_response_add_header_do(struct wsgi_request *wsgi_req, char *key, uint16_t key_len, char *value, uint16_t value_len) {
// collect the header ?
// pull/collect the header ?
struct uwsgi_string_list *usl = NULL;
uwsgi_foreach(usl, uwsgi.pull_headers) {
if (!uwsgi_strnicmp(key, key_len, usl->value, usl->custom)) {
if (!uwsgi_req_append(wsgi_req, usl->custom_ptr, usl->custom2, value, value_len)) {
wsgi_req->write_errors++ ; return -1;
}
return 0;
}
}
uwsgi_foreach(usl, uwsgi.collect_headers) {
if (!uwsgi_strnicmp(key, key_len, usl->value, usl->custom)) {
if (!uwsgi_req_append(wsgi_req, usl->custom_ptr, usl->custom2, value, value_len)) {
+1 -1
View File
@@ -328,7 +328,7 @@ static void asyncio_loop() {
uwsgi.schedule_fix = uwsgi_asyncio_schedule_fix;
}
#ifndef UWSGI_PYTHREE
#ifndef PYTHREE
PyObject *asyncio = PyImport_ImportModule("trollius");
#else
PyObject *asyncio = PyImport_ImportModule("asyncio");
+4 -1
View File
@@ -71,7 +71,9 @@ static void carbon_post_init() {
u_server->errors = 0;
char *p, *ctx = NULL;
uwsgi_foreach_token(usl->value, ":", p, ctx) {
// make a copy to not clobber argv
char *tmp = uwsgi_str(usl->value);
uwsgi_foreach_token(tmp, ":", p, ctx) {
if (!u_server->hostname) {
u_server->hostname = uwsgi_str(p);
}
@@ -81,6 +83,7 @@ static void carbon_post_init() {
else
break;
}
free(tmp);
if (!u_server->hostname || !u_server->port) {
uwsgi_log("[carbon] invalid carbon server address (%s)\n", usl->value);
usl = usl->next;
+6
View File
@@ -75,6 +75,12 @@ void uwsgi_cr_peer_reset(struct corerouter_peer *peer) {
peer->hook_write = NULL;
}
if (peer->is_buffering) {
if (peer->buffering_fd != -1) {
close(peer->buffering_fd);
}
}
peer->failed = 0;
peer->soopt = 0;
peer->timed_out = 0;
+5 -2
View File
@@ -180,8 +180,8 @@ struct corerouter_peer {
uint16_t retries;
// parsed key
char *key;
uint16_t key_len;
char key[0xff];
uint8_t key_len;
uint8_t modifier1;
uint8_t modifier2;
@@ -192,7 +192,10 @@ struct corerouter_peer {
int current_timeout;
ssize_t (*flush)(struct corerouter_peer *);
int is_flushing;
int is_buffering;
int buffering_fd;
};
struct uwsgi_corerouter {
+2
View File
@@ -53,6 +53,7 @@ int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *ucr, struct coreroute
peer->instance_address = peer->un->name;
peer->instance_address_len = peer->un->len;
peer->modifier1 = peer->un->modifier1;
peer->modifier2 = peer->un->modifier2;
}
else if (ucr->cheap && !ucr->i_am_cheap && uwsgi_no_subscriptions(ucr->subscriptions)) {
uwsgi_gateway_go_cheap(ucr->name, ucr->queue, &ucr->i_am_cheap);
@@ -88,6 +89,7 @@ split:
peer->instance_address = peer->un->name;
peer->instance_address_len = peer->un->len;
peer->modifier1 = peer->un->modifier1;
peer->modifier2 = peer->un->modifier2;
}
else if (ucr->cheap && !ucr->i_am_cheap && uwsgi_no_subscriptions(ucr->subscriptions)) {
uwsgi_gateway_go_cheap(ucr->name, ucr->queue, &ucr->i_am_cheap);
+17 -26
View File
@@ -2,13 +2,6 @@
#define AMQP_CONNECTION_HEADER "AMQP\0\0\x09\x01"
#ifdef __BIG_ENDIAN__
#define ntohll(x) x
#else
#define ntohll(x) ( ( (uint64_t)(ntohl( (uint32_t)((x << 32) >> 32) )) << 32) | ntohl( ((uint32_t)(x >> 32)) ) )
#endif
#define htonll(x) ntohll(x)
#define amqp_send(a, b, c) if (send(a, b, c, 0) < 0) { uwsgi_error("send()"); return -1; }
struct amqp_frame_header {
@@ -69,13 +62,8 @@ static char *amqp_get_long(char *ptr, char *watermark, uint32_t *lv) {
static char *amqp_get_longlong(char *ptr, char *watermark, uint64_t *llv) {
uint64_t tmp_longlong;
if (ptr+8 > watermark) return NULL;
memcpy(&tmp_longlong, ptr, 8);
*llv = ntohll(tmp_longlong);
*llv = uwsgi_be64(ptr);
return ptr+8;
}
@@ -85,26 +73,29 @@ static int amqp_send_ack(int fd, uint64_t delivery_tag) {
uint32_t size = 4 + 8 + 1;
size = htonl(size);
struct uwsgi_buffer *ub = uwsgi_buffer_new(64);
// send type and channel
amqp_send(fd, "\1\0\1", 3);
if (uwsgi_buffer_append(ub, "\1\0\1", 3)) goto end;
// send size
amqp_send(fd, &size, 4);
// send class 60 method 80
amqp_send(fd, "\x00\x3C\x00\x50", 4);
if (uwsgi_buffer_u32be(ub, size)) goto end;
// send class 60 method 80
if (uwsgi_buffer_append(ub, "\x00\x3C\x00\x50", 4)) goto end;
// set delivery_tag
delivery_tag = htonll(delivery_tag);
amqp_send(fd, &delivery_tag, 8);
if (uwsgi_buffer_u64be(ub, delivery_tag)) goto end;
if (uwsgi_buffer_append(ub, "\0\xCE", 2)) goto end;
// empty bits
amqp_send(fd, "\0", 1);
// send buffer to socket
if (write(fd, ub->buf, ub->pos) < 0) {
uwsgi_error("amqp_send_ack()/write()");
goto end;
}
// send frame-end
amqp_send(fd, "\xCE", 1);
uwsgi_buffer_destroy(ub);
return 0;
end:
uwsgi_buffer_destroy(ub);
return -1;
}
char *uwsgi_amqp_consume(int fd, uint64_t *msgsize, char **routing_key) {
+118 -15
View File
@@ -16,6 +16,8 @@ extern struct uwsgi_server uwsgi;
struct fastrouter_session {
struct corerouter_session session;
int has_key;
uint64_t content_length;
uint64_t buffered;
};
static struct uwsgi_option fastrouter_options[] = {
@@ -42,7 +44,7 @@ static struct uwsgi_option fastrouter_options[] = {
{"fastrouter-timeout", required_argument, 0, "set fastrouter timeout", uwsgi_opt_set_int, &ufr.cr.socket_timeout, 0},
{"fastrouter-post-buffering", required_argument, 0, "enable fastrouter post buffering", uwsgi_opt_set_64bit, &ufr.cr.post_buffering, 0},
{"fastrouter-post-buffering-dir", required_argument, 0, "put fastrouter buffered files to the specified directory", uwsgi_opt_set_str, &ufr.cr.pb_base_dir, 0},
{"fastrouter-post-buffering-dir", required_argument, 0, "put fastrouter buffered files to the specified directory (noop, use TMPDIR env)", uwsgi_opt_set_str, &ufr.cr.pb_base_dir, 0},
{"fastrouter-stats", required_argument, 0, "run the fastrouter stats server", uwsgi_opt_set_str, &ufr.cr.stats_server, 0},
{"fastrouter-stats-server", required_argument, 0, "run the fastrouter stats server", uwsgi_opt_set_str, &ufr.cr.stats_server, 0},
@@ -66,21 +68,27 @@ static void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vall
//uwsgi_log("%.*s = %.*s\n", keylen, key, vallen, val);
if (!uwsgi_strncmp("SERVER_NAME", 11, key, keylen) && !peer->key_len) {
peer->key = val;
peer->key_len = vallen;
if (vallen <= 0xff) {
memcpy(peer->key, val, vallen);
peer->key_len = vallen;
}
return;
}
if (!uwsgi_strncmp("HTTP_HOST", 9, key, keylen) && !fr->has_key) {
peer->key = val;
peer->key_len = vallen;
if (vallen <= 0xff) {
memcpy(peer->key, val, vallen);
peer->key_len = vallen;
}
return;
}
if (!uwsgi_strncmp("UWSGI_FASTROUTER_KEY", 20, key, keylen)) {
fr->has_key = 1;
peer->key = val;
peer->key_len = vallen;
if (vallen <= 0xff) {
fr->has_key = 1;
memcpy(peer->key, val, vallen);
peer->key_len = vallen;
}
return;
}
@@ -97,6 +105,12 @@ static void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vall
}
return;
}
if (ufr.cr.post_buffering > 0) {
if (!uwsgi_strncmp("CONTENT_LENGTH", 14, key, keylen)) {
fr->content_length = uwsgi_str_num(val, vallen);
}
}
}
// writing client body to the instance
@@ -157,6 +171,23 @@ static ssize_t fr_instance_read(struct corerouter_peer *peer) {
return len;
}
static ssize_t fr_instance_sendfile(struct corerouter_peer *peer) {
struct fastrouter_session *fr = (struct fastrouter_session *) peer->session;
ssize_t len = uwsgi_sendfile_do(peer->fd, peer->session->main_peer->buffering_fd, fr->buffered, fr->content_length - fr->buffered);
if (len < 0) {
cr_try_again;
uwsgi_cr_error(peer, "fr_instance_sendfile()/sendfile()");
return -1;
}
if (len == 0) return 0;
fr->buffered += len;
if (peer != peer->session->main_peer && peer->un) peer->un->rx+=len;
if (fr->buffered >= fr->content_length) {
cr_reset_hooks(peer);
}
return len;
}
// send the uwsgi request header and vars
static ssize_t fr_instance_send_request(struct corerouter_peer *peer) {
ssize_t len = cr_write(peer, "fr_instance_send_request()");
@@ -167,9 +198,16 @@ static ssize_t fr_instance_send_request(struct corerouter_peer *peer) {
if (cr_write_complete(peer)) {
// reset the original read buffer
peer->out->pos = 0;
// start waiting for body
peer->session->main_peer->last_hook_read = fr_read_body;
cr_reset_hooks(peer);
if (!peer->session->main_peer->is_buffering) {
// start waiting for body
peer->session->main_peer->last_hook_read = fr_read_body;
cr_reset_hooks(peer);
}
else {
peer->hook_write = fr_instance_sendfile;
// stop reading from the client
peer->session->main_peer->last_hook_read = NULL;
}
}
return len;
@@ -184,8 +222,8 @@ static ssize_t fr_instance_connected(struct corerouter_peer *peer) {
peer->can_retry = 0;
// fix modifiers
peer->in->buf[0] = peer->session->main_peer->modifier1;
peer->in->buf[3] = peer->session->main_peer->modifier2;
peer->in->buf[0] = peer->modifier1;
peer->in->buf[3] = peer->modifier2;
// prepare to write the uwsgi packet
peer->out = peer->session->main_peer->in;
@@ -197,20 +235,73 @@ static ssize_t fr_instance_connected(struct corerouter_peer *peer) {
// called after receaving the uwsgi header (read vars)
static ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) {
struct fastrouter_session *fr = (struct fastrouter_session *) main_peer->session;
struct corerouter_peer *new_peer = NULL;
ssize_t len = 0;
struct uwsgi_header *uh = (struct uwsgi_header *) main_peer->in->buf;
// better to store it as the original buf address could change
uint16_t pktsize = uh->pktsize;
// are we buffering ?
if (main_peer->is_buffering) {
// memory or disk ?
if (fr->content_length <= ufr.cr.post_buffering) {
// increase buffer if needed
if (uwsgi_buffer_fix(main_peer->in, pktsize+4+fr->content_length))
return -1;
len = cr_read_exact(main_peer, pktsize+4+fr->content_length, "fr_recv_uwsgi_vars()");
if (!len) return 0;
// whole body read ?
if (main_peer->in->pos == (size_t)(pktsize+4+fr->content_length)) {
main_peer->is_buffering = 0;
goto done;
}
return len;
}
// first round ?
if (main_peer->buffering_fd == -1) {
main_peer->buffering_fd = uwsgi_tmpfd();
if (main_peer->buffering_fd < 0) return -1;
}
char buf[32768];
size_t remains = fr->content_length - fr->buffered;
ssize_t rlen = read(main_peer->fd, buf, UMIN(32768, remains));
if (rlen < 0) {
cr_try_again;
uwsgi_cr_error(main_peer, "fr_recv_uwsgi_vars()/read()");
return -1;
}
if (rlen == 0) return 0;
fr->buffered += rlen;
if (write(main_peer->buffering_fd, buf, rlen) != rlen) {
uwsgi_cr_error(main_peer, "fr_recv_uwsgi_vars()/write()");
return -1;
}
// have we done ?
if (fr->buffered >= fr->content_length) {
fr->buffered = 0;
len = rlen;
goto done;
}
return rlen;
}
// increase buffer if needed
if (uwsgi_buffer_fix(main_peer->in, pktsize+4))
return -1;
ssize_t len = cr_read_exact(main_peer, pktsize+4, "fr_recv_uwsgi_vars()");
len = cr_read_exact(main_peer, pktsize+4, "fr_recv_uwsgi_vars()");
if (!len) return 0;
// headers received, ready to choose the instance
if (main_peer->in->pos == (size_t)(pktsize+4)) {
struct uwsgi_corerouter *ucr = main_peer->session->corerouter;
struct corerouter_peer *new_peer = uwsgi_cr_peer_add(main_peer->session);
new_peer = uwsgi_cr_peer_add(main_peer->session);
new_peer->last_hook_read = fr_instance_read;
// find the hostname
@@ -237,6 +328,18 @@ static ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) {
return -1;
}
// buffering ?
if (ufr.cr.post_buffering > 0 && fr->content_length > 0) {
main_peer->is_buffering = 1;
main_peer->buffering_fd = -1;
return len;
}
done:
if (!new_peer) {
new_peer = main_peer->session->peers;
}
new_peer->can_retry = 1;
cr_connect(new_peer, fr_instance_connected);
+42
View File
@@ -100,6 +100,31 @@ retry:
return Py_None;
}
PyObject *py_uwsgi_gevent_int(PyObject *self, PyObject *args) {
uwsgi_log("Brutally killing worker %d (pid: %d)...\n", uwsgi.mywid, uwsgi.mypid);
uwsgi.workers[uwsgi.mywid].manage_next_request = 0;
uwsgi_log_verbose("stopping gevent signals watchers for worker %d (pid: %d)...\n", uwsgi.mywid, uwsgi.mypid);
PyObject_CallMethod(ugevent.my_signal_watcher, "stop", NULL);
PyObject_CallMethod(ugevent.signal_watcher, "stop", NULL);
uwsgi_log_verbose("stopping gevent sockets watchers for worker %d (pid: %d)...\n", uwsgi.mywid, uwsgi.mypid);
int i,count = uwsgi_count_sockets(uwsgi.sockets);
for(i=0;i<count;i++) {
PyObject_CallMethod(ugevent.watchers[i], "stop", NULL);
}
uwsgi_log_verbose("main gevent watchers stopped for worker %d (pid: %d)...\n", uwsgi.mywid, uwsgi.mypid);
if (!ugevent.wait_for_hub) {
PyObject_CallMethod(ugevent.ctrl_gl, "kill", NULL);
}
Py_INCREF(Py_None);
return Py_None;
}
static void uwsgi_gevent_gbcw() {
// already running
@@ -333,6 +358,7 @@ PyMethodDef uwsgi_gevent_signal_def[] = { {"uwsgi_gevent_signal", py_uwsgi_geven
PyMethodDef uwsgi_gevent_my_signal_def[] = { {"uwsgi_gevent_my_signal", py_uwsgi_gevent_my_signal, METH_VARARGS, ""} };
PyMethodDef uwsgi_gevent_signal_handler_def[] = { {"uwsgi_gevent_signal_handler", py_uwsgi_gevent_signal_handler, METH_VARARGS, ""} };
PyMethodDef uwsgi_gevent_unix_signal_handler_def[] = { {"uwsgi_gevent_unix_signal_handler", py_uwsgi_gevent_graceful, METH_VARARGS, ""} };
PyMethodDef uwsgi_gevent_unix_signal_int_handler_def[] = { {"uwsgi_gevent_unix_signal_int_handler", py_uwsgi_gevent_int, METH_VARARGS, ""} };
PyMethodDef uwsgi_gevent_ctrl_gl_def[] = { {"uwsgi_gevent_ctrl_gl_handler", py_uwsgi_gevent_ctrl_gl, METH_VARARGS, ""} };
static void gil_gevent_get() {
@@ -488,6 +514,22 @@ static void gevent_loop() {
python_call(ugevent.signal, ge_signal_tuple, 0, NULL);
// map SIGINT/SIGTERM with gevent.signal
ge_signal_tuple = PyTuple_New(2);
PyTuple_SetItem(ge_signal_tuple, 0, PyInt_FromLong(SIGINT));
PyObject *uwsgi_gevent_unix_signal_int_handler = PyCFunction_New(uwsgi_gevent_unix_signal_int_handler_def, NULL);
Py_INCREF(uwsgi_gevent_unix_signal_int_handler);
PyTuple_SetItem(ge_signal_tuple, 1, uwsgi_gevent_unix_signal_int_handler);
python_call(ugevent.signal, ge_signal_tuple, 0, NULL);
ge_signal_tuple = PyTuple_New(2);
PyTuple_SetItem(ge_signal_tuple, 0, PyInt_FromLong(SIGTERM));
PyTuple_SetItem(ge_signal_tuple, 1, uwsgi_gevent_unix_signal_int_handler);
python_call(ugevent.signal, ge_signal_tuple, 0, NULL);
PyObject *wait_for_me = ugevent.hub;
if (!ugevent.wait_for_hub) {
+33 -15
View File
@@ -141,9 +141,14 @@ static int http_add_uwsgi_header(struct corerouter_peer *peer, char *hh, size_t
}
if (!uwsgi_strncmp("HOST", 4, hh, keylen)) {
peer->key = val;
peer->key_len = vallen;
if (uhttp.server_name_as_http_host && uwsgi_buffer_append_keyval(out, "SERVER_NAME", 11, peer->key, peer->key_len)) return -1;
if (vallen <= 0xff) {
memcpy(peer->key, val, vallen);
peer->key_len = vallen;
if (uhttp.server_name_as_http_host && uwsgi_buffer_append_keyval(out, "SERVER_NAME", 11, peer->key, peer->key_len)) return -1;
}
else {
return -1;
}
}
else if (!uwsgi_strncmp("CONTENT_LENGTH", 14, hh, keylen)) {
@@ -166,8 +171,10 @@ static int http_add_uwsgi_header(struct corerouter_peer *peer, char *hh, size_t
}
}
else if (peer->key == uwsgi.hostname && hr->raw_body && !uwsgi_strncmp("ICE_URL", 7, hh, keylen)) {
peer->key = val;
peer->key_len = vallen;
if (vallen <= 0xff) {
memcpy(peer->key, val, vallen);
peer->key_len = vallen;
}
}
#ifdef UWSGI_ZLIB
@@ -345,7 +352,7 @@ int http_headers_parse(struct corerouter_peer *peer) {
// SERVER_NAME
if (!uhttp.server_name_as_http_host && uwsgi_buffer_append_keyval(out, "SERVER_NAME", 11, uwsgi.hostname, uwsgi.hostname_len)) return -1;
peer->key = uwsgi.hostname;
memcpy(peer->key, uwsgi.hostname, uwsgi.hostname_len);
peer->key_len = uwsgi.hostname_len;
// SERVER_PORT
@@ -461,10 +468,12 @@ int http_headers_parse(struct corerouter_peer *peer) {
if (uwsgi_starts_with("rtsp://", 7, hr->path_info, hr->path_info_len)) {
char *slash = memchr(hr->path_info + 7, '/', hr->path_info_len - 7);
if (!slash) return -1;
peer->key = hr->path_info + 7;
peer->key_len = slash - (hr->path_info + 7);
// override PATH_INFO
if (uwsgi_buffer_append_keyval(out, "PATH_INFO", 9, slash, hr->path_info_len - (7 + peer->key_len))) return -1;
if (slash - (hr->path_info + 7) <= 0xff) {
peer->key_len = slash - (hr->path_info + 7);
memcpy(peer->key, hr->path_info + 7, peer->key_len);
// override PATH_INFO
if (uwsgi_buffer_append_keyval(out, "PATH_INFO", 9, slash, hr->path_info_len - (7 + peer->key_len))) return -1;
}
}
}
@@ -831,10 +840,16 @@ ssize_t http_parse(struct corerouter_peer *main_peer) {
if (new_peer->instance_address_len == 0)
return -1;
// fix modifiers
if (uhttp.modifier1)
new_peer->modifier1 = uhttp.modifier1;
if (uhttp.modifier2)
new_peer->modifier2 = uhttp.modifier2;
uint16_t pktsize = new_peer->out->pos-4;
// fix modifiers
new_peer->out->buf[0] = new_peer->session->main_peer->modifier1;
new_peer->out->buf[3] = new_peer->session->main_peer->modifier2;
new_peer->out->buf[0] = new_peer->modifier1;
new_peer->out->buf[3] = new_peer->modifier2;
// fix pktsize
new_peer->out->buf[1] = (uint8_t) (pktsize & 0xff);
new_peer->out->buf[2] = (uint8_t) ((pktsize >> 8) & 0xff);
@@ -852,6 +867,12 @@ ssize_t http_parse(struct corerouter_peer *main_peer) {
if (uwsgi_buffer_append(new_peer->out, main_peer->in->buf + hr->headers_size + 1, hr->remains)) return -1;
}
if (new_peer->modifier1 == 123) {
// reset modifier1 to 0
new_peer->out->buf[0] = 0;
hr->raw_body = 1;
}
if (hr->websockets > 2 && hr->websocket_key_len > 0) {
hr->raw_body = 1;
}
@@ -976,9 +997,6 @@ int http_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket
// set the retry hook
cs->retry = hr_retry;
struct http_session *hr = (struct http_session *) cs;
// set the modifier1
cs->main_peer->modifier1 = uhttp.modifier1;
cs->main_peer->modifier2 = uhttp.modifier2;
// default hook
cs->main_peer->last_hook_read = hr_read;
+26 -11
View File
@@ -10,6 +10,14 @@
extern struct uwsgi_http uhttp;
// taken from nginx
static void hr_ssl_clear_errors() {
while (ERR_peek_error()) {
(void) ERR_get_error();
}
ERR_clear_error();
}
void uwsgi_opt_https(char *opt, char *value, void *cr) {
struct uwsgi_corerouter *ucr = (struct uwsgi_corerouter *) cr;
char *client_ca = NULL;
@@ -172,9 +180,9 @@ int hr_https_add_vars(struct http_session *hr, struct corerouter_peer *peer, str
if (uwsgi_buffer_append_keyval(out, "HTTPS", 5, "on", 2)) return -1;
#ifdef SSL_CTRL_SET_TLSEXT_HOSTNAME
const char *servername = SSL_get_servername(hr->ssl, TLSEXT_NAMETYPE_host_name);
if (servername) {
peer->key = (char *) servername;
peer->key_len = strlen(servername);
if (servername && strlen(servername) <= 0xff) {
peer->key_len = strlen(servername);
memcpy(peer->key, servername, peer->key_len) ;
}
#endif
hr->ssl_client_cert = SSL_get_peer_certificate(hr->ssl);
@@ -234,7 +242,7 @@ void hr_session_ssl_close(struct corerouter_session *cs) {
#endif
// clear the errors (otherwise they could be propagated)
ERR_clear_error();
hr_ssl_clear_errors();
SSL_free(hr->ssl);
}
@@ -270,6 +278,8 @@ ssize_t hr_ssl_write(struct corerouter_peer *main_peer) {
struct corerouter_session *cs = main_peer->session;
struct http_session *hr = (struct http_session *) cs;
hr_ssl_clear_errors();
int ret = SSL_write(hr->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;
@@ -294,9 +304,11 @@ ssize_t hr_ssl_write(struct corerouter_peer *main_peer) {
}
return ret;
}
if (ret == 0) return 0;
int err = SSL_get_error(hr->ssl, ret);
if (err == SSL_ERROR_ZERO_RETURN || err == 0) return 0;
if (err == SSL_ERROR_WANT_READ) {
cr_reset_hooks_and_read(main_peer, hr_ssl_write);
return 1;
@@ -322,6 +334,8 @@ ssize_t hr_ssl_read(struct corerouter_peer *main_peer) {
struct corerouter_session *cs = main_peer->session;
struct http_session *hr = (struct http_session *) cs;
hr_ssl_clear_errors();
// try to always leave 4k available
if (uwsgi_buffer_ensure(main_peer->in, uwsgi.page_size)) return -1;
int ret = SSL_read(hr->ssl, main_peer->in->buf + main_peer->in->pos, main_peer->in->len - main_peer->in->pos);
@@ -350,9 +364,11 @@ ssize_t hr_ssl_read(struct corerouter_peer *main_peer) {
#endif
return http_parse(main_peer);
}
if (ret == 0) return 0;
int err = SSL_get_error(hr->ssl, ret);
if (err == SSL_ERROR_ZERO_RETURN || err == 0) return 0;
if (err == SSL_ERROR_WANT_READ) {
cr_reset_hooks_and_read(main_peer, hr_ssl_read);
return 1;
@@ -381,18 +397,17 @@ ssize_t hr_ssl_shutdown(struct corerouter_peer *peer) {
struct corerouter_session *cs = peer->session;
struct http_session *hr = (struct http_session *) cs;
int ret = SSL_shutdown(hr->ssl);
if (ret < 0) return -1;
if (ret == 1) return 0;
hr_ssl_clear_errors();
int ret = SSL_shutdown(hr->ssl);
int err = 0;
if (ERR_peek_error()) {
if (ret != 1 && ERR_peek_error()) {
err = SSL_get_error(hr->ssl, ret);
}
// no error, close the connection
if (err == 0 || err == SSL_ERROR_ZERO_RETURN) return -1;
if (ret == 1 || err == 0 || err == SSL_ERROR_ZERO_RETURN) return 0;
if (err == SSL_ERROR_WANT_READ) {
if (uwsgi_cr_set_hooks(peer, hr_ssl_shutdown, NULL)) return -1;
+4 -2
View File
@@ -556,8 +556,10 @@ static ssize_t spdy_inflate_http_headers(struct http_session *hr) {
}
if (!uwsgi_strncmp(cgi_name, nk_len, "HTTP_HOST", 9)) {
new_peer->key = new_peer->out->buf + (new_peer->out->pos - v_len);
new_peer->key_len = v_len;
if (v_len <= 0xff) {
memcpy(new_peer->key, new_peer->out->buf + (new_peer->out->pos - v_len), v_len);
new_peer->key_len = v_len;
}
}
else if (!uwsgi_strncmp(cgi_name, nk_len, "REQUEST_URI", 11)) {
char *path_info = new_peer->out->buf + (new_peer->out->pos - v_len);
+10 -5
View File
@@ -543,7 +543,7 @@ static int uwsgi_mono_request(struct wsgi_request *wsgi_req) {
wsgi_req->app_id = uwsgi_get_app_id(NULL, key, key_len, mono_plugin.modifier1);
// if it is -1, try to load a dynamic app
if (wsgi_req->app_id == -1) {
if (wsgi_req->app_id == -1 && key_len > 0) {
if (uwsgi.threads > 1) {
pthread_mutex_lock(&umono.lock_loader);
}
@@ -562,10 +562,15 @@ static int uwsgi_mono_request(struct wsgi_request *wsgi_req) {
if (wsgi_req->app_id == -1) {
uwsgi_500(wsgi_req);
uwsgi_log("--- unable to find Mono/ASP.NET application ---\n");
// nothing to clear/free
return UWSGI_OK;
if (!uwsgi.no_default_app && uwsgi.default_app > -1 && uwsgi_apps[uwsgi.default_app].modifier1 == mono_plugin.modifier1) {
wsgi_req->app_id = uwsgi.default_app;
}
else {
uwsgi_500(wsgi_req);
uwsgi_log("--- unable to find Mono/ASP.NET application ---\n");
// nothing to clear/free
return UWSGI_OK;
}
}
struct uwsgi_app *app = &uwsgi_apps[wsgi_req->app_id];
+11 -3
View File
@@ -77,7 +77,7 @@ static int sapi_uwsgi_ub_write(const char *str, uint str_length TSRMLS_DC)
return str_length;
}
static int sapi_uwsgi_send_headers(sapi_headers_struct *sapi_headers)
static int sapi_uwsgi_send_headers(sapi_headers_struct *sapi_headers TSRMLS_DC)
{
sapi_header_struct *h;
zend_llist_position pos;
@@ -135,7 +135,7 @@ static int sapi_uwsgi_read_post(char *buffer, uint count_bytes TSRMLS_DC)
}
static char *sapi_uwsgi_read_cookies(void)
static char *sapi_uwsgi_read_cookies(TSRMLS_D)
{
uint16_t len = 0;
struct wsgi_request *wsgi_req = (struct wsgi_request *) SG(server_context);
@@ -518,7 +518,7 @@ static int php_uwsgi_startup(sapi_module_struct *sapi_module)
}
}
static void sapi_uwsgi_log_message(char *message) {
static void sapi_uwsgi_log_message(char *message TSRMLS_DC) {
uwsgi_log("%s\n", message);
}
@@ -559,6 +559,10 @@ int uwsgi_php_init(void) {
struct uwsgi_string_list *pset = uphp.set;
struct uwsgi_string_list *append_config = uphp.append_config;
#ifdef ZTS
tsrm_startup(1, 1, 0, NULL);
#endif
sapi_startup(&uwsgi_sapi_module);
// applying custom options
@@ -663,6 +667,10 @@ int uwsgi_php_request(struct wsgi_request *wsgi_req) {
zend_file_handle file_handle;
#ifdef ZTS
TSRMLS_FETCH();
#endif
SG(server_context) = (void *) wsgi_req;
if (uwsgi_parse_vars(wsgi_req)) {
+54 -6
View File
@@ -25,11 +25,12 @@ XS(XS_error) {
psgi_check_args(0);
if (uwsgi.threads > 1) {
ST(0) = sv_bless(newRV(sv_newmortal()), ((HV **)wi->error)[wsgi_req->async_id]);
ST(0) = sv_bless(newRV_noinc(newSV(0)), ((HV **)wi->error)[wsgi_req->async_id]);
}
else {
ST(0) = sv_bless(newRV(sv_newmortal()), ((HV **)wi->error)[0]);
ST(0) = sv_bless(newRV_noinc(newSV(0)), ((HV **)wi->error)[0]);
}
sv_2mortal(ST(0));
XSRETURN(1);
}
@@ -41,11 +42,12 @@ XS(XS_input) {
psgi_check_args(0);
if (uwsgi.threads > 1) {
ST(0) = sv_bless(newRV(sv_newmortal()), ((HV **)wi->input)[wsgi_req->async_id]);
ST(0) = sv_bless(newRV_noinc(newSV(0)), ((HV **)wi->input)[wsgi_req->async_id]);
}
else {
ST(0) = sv_bless(newRV(sv_newmortal()), ((HV **)wi->input)[0]);
ST(0) = sv_bless(newRV_noinc(newSV(0)), ((HV **)wi->input)[0]);
}
sv_2mortal(ST(0));
XSRETURN(1);
}
@@ -80,11 +82,12 @@ XS(XS_stream)
SvREFCNT_dec(response);
if (uwsgi.threads > 1) {
ST(0) = sv_bless(newRV(sv_newmortal()), ((HV **)wi->stream)[wsgi_req->async_id]);
ST(0) = sv_bless(newRV_noinc(newSV(0)), ((HV **)wi->stream)[wsgi_req->async_id]);
}
else {
ST(0) = sv_bless(newRV(sv_newmortal()), ((HV **)wi->stream)[0]);
ST(0) = sv_bless(newRV_noinc(newSV(0)), ((HV **)wi->stream)[0]);
}
sv_2mortal(ST(0));
XSRETURN(1);
}
else {
@@ -271,6 +274,51 @@ nonworker:
newCONSTSUB(stash, "SPOOL_RETRY", newSViv(-1));
newCONSTSUB(stash, "SPOOL_IGNORE", newSViv(0));
HV *_opts = newHV();
int i;
for (i = 0; i < uwsgi.exported_opts_cnt; i++) {
if (hv_exists(_opts, uwsgi.exported_opts[i]->key, strlen(uwsgi.exported_opts[i]->key))) {
SV **value = hv_fetch(_opts, uwsgi.exported_opts[i]->key, strlen(uwsgi.exported_opts[i]->key), 0);
// last resort !!!
if (!value) {
uwsgi_log("[perl] WARNING !!! unable to build uwsgi::opt hash !!!\n");
goto end;
}
if (SvTYPE(SvRV(*value)) == SVt_PVAV) {
if (uwsgi.exported_opts[i]->value == NULL) {
av_push((AV *)SvRV(*value), newSViv(1));
}
else {
av_push((AV *)SvRV(*value), newSVpv(uwsgi.exported_opts[i]->value, 0));
}
}
else {
AV *_opt_a = newAV();
av_push(_opt_a, SvREFCNT_inc(*value));
if (uwsgi.exported_opts[i]->value == NULL) {
av_push(_opt_a, newSViv(1));
}
else {
av_push(_opt_a, newSVpv(uwsgi.exported_opts[i]->value, 0));
}
hv_store(_opts, uwsgi.exported_opts[i]->key, strlen(uwsgi.exported_opts[i]->key), newRV_inc((SV *) _opt_a), 0);
}
}
else {
if (uwsgi.exported_opts[i]->value == NULL) {
hv_store(_opts, uwsgi.exported_opts[i]->key, strlen(uwsgi.exported_opts[i]->key), newSViv(1), 0);
}
else {
hv_store(_opts, uwsgi.exported_opts[i]->key, strlen(uwsgi.exported_opts[i]->key), newSVpv(uwsgi.exported_opts[i]->value, 0), 0);
}
}
}
newCONSTSUB(stash, "opt", newRV_inc((SV *) _opts));
end:
init_perl_embedded_module();
}
+31 -10
View File
@@ -114,7 +114,7 @@ SV *uwsgi_perl_obj_new(char *class, size_t class_len) {
SPAGAIN;
newobj = POPs;
newobj = SvREFCNT_inc(POPs);
PUTBACK;
FREETMPS;
LEAVE;
@@ -661,15 +661,21 @@ void uwsgi_perl_after_request(struct wsgi_request *wsgi_req) {
if (SvTRUE(*harakiri)) wsgi_req->async_plagued = 1;
}
// Free the $env hash
SvREFCNT_dec(wsgi_req->async_environ);
// async plagued could be defined in other areas...
if (wsgi_req->async_plagued) {
uwsgi_log("*** psgix.harakiri.commit requested ***\n");
// Before we call exit(0) we'll run the
// uwsgi_perl_atexit() hook which'll properly tear
// down the interpreter.
// mark the request as ended (otherwise the atexit hook will be skipped)
uwsgi.workers[uwsgi.mywid].cores[wsgi_req->async_id].in_request = 0;
goodbye_cruel_world();
}
// clear the env
SvREFCNT_dec(wsgi_req->async_environ);
// now we can check for changed files
if (uperl.auto_reload) {
time_t now = uwsgi_now();
@@ -803,24 +809,39 @@ void uwsgi_perl_run_hook(SV *hook) {
}
static void uwsgi_perl_atexit() {
int i;
if (uwsgi.mywid == 0) goto realstuff;
// if hijacked do not run atexit hooks
// if hijacked do not run atexit hooks -- TODO: explain why
// not.
if (uwsgi.workers[uwsgi.mywid].hijacked)
return;
goto destroyperl;
// if busy do not run atexit hooks
// if busy do not run atexit hooks (as this part could be called in a signal handler
// while a subroutine is running)
if (uwsgi_worker_is_busy(uwsgi.mywid))
return;
// managing atexit in async mode is a real pain...skip it for now
if (uwsgi.async > 1)
return;
realstuff:
if (uperl.atexit) {
uwsgi_perl_run_hook(uperl.atexit);
}
destroyperl:
// We must free our perl context(s) so any DESTROY hooks
// etc. will run.
for(i=0;i<uwsgi.threads;i++) {
PERL_SET_CONTEXT(uperl.main[i]);
// Destroy the PerlInterpreter, see "perldoc perlembed"
perl_destruct(uperl.main[i]);
perl_free(uperl.main[i]);
}
PERL_SYS_TERM();
free(uperl.main);
}
static uint64_t uwsgi_perl_rpc(void *func, uint8_t argc, char **argv, uint16_t argvs[], char **buffer) {
+8
View File
@@ -318,6 +318,13 @@ XS(XS_alarm) {
XSRETURN_UNDEF;
}
XS(XS_worker_id) {
dXSARGS;
psgi_check_args(0);
ST(0) = newSViv(uwsgi.mywid);
XSRETURN(1);
}
XS(XS_async_connect) {
dXSARGS;
@@ -1042,5 +1049,6 @@ void init_perl_embedded_module() {
psgi_xs(spool);
psgi_xs(add_var);
psgi_xs(worker_id);
}
+22 -5
View File
@@ -66,17 +66,34 @@ static int uwsgi_pypy_init() {
}
else {
if (upypy.home) {
// first try with /bin way:
#ifdef __CYGWIN__
char *libpath = uwsgi_concat2(upypy.home, "/libpypy-c.dll");
char *libpath = uwsgi_concat2(upypy.home, "/bin/libpypy-c.dll");
#elif defined(__APPLE__)
char *libpath = uwsgi_concat2(upypy.home, "/libpypy-c.dylib");
char *libpath = uwsgi_concat2(upypy.home, "/bin/libpypy-c.dylib");
#else
char *libpath = uwsgi_concat2(upypy.home, "/libpypy-c.so");
char *libpath = uwsgi_concat2(upypy.home, "/bin/libpypy-c.so");
#endif
if (uwsgi_file_exists(libpath)) {
upypy.handler = dlopen(libpath, RTLD_NOW | RTLD_GLOBAL);
upypy.handler = dlopen(libpath, RTLD_NOW | RTLD_GLOBAL);
}
free(libpath);
// fallback to old-style way
if (!upypy.handler) {
#ifdef __CYGWIN__
char *libpath = uwsgi_concat2(upypy.home, "/libpypy-c.dll");
#elif defined(__APPLE__)
char *libpath = uwsgi_concat2(upypy.home, "/libpypy-c.dylib");
#else
char *libpath = uwsgi_concat2(upypy.home, "/libpypy-c.so");
#endif
if (uwsgi_file_exists(libpath)) {
upypy.handler = dlopen(libpath, RTLD_NOW | RTLD_GLOBAL);
}
free(libpath);
}
free(libpath);
}
// fallback to standard library search path
if (!upypy.handler) {
+8
View File
@@ -1425,7 +1425,15 @@ void *uwsgi_python_autoreloader_thread(void *foobar) {
int found = 0;
struct uwsgi_string_list *usl = up.auto_reload_ignore;
while(usl) {
#ifdef PYTHREE
PyObject *zero = PyUnicode_AsUTF8String(mod_name);
char *str_mod_name = PyString_AsString(zero);
int ret_cmp = strcmp(usl->value, str_mod_name);
Py_DECREF(zero);
if (!ret_cmp) {
#else
if (!strcmp(usl->value, PyString_AsString(mod_name))) {
#endif
found = 1;
break;
}
+1 -1
View File
@@ -245,7 +245,7 @@ static int rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_ga
peer->last_hook_read = rr_instance_read;
// use the address as hostname
peer->key = cs->ugs->name;
memcpy(peer->key, cs->ugs->name, cs->ugs->name_len);
peer->key_len = cs->ugs->name_len;
// the mapper hook
+9
View File
@@ -30,6 +30,9 @@ static int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_r
if (ur->custom & 0x02) {
ub = uwsgi_buffer_new(uwsgi.page_size);
}
else if (ur->custom & 0x04) {
ub = uwsgi_to_http_dumb(wsgi_req, ur->data2, ur->data2_len, ub_url ? ub_url->buf : NULL, ub_url ? ub_url->pos : 0);
}
else {
ub = uwsgi_to_http(wsgi_req, ur->data2, ur->data2_len, ub_url ? ub_url->buf : NULL, ub_url ? ub_url->pos : 0);
}
@@ -133,10 +136,16 @@ static int uwsgi_router_http_connect(struct uwsgi_route *ur, char *args) {
return uwsgi_router_http(ur, args);
}
static int uwsgi_router_httpdumb(struct uwsgi_route *ur, char *args) {
ur->custom = 0x04;
return uwsgi_router_http(ur, args);
}
static void router_http_register(void) {
uwsgi_register_router("http", uwsgi_router_http);
uwsgi_register_router("httpdumb", uwsgi_router_httpdumb);
uwsgi_register_router("proxyhttp", uwsgi_router_proxyhttp);
uwsgi_register_router("httpconnect", uwsgi_router_http_connect);
uwsgi_register_router("proxyhttpconnect", uwsgi_router_proxyhttp_connect);
+3 -3
View File
@@ -288,15 +288,15 @@ static ssize_t sr_read(struct corerouter_peer *main_peer) {
// set default peer hook
peer->last_hook_read = sr_instance_read;
// use the address as hostname
peer->key = cs->ugs->name;
memcpy(peer->key, cs->ugs->name, cs->ugs->name_len);
peer->key_len = cs->ugs->name_len;
#ifdef SSL_CTRL_SET_TLSEXT_HOSTNAME
if (usr.sni) {
const char *servername = SSL_get_servername(sr->ssl, TLSEXT_NAMETYPE_host_name);
if (servername) {
peer->key = (char *) servername;
if (servername && strlen(servername) <= 0xff) {
peer->key_len = strlen(servername);
memcpy(peer->key, servername, peer->key_len);
}
}
#endif
+1 -1
View File
@@ -10,7 +10,7 @@ ssize_t uwsgi_systemd_logger(struct uwsgi_logger *ul, char *message, size_t len)
for(i=0;i<len;i++) {
if (message[i] == '\n') {
message[i] = 0;
sd_journal_print(LOG_INFO, base);
sd_journal_print(LOG_INFO, "%s", base);
base = message+i+1;
}
}
+90 -6
View File
@@ -480,6 +480,96 @@ static void uwsgi_httpize_var(char *buf, size_t len) {
}
}
struct uwsgi_buffer *uwsgi_to_http_dumb(struct wsgi_request *wsgi_req, char *host, uint16_t host_len, char *uri, uint16_t uri_len) {
struct uwsgi_buffer *ub = uwsgi_buffer_new(4096);
if (uwsgi_buffer_append(ub, wsgi_req->method, wsgi_req->method_len)) goto clear;
if (uwsgi_buffer_append(ub, " ", 1)) goto clear;
if (uri_len && uri) {
if (uwsgi_buffer_append(ub, uri, uri_len)) goto clear;
}
else {
if (uwsgi_buffer_append(ub, wsgi_req->uri, wsgi_req->uri_len)) goto clear;
}
if (uwsgi_buffer_append(ub, " ", 1)) goto clear;
if (uwsgi_buffer_append(ub, wsgi_req->protocol, wsgi_req->protocol_len)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
int i;
char *x_forwarded_for = NULL;
size_t x_forwarded_for_len = 0;
// start adding headers
for(i=0;i<wsgi_req->var_cnt;i++) {
if (!uwsgi_starts_with(wsgi_req->hvec[i].iov_base, wsgi_req->hvec[i].iov_len, "HTTP_", 5)) {
char *header = wsgi_req->hvec[i].iov_base+5;
size_t header_len = wsgi_req->hvec[i].iov_len-5;
if (host && !uwsgi_strncmp(header, header_len, "HOST", 4)) goto next;
if (!uwsgi_strncmp(header, header_len, "X_FORWARDED_FOR", 15)) {
x_forwarded_for = wsgi_req->hvec[i+1].iov_base;
x_forwarded_for_len = wsgi_req->hvec[i+1].iov_len;
goto next;
}
if (uwsgi_buffer_append(ub, header, header_len)) goto clear;
// transofmr uwsgi var to http header
uwsgi_httpize_var((ub->buf+ub->pos) - header_len, header_len);
if (uwsgi_buffer_append(ub, ": ", 2)) goto clear;
if (uwsgi_buffer_append(ub, wsgi_req->hvec[i+1].iov_base, wsgi_req->hvec[i+1].iov_len)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
}
next:
i++;
}
// append custom Host (if needed)
if (host) {
if (uwsgi_buffer_append(ub, "Host: ", 6)) goto clear;
if (uwsgi_buffer_append(ub, host, host_len)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
}
if (wsgi_req->content_type_len > 0) {
if (uwsgi_buffer_append(ub, "Content-Type: ", 14)) goto clear;
if (uwsgi_buffer_append(ub, wsgi_req->content_type, wsgi_req->content_type_len)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
}
if (wsgi_req->post_cl > 0) {
if (uwsgi_buffer_append(ub, "Content-Length: ", 16)) goto clear;
if (uwsgi_buffer_num64(ub, (int64_t) wsgi_req->post_cl)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n", 2)) goto clear;
}
// append required headers
if (uwsgi_buffer_append(ub, "X-Forwarded-For: ", 17)) goto clear;
if (x_forwarded_for_len > 0) {
if (uwsgi_buffer_append(ub, x_forwarded_for, x_forwarded_for_len)) goto clear;
if (uwsgi_buffer_append(ub, ", ", 2)) goto clear;
}
if (uwsgi_buffer_append(ub, wsgi_req->remote_addr, wsgi_req->remote_addr_len)) goto clear;
if (uwsgi_buffer_append(ub, "\r\n\r\n", 4)) goto clear;
return ub;
clear:
uwsgi_buffer_destroy(ub);
return NULL;
}
struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *wsgi_req, char *host, uint16_t host_len, char *uri, uint16_t uri_len) {
struct uwsgi_buffer *ub = uwsgi_buffer_new(4096);
@@ -513,14 +603,8 @@ struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *wsgi_req, char *host, ui
// remove dangerous headers
if (!uwsgi_strncmp(header, header_len, "CONNECTION", 10)) goto next;
if (!uwsgi_strncmp(header, header_len, "KEEP_ALIVE", 10)) goto next;
if (!uwsgi_strncmp(header, header_len, "ACCEPT_ENCODING", 15)) goto next;
if (!uwsgi_strncmp(header, header_len, "TE", 2)) goto next;
if (!uwsgi_strncmp(header, header_len, "TRAILER", 7)) goto next;
if (!uwsgi_strncmp(header, header_len, "IF_MATCH", 8)) goto next;
if (!uwsgi_strncmp(header, header_len, "IF_MODIFIED_SINCE", 17)) goto next;
if (!uwsgi_strncmp(header, header_len, "IF_RANGE", 8)) goto next;
if (!uwsgi_strncmp(header, header_len, "IF_UNMODIFIED_SINCE", 19)) goto next;
if (!uwsgi_strncmp(header, header_len, "IF_NONE_MATCH", 13)) goto next;
if (!uwsgi_strncmp(header, header_len, "X_FORWARDED_FOR", 15)) {
x_forwarded_for = wsgi_req->hvec[i+1].iov_base;
x_forwarded_for_len = wsgi_req->hvec[i+1].iov_len;
+10
View File
@@ -0,0 +1,10 @@
#uwsgi --psgi t/perl/active_workers_signal.pl -s :3031 --perl-no-plack --timer "17 3" -p 8 --cheap --idle 10
my $handler = sub {
print "hello i am the signal handler on worker ".uwsgi::worker_id()."\n";
};
uwsgi::register_signal(17, 'active-workers', $handler);
my $app = sub {
};
+19
View File
@@ -0,0 +1,19 @@
use strict;
use warnings;
{
package psgix::harakiri::tester;
sub DESTROY { print STDERR "$$: Calling DESTROY\n" }
}
sub {
my $env = shift;
die "PANIC: We should support psgix.harakiri here" unless $env->{'psgix.harakiri'};
$env->{'psgix.harakiri.tester'} = bless {} => 'psgix::harakiri::tester';
my $harakiri = $env->{QUERY_STRING};
$env->{'psgix.harakiri.commit'} = $harakiri ? 1 : 0;
return [200, [], [ $harakiri ? "We are about to destroy ourselves\n" : "We will live for another request\n" ]];
}
+1 -1
View File
@@ -2,7 +2,7 @@ Gem::Specification.new do |s|
s.name = 'uwsgi'
s.license = 'GPL-2'
s.version = `python -c "import uwsgiconfig as uc; print uc.uwsgi_version"`.sub(/-dev-.*/,'')
s.date = '2014-10-26'
s.date = '2014-12-30'
s.summary = "uWSGI"
s.description = "The uWSGI server for Ruby/Rack"
s.authors = ["Unbit"]
+5
View File
@@ -2753,6 +2753,10 @@ struct uwsgi_server {
#endif
struct uwsgi_string_list *hook_post_fork;
// uWSGI 2.0.9
char *subscribe_with_modifier1;
struct uwsgi_string_list *pull_headers;
};
struct uwsgi_rpc {
@@ -4182,6 +4186,7 @@ struct uwsgi_buffer *uwsgi_buffer_from_file(char *);
ssize_t uwsgi_buffer_write_simple(struct wsgi_request *, struct uwsgi_buffer *);
struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *, char *, uint16_t, char *, uint16_t);
struct uwsgi_buffer *uwsgi_to_http_dumb(struct wsgi_request *, char *, uint16_t, char *, uint16_t);
ssize_t uwsgi_pipe(int, int, int);
ssize_t uwsgi_pipe_sized(int, int, size_t, int);
+3 -1
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '2.0.8'
uwsgi_version = '2.0.9'
import os
import re
@@ -1171,6 +1171,8 @@ class uConf(object):
else:
jsonconf = spcall("pkg-config --cflags yajl")
if jsonconf:
if jsonconf.endswith('include/yajl'):
jsonconf = jsonconf.rstrip('yajl')
self.cflags.append(jsonconf)
self.cflags.append("-DUWSGI_JSON")
self.gcc_list.append('core/json')