Compare commits

..
25 Commits
Author SHA1 Message Date
Unbit 0f397eeb55 fixed #305 2013-06-05 11:14:47 +02:00
Unbit bf929398a9 updated gemspec for 1.9.12 2013-06-05 09:19:57 +02:00
Unbit 4f8a0f753e added router_hash and extended support for memcached and redis 2013-06-05 06:26:39 +02:00
Unbit 7db5da9566 fixed support for newer cygwin 2013-06-04 17:13:18 +02:00
Unbit 4d1389805c added redis router code 2013-06-04 13:20:47 +02:00
Unbit ee87e2b946 memcached and redis routers are now builtin by default 2013-06-04 13:07:51 +02:00
Unbit 6356d78768 completed memcachedstore 2013-06-04 12:04:42 +02:00
Unbit afdf7e740e prepare for memcachedstore 2013-06-04 11:40:19 +02:00
Unbit 51ce3ddf78 added offloading to memcached internal router 2013-06-04 11:18:02 +02:00
Unbit 192186056e added pipe offload engine 2013-06-04 10:34:19 +02:00
Unbit 336a7febd7 prepate for 1.9.12 2013-06-03 15:35:17 +02:00
Unbit 63e73d1b73 added --touch-signal 2013-06-03 11:15:39 +02:00
Unbit a02ab9a723 automatically search for 'application' method in jwsgi 2013-06-03 06:33:48 +02:00
Unbit 95f15ab9dd added a note for python logger 2013-06-02 19:20:35 +02:00
Unbit 5d0ea0066a fixed fastcgi body read 2013-06-02 12:52:13 +02:00
Unbit 9ff6b2da9b improved tmp file generation and first implementation of offload transformation 2013-06-02 11:12:28 +02:00
Roberta Giordano 6ac7e27847 added seek() support to jvm input 2013-06-01 19:53:29 +02:00
Roberta Giordano 4ed1c5e747 Merge branch 'master' of https://github.com/unbit/uwsgi 2013-06-01 19:42:29 +02:00
Roberta Giordano bae3bc096f added readLine featute to jvm input 2013-06-01 19:42:11 +02:00
unbit 9d73683de2 Merge pull request #298 from prymitive/csfix
try to escape json
2013-06-01 08:15:19 -07:00
Roberta Giordano f5fec0ab71 added support for jvm list and uwsgi.opt 2013-06-01 16:32:39 +02:00
Łukasz Mierzwa 9909ab10c4 try to escape json 2013-06-01 15:53:35 +02:00
Roberta Giordano 65cff01fee implemented non-static JWSGI apps 2013-06-01 12:20:19 +02:00
unbit 10190a3012 Merge pull request #301 from jimfunk/master
Fix signed/unsigned comparison in emperor_amqp plugin
2013-05-31 21:56:23 -07:00
James Oakley 79ab696f4e Fix signed/unsigned comparison 2013-05-31 13:24:49 -07:00
28 changed files with 1302 additions and 68 deletions
+1 -1
View File
@@ -16,7 +16,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, sslrouter, spooler, cheaper_busyness, symcall, transformation_tofile, transformation_gzip, transformation_chunked
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, spooler, cheaper_busyness, symcall, transformation_tofile, transformation_gzip, transformation_chunked, transformation_offload, router_memcached, router_redis, router_hash
as_shared_library = false
locking = auto
+18
View File
@@ -485,6 +485,18 @@ int master_loop(char **argv, char **environ) {
usl = usl->next;
}
uwsgi_check_touches(uwsgi.touch_exec);
// update signal touches
usl = uwsgi.touch_signal;
while(usl) {
char *space = strchr(usl->value, ' ');
if (space) {
*space = 0;
usl->len = strlen(usl->value);
usl->custom_ptr = space + 1;
}
usl = usl->next;
}
uwsgi_check_touches(uwsgi.touch_signal);
// setup cheaper algos (can be stacked)
uwsgi.cheaper_algo = uwsgi_cheaper_algo_spare;
@@ -769,6 +781,12 @@ int master_loop(char **argv, char **environ) {
uwsgi_log_verbose("[uwsgi-touch-exec] running %s\n", touched);
}
}
touched = uwsgi_check_touches(uwsgi.touch_signal);
if (touched) {
uint8_t signum = atoi(touched);
uwsgi_route_signal(signum);
uwsgi_log_verbose("[uwsgi-touch-signal] raising %u\n", signum);
}
}
continue;
+16 -2
View File
@@ -696,8 +696,17 @@ struct uwsgi_stats *uwsgi_master_generate_stats() {
while (ud) {
if (uwsgi_stats_object_open(us))
goto end;
if (uwsgi_stats_keyval_comma(us, "cmd", ud->command))
// allocate 2x the size of original command
// in case we need to escape all chars
char *cmd = uwsgi_malloc(strlen(ud->command)*2);
escape_json(ud->command, strlen(ud->command), cmd);
if (uwsgi_stats_keyval_comma(us, "cmd", cmd)) {
free(cmd);
goto end;
}
free(cmd);
if (uwsgi_stats_keylong_comma(us, "pid", (unsigned long long) (ud->pid < 0) ? 0 : ud->pid))
goto end;
if (uwsgi_stats_keylong(us, "respawns", (unsigned long long) ud->respawns ? 0 : ud->respawns))
@@ -1106,8 +1115,13 @@ struct uwsgi_stats *uwsgi_master_generate_stats() {
if (uwsgi_stats_keyslong_comma(us, "week", (long long) ucron->week))
goto end;
if (uwsgi_stats_keyval_comma(us, "command", ucron->command))
char *cmd = uwsgi_malloc(strlen(ucron->command)*2);
escape_json(ucron->command, strlen(ucron->command), cmd);
if (uwsgi_stats_keyval_comma(us, "command", cmd)) {
free(cmd);
goto end;
}
free(cmd);
if (uwsgi_stats_keylong_comma(us, "unique", (unsigned long long) ucron->unique))
goto end;
+92 -2
View File
@@ -58,6 +58,22 @@ static int uwsgi_offload_enqueue(struct wsgi_request *wsgi_req, struct uwsgi_off
return 0;
}
/*
pipe offload engine:
fd -> the file descriptor to read from
len -> amount of data to transfer
*/
static int u_offload_pipe_prepare(struct wsgi_request *wsgi_req, struct uwsgi_offload_request *uor) {
if (uor->fd < 0 || !uor->len) {
return -1;
}
return 0;
}
/*
memory offload engine:
@@ -66,7 +82,7 @@ static int uwsgi_offload_enqueue(struct wsgi_request *wsgi_req, struct uwsgi_off
*/
int u_offload_memory_prepare(struct wsgi_request *wsgi_req, struct uwsgi_offload_request *uor) {
static int u_offload_memory_prepare(struct wsgi_request *wsgi_req, struct uwsgi_offload_request *uor) {
if (!uor->buf || !uor->len) {
return -1;
@@ -83,7 +99,7 @@ int u_offload_memory_prepare(struct wsgi_request *wsgi_req, struct uwsgi_offload
*/
int u_offload_transfer_prepare(struct wsgi_request *wsgi_req, struct uwsgi_offload_request *uor) {
static int u_offload_transfer_prepare(struct wsgi_request *wsgi_req, struct uwsgi_offload_request *uor) {
if (!uor->name) {
return -1;
@@ -343,6 +359,71 @@ static int u_offload_sendfile_do(struct uwsgi_thread *ut, struct uwsgi_offload_r
}
/*
pipe offloading
status:
0 -> waiting for data on fd
1 -> waiting for write to s
*/
static int u_offload_pipe_do(struct uwsgi_thread *ut, struct uwsgi_offload_request *uor, int fd) {
ssize_t rlen;
// setup
if (fd == -1) {
event_queue_add_fd_read(ut->queue, uor->fd);
return 0;
}
switch(uor->status) {
// read event from fd
case 0:
if (!uor->buf) {
uor->buf = uwsgi_malloc(4096);
}
rlen = read(uor->fd, uor->buf, 4096);
if (rlen > 0) {
uor->to_write = rlen;
uor->pos = 0;
if (event_queue_del_fd(ut->queue, uor->fd, event_queue_read())) return -1;
if (event_queue_add_fd_write(ut->queue, uor->s)) return -1;
uor->status = 1;
return 0;
}
if (rlen < 0) {
uwsgi_offload_retry
uwsgi_error("u_offload_pipe_do() -> read()");
}
return -1;
// write event on s
case 1:
rlen = write(uor->s, uor->buf + uor->pos, uor->to_write);
if (rlen > 0) {
uor->to_write -= rlen;
uor->pos += rlen;
if (uor->to_write == 0) {
if (event_queue_del_fd(ut->queue, uor->s, event_queue_write())) return -1;
if (event_queue_add_fd_read(ut->queue, uor->fd)) return -1;
uor->status = 0;
}
return 0;
}
else if (rlen < 0) {
uwsgi_offload_retry
uwsgi_error("u_offload_pipe_do() -> write()");
}
return -1;
default:
break;
}
return -1;
}
/*
the offload task starts soon after the call to connect()
@@ -541,6 +622,7 @@ void uwsgi_offload_engines_register_all() {
uwsgi.offload_engine_sendfile = uwsgi_offload_register_engine("sendfile", u_offload_sendfile_prepare, u_offload_sendfile_do);
uwsgi.offload_engine_transfer = uwsgi_offload_register_engine("transfer", u_offload_transfer_prepare, u_offload_transfer_do);
uwsgi.offload_engine_memory = uwsgi_offload_register_engine("memory", u_offload_memory_prepare, u_offload_memory_do);
uwsgi.offload_engine_pipe = uwsgi_offload_register_engine("pipe", u_offload_pipe_prepare, u_offload_pipe_do);
}
int uwsgi_offload_request_sendfile_do(struct wsgi_request *wsgi_req, int fd, size_t len) {
@@ -566,3 +648,11 @@ int uwsgi_offload_request_memory_do(struct wsgi_request *wsgi_req, char *buf, si
uor.len = len;
return uwsgi_offload_run(wsgi_req, &uor, NULL);
}
int uwsgi_offload_request_pipe_do(struct wsgi_request *wsgi_req, int fd, size_t len) {
struct uwsgi_offload_request uor;
uwsgi_offload_setup(uwsgi.offload_engine_pipe, &uor, wsgi_req, 1);
uor.fd = fd;
uor.len = len;
return uwsgi_offload_run(wsgi_req, &uor, NULL);
}
+2 -2
View File
@@ -478,9 +478,9 @@ int uwsgi_postbuffer_do_in_disk(struct wsgi_request *wsgi_req) {
int upload_progress_fd = -1;
char *upload_progress_filename = NULL;
wsgi_req->post_file = tmpfile();
wsgi_req->post_file = uwsgi_tmpfile();
if (!wsgi_req->post_file) {
uwsgi_error("uwsgi_postbuffer_do_in_disk()/tmpfile()");
uwsgi_error("uwsgi_postbuffer_do_in_disk()/uwsgi_tmpfile()");
return -1;
}
+36
View File
@@ -37,3 +37,39 @@ ssize_t uwsgi_sendfile_do(int sockfd, int filefd, size_t pos, size_t len) {
#endif
}
/*
simple non blocking sendfile implementation
(generally used as fallback)
*/
int uwsgi_simple_sendfile(struct wsgi_request *wsgi_req, int fd, size_t pos, size_t len) {
wsgi_req->write_pos = 0;
for(;;) {
int ret = wsgi_req->socket->proto_sendfile(wsgi_req, fd, pos, len);
if (ret < 0) {
if (!uwsgi.ignore_write_errors) {
uwsgi_error("uwsgi_simple_sendfile()");
}
wsgi_req->write_errors++;
return -1;
}
if (ret == UWSGI_OK) {
break;
}
ret = uwsgi_wait_write_req(wsgi_req);
if (ret < 0) {
wsgi_req->write_errors++;
return -1;
}
if (ret == 0) {
uwsgi_log("uwsgi_simple_sendfile() TIMEOUT !!!\n");
wsgi_req->write_errors++;
return -1;
}
}
return 0;
}
+8 -5
View File
@@ -132,6 +132,12 @@ void uwsgi_free_transformations(struct wsgi_request *wsgi_req) {
if (current_ut->chunk) {
uwsgi_buffer_destroy(current_ut->chunk);
}
if (current_ut->ub) {
uwsgi_buffer_destroy(current_ut->ub);
}
if (current_ut->fd > -1) {
close(current_ut->fd);
}
ut = ut->next;
free(current_ut);
}
@@ -144,12 +150,9 @@ struct uwsgi_transformation *uwsgi_add_transformation(struct wsgi_request *wsgi_
ut = ut->next;
}
ut = uwsgi_malloc(sizeof(struct uwsgi_transformation));
ut = uwsgi_calloc(sizeof(struct uwsgi_transformation));
ut->func = func;
ut->is_final = 0;
ut->next = NULL;
ut->chunk = NULL;
ut->can_stream = 0;
ut->fd = -1;
ut->data = data;
if (old_ut) {
+44 -6
View File
@@ -3065,6 +3065,36 @@ void escape_shell_arg(char *src, size_t len, char *dst) {
*ptr++ = 0;
}
void escape_json(char *src, size_t len, char *dst) {
size_t i;
char *ptr = dst;
for (i = 0; i < len; i++) {
if (src[i] == '\t') {
*ptr++ = '\\';
*ptr++ = 't';
}
else if (src[i] == '\n') {
*ptr++ = '\\';
*ptr++ = 'n';
}
else if (src[i] == '\r') {
*ptr++ = '\\';
*ptr++ = 'r';
}
else if (src[i] == '"') {
*ptr++ = '\\';
*ptr++ = '"';
}
else {
*ptr++ = src[i];
}
}
*ptr++ = 0;
}
void http_url_decode(char *buf, uint16_t * len, char *dst) {
uint16_t i;
@@ -3221,14 +3251,22 @@ char *uwsgi_chomp(char *str) {
return str;
}
char *uwsgi_tmpname(char *base, char *id) {
char *template = uwsgi_concat3(base, "/", id);
if (mkstemp(template) < 0) {
free(template);
return NULL;
int uwsgi_tmpfd() {
char *tmpdir = getenv("TMPDIR");
if (!tmpdir) {
tmpdir = "/tmp";
}
char *template = uwsgi_concat2(tmpdir, "/uwsgiXXXXXX");
int fd = mkstemp(template);
unlink(template);
free(template);
return fd;
}
return template;
FILE *uwsgi_tmpfile() {
int fd = uwsgi_tmpfd();
if (fd < 0) return NULL;
return fdopen(fd, "w+");
}
int uwsgi_file_to_string_list(char *filename, struct uwsgi_string_list **list) {
+1
View File
@@ -349,6 +349,7 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"touch-logrotate", required_argument, 0, "trigger logrotation if the specified file is modified/touched", uwsgi_opt_add_string_list, &uwsgi.touch_logrotate, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
{"touch-logreopen", required_argument, 0, "trigger log reopen if the specified file is modified/touched", uwsgi_opt_add_string_list, &uwsgi.touch_logreopen, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
{"touch-exec", required_argument, 0, "run command when the specified file is modified/touched (syntax: file command)", uwsgi_opt_add_string_list, &uwsgi.touch_exec, UWSGI_OPT_MASTER},
{"touch-signal", required_argument, 0, "signal when the specified file is modified/touched (syntax: file signal)", uwsgi_opt_add_string_list, &uwsgi.touch_signal, UWSGI_OPT_MASTER},
{"propagate-touch", no_argument, 0, "over-engineering option for system with flaky signal mamagement", uwsgi_opt_true, &uwsgi.propagate_touch, 0},
{"limit-post", required_argument, 0, "limit request body", uwsgi_opt_set_64bit, &uwsgi.limit_post, 0},
{"no-orphans", no_argument, 0, "automatically kill workers if master dies (can be dangerous for availability)", uwsgi_opt_true, &uwsgi.no_orphans, 0},
+32
View File
@@ -408,3 +408,35 @@ int uwsgi_simple_wait_write_hook(int fd, int timeout) {
return ret;
}
/*
simplified write to client
(generally used as fallback)
*/
int uwsgi_simple_write(struct wsgi_request *wsgi_req, char *buf, size_t len) {
wsgi_req->write_pos = 0;
for(;;) {
int ret = wsgi_req->socket->proto_write(wsgi_req, buf, len);
if (ret < 0) {
if (!uwsgi.ignore_write_errors) {
uwsgi_error("uwsgi_simple_write()");
}
wsgi_req->write_errors++;
return -1;
}
if (ret == UWSGI_OK) {
break;
}
ret = uwsgi_wait_write_req(wsgi_req);
if (ret < 0) { wsgi_req->write_errors++; return -1;}
if (ret == 0) {
uwsgi_log("uwsgi_simple_write() TIMEOUT !!!\n");
wsgi_req->write_errors++;
return -1;
}
}
return 0;
}
+2 -1
View File
@@ -453,7 +453,8 @@ static int amqp_wait_connection_tune(int fd) {
static char *amqp_simple_get_frame(int fd, struct amqp_frame_header *fh) {
char *ptr = (char *) fh;
ssize_t len = 0, rlen;
size_t len = 0;
ssize_t rlen;
while(len < 7) {
rlen = recv(fd, ptr, 7-len, 0);
+5
View File
@@ -25,6 +25,7 @@ struct uwsgi_jvm {
jclass input_stream_class;
jclass file_class;
jclass hashmap_class;
jclass list_class;
jclass set_class;
jclass iterator_class;
@@ -63,6 +64,10 @@ jmethodID uwsgi_jvm_get_static_method_id_quiet(jclass, char *, char *);
jobject uwsgi_jvm_str(char *, size_t);
jobject uwsgi_jvm_hashmap(void);
int uwsgi_jvm_hashmap_put(jobject, jobject, jobject);
int uwsgi_jvm_hashmap_has(jobject, jobject);
jobject uwsgi_jvm_list(void);
int uwsgi_jvm_list_add(jobject, jobject);
jobject uwsgi_jvm_call_object(jobject, jmethodID, ...);
jobject uwsgi_jvm_call_object_static(jclass, jmethodID, ...);
+120
View File
@@ -215,15 +215,42 @@ JNIEXPORT jint JNICALL uwsgi_jvm_request_body_read_bytearray(JNIEnv *env, jobjec
return rlen;
}
JNIEXPORT jint JNICALL uwsgi_jvm_request_body_readline_bytearray(JNIEnv *env, jobject o, jobject b) {
struct wsgi_request *wsgi_req = current_wsgi_req();
ssize_t rlen = 0;
size_t len = uwsgi_jvm_array_len(b);
char *chunk = uwsgi_request_body_readline(wsgi_req, len, &rlen);
if (!chunk) {
uwsgi_jvm_throw_io("error reading request body");
return -1;
}
if (chunk == uwsgi.empty) {
return -1;
}
char *buf = (char *) (*ujvm_env)->GetByteArrayElements(ujvm_env, b, JNI_FALSE);
if (!buf) return -1;
memcpy(buf, chunk, rlen);
(*ujvm_env)->ReleaseByteArrayElements(ujvm_env, b, (jbyte *) buf, 0);
return rlen;
}
JNIEXPORT jint JNICALL uwsgi_jvm_request_body_available(JNIEnv *env, jobject o) {
struct wsgi_request *wsgi_req = current_wsgi_req();
return (jint) (wsgi_req->post_cl - wsgi_req->post_pos);
}
JNIEXPORT void JNICALL uwsgi_jvm_request_body_seek(JNIEnv *env, jobject o, jint pos) {
struct wsgi_request *wsgi_req = current_wsgi_req();
uwsgi_request_body_seek(wsgi_req, pos);
}
static JNINativeMethod uwsgi_jvm_request_body_methods[] = {
{"read", "()I", (void *) &uwsgi_jvm_request_body_read},
{"read", "([B)I", (void *) &uwsgi_jvm_request_body_read_bytearray},
{"readLine", "([B)I", (void *) &uwsgi_jvm_request_body_readline_bytearray},
{"available", "()I", (void *) &uwsgi_jvm_request_body_available},
{"seek", "(I)V", (void *) &uwsgi_jvm_request_body_seek},
};
static struct uwsgi_option uwsgi_jvm_options[] = {
@@ -684,6 +711,23 @@ void uwsgi_jvm_local_unref(jobject obj) {
(*ujvm_env)->DeleteLocalRef(ujvm_env, obj);
}
jobject uwsgi_jvm_list() {
// optimization
static jmethodID mid = 0;
if (!mid) {
mid = uwsgi_jvm_get_method_id(ujvm.list_class, "<init>", "()V");
if (!mid) return NULL;
}
jobject ll = (*ujvm_env)->NewObject(ujvm_env, ujvm.list_class, mid);
if (uwsgi_jvm_exception()) {
return NULL;
}
return ll;
}
jobject uwsgi_jvm_hashmap() {
// optimization
static jmethodID mid = 0;
@@ -712,6 +756,18 @@ int uwsgi_jvm_hashmap_put(jobject hm, jobject key, jobject value) {
return uwsgi_jvm_call(hm, mid, key, value);
}
int uwsgi_jvm_list_add(jobject ll, jobject value) {
// optimization
static jmethodID mid = 0;
if (!mid) {
mid = uwsgi_jvm_get_method_id(ujvm.list_class, "add", "(Ljava/lang/Object;)Z");
if (!mid) return -1;
}
return uwsgi_jvm_call(ll, mid, value);
}
jobject uwsgi_jvm_hashmap_get(jobject hm, jobject key) {
// optimization
static jmethodID mid = 0;
@@ -724,6 +780,21 @@ jobject uwsgi_jvm_hashmap_get(jobject hm, jobject key) {
return uwsgi_jvm_call_object(hm, mid, key);
}
int uwsgi_jvm_hashmap_has(jobject hm, jobject key) {
// optimization
static jmethodID mid = 0;
if (!mid) {
mid = uwsgi_jvm_get_method_id(ujvm.hashmap_class, "containsKey", "(Ljava/lang/Object;)Z");
if (!mid) return 0;
}
if (uwsgi_jvm_call_bool(hm, mid, key)) {
return 1;
}
return 0;
}
jobject uwsgi_jvm_iterator(jobject set) {
// optimization
static jmethodID mid = 0;
@@ -1010,6 +1081,9 @@ static void uwsgi_jvm_create(void) {
ujvm.hashmap_class = uwsgi_jvm_class("java/util/HashMap");
if (!ujvm.hashmap_class) exit(1);
ujvm.list_class = uwsgi_jvm_class("java/util/ArrayList");
if (!ujvm.list_class) exit(1);
ujvm.set_class = uwsgi_jvm_class("java/util/Set");
if (!ujvm.set_class) exit(1);
@@ -1026,6 +1100,52 @@ static void uwsgi_jvm_create(void) {
if (!uwsgi_class) {
exit(1);
}
/*
start filling uwsgi.opt
*/
jfieldID opt_fid = (*ujvm_env)->GetStaticFieldID(ujvm_env, uwsgi_class, "opt", "Ljava/util/HashMap;");
if (uwsgi_jvm_exception()) {
exit(1);
}
jobject opt_hm = uwsgi_jvm_hashmap();
int j;
for (j = 0; j < uwsgi.exported_opts_cnt; j++) {
jstring j_opt_key = uwsgi_jvm_str(uwsgi.exported_opts[j]->key, 0);
if (uwsgi_jvm_hashmap_has(opt_hm, (jobject) j_opt_key)) {
jobject j_opt_value = uwsgi_jvm_hashmap_get(opt_hm, (jobject) j_opt_key);
if (uwsgi_jvm_object_is_instance(j_opt_value, ujvm.list_class)) {
if (uwsgi.exported_opts[j]->value == NULL) {
uwsgi_jvm_list_add(j_opt_value, (jobject) 1);
}
else {
uwsgi_jvm_list_add(j_opt_value, uwsgi_jvm_str(uwsgi.exported_opts[j]->value, 0));
}
}
else {
jobject ll = uwsgi_jvm_list();
uwsgi_jvm_list_add(ll, j_opt_value);
if (uwsgi.exported_opts[j]->value == NULL) {
uwsgi_jvm_list_add(ll, (jobject) 1);
}
else {
uwsgi_jvm_list_add(ll, uwsgi_jvm_str(uwsgi.exported_opts[j]->value, 0));
}
uwsgi_jvm_hashmap_put(opt_hm, j_opt_key, ll);
}
}
else {
if (uwsgi.exported_opts[j]->value == NULL) {
uwsgi_jvm_hashmap_put(opt_hm, j_opt_key, (jobject) 1);
}
else {
uwsgi_jvm_hashmap_put(opt_hm, j_opt_key, uwsgi_jvm_str(uwsgi.exported_opts[j]->value, 0));
}
}
}
(*ujvm_env)->SetStaticObjectField(ujvm_env, uwsgi_class, opt_fid, opt_hm);
(*ujvm_env)->RegisterNatives(ujvm_env, uwsgi_class, uwsgi_jvm_api_methods, sizeof(uwsgi_jvm_api_methods)/sizeof(uwsgi_jvm_api_methods[0]));
if (uwsgi_jvm_exception()) {
exit(1);
+5
View File
@@ -1,11 +1,16 @@
import java.io.*;
import java.util.*;
public class uwsgi {
static HashMap<String,Object> opt;
public static class RequestBody extends InputStream {
public native int read();
public native int read(byte[] b);
public native int readLine(byte[] b);
public native int available();
public native void seek(int pos);
}
public interface SignalHandler {
+33 -10
View File
@@ -9,6 +9,7 @@ struct uwsgi_jwsgi {
char *app;
jmethodID app_mid;
jclass app_class;
jobject app_instance;
} ujwsgi;
static struct uwsgi_option uwsgi_jwsgi_options[] = {
@@ -20,7 +21,16 @@ static int uwsgi_jwsgi_add_request_item(jobject hm, char *key, uint16_t key_len,
jobject j_key = uwsgi_jvm_str(key, key_len);
if (!j_key) return -1;
jobject j_value = uwsgi_jvm_str(value, value_len);
jobject j_value = NULL;
// avoid clobbering vars
if (value_len > 0) {
j_value = uwsgi_jvm_str(value, value_len);
}
else {
char *tmp = uwsgi_str("");
j_value = uwsgi_jvm_str(tmp, 0);
free(tmp);
}
if (!j_value) {
uwsgi_jvm_local_unref(j_value);
return -1;
@@ -72,7 +82,12 @@ static int uwsgi_jwsgi_request(struct wsgi_request *wsgi_req) {
if (uwsgi_jwsgi_add_request_input(hm, "jwsgi.input", 11)) goto end;
response = uwsgi_jvm_call_object_static(ujwsgi.app_class, ujwsgi.app_mid, hm);
if (!ujwsgi.app_instance) {
response = uwsgi_jvm_call_object_static(ujwsgi.app_class, ujwsgi.app_mid, hm);
}
else {
response = uwsgi_jvm_call_object(ujwsgi.app_instance, ujwsgi.app_mid, hm);
}
if (!response) goto end;
if (uwsgi_jvm_array_len(response) != 3) {
@@ -135,23 +150,31 @@ static int uwsgi_jwsgi_setup() {
char *app = uwsgi_str(ujwsgi.app);
char *method = "application";
char *colon = strchr(app, ':');
if (!colon) {
uwsgi_log("invalid JWSGI app definition, must be class:method\n");
exit(1);
if (colon) {
*colon = 0;
method = colon + 1;
}
*colon = 0;
ujwsgi.app_class = uwsgi_jvm_class(app);
if (!ujwsgi.app_class) {
exit(1);
}
ujwsgi.app_mid = uwsgi_jvm_get_static_method_id(ujwsgi.app_class, colon+1, "(Ljava/util/HashMap;)[Ljava/lang/Object;");
if (!ujwsgi.app_mid) {
exit(1);
ujwsgi.app_mid = uwsgi_jvm_get_static_method_id_quiet(ujwsgi.app_class, method, "(Ljava/util/HashMap;)[Ljava/lang/Object;");
if (uwsgi_jvm_exception() || !ujwsgi.app_mid) {
jmethodID mid = uwsgi_jvm_get_method_id(ujwsgi.app_class, "<init>", "()V");
if (uwsgi_jvm_exception() || !mid) exit(1);
ujwsgi.app_instance = (*ujvm_env)->NewObject(ujvm_env, ujwsgi.app_class, mid);
if (uwsgi_jvm_exception() || !ujwsgi.app_instance) {
exit(1);
}
ujwsgi.app_mid = uwsgi_jvm_get_method_id(ujwsgi.app_class, method, "(Ljava/util/HashMap;)[Ljava/lang/Object;");
if (uwsgi_jvm_exception() || !ujwsgi.app_mid) {
exit(1);
}
}
uwsgi_log("JWSGI app \"%s\" loaded\n", ujwsgi.app);
+6 -1
View File
@@ -1756,7 +1756,12 @@ void uwsgi_python_harakiri(int wid) {
}
}
/*
# you can use this logger to offload logging to python
# be sure to configure it to not log to stderr otherwise you will generate a loop
import logging
logging.basicConfig(filename='/tmp/pippo.log')
*/
ssize_t uwsgi_python_logger(struct uwsgi_logger *ul, char *message, size_t len) {
if (!Py_IsInitialized()) return -1;
-1
View File
@@ -193,7 +193,6 @@ static int uwsgi_routing_func_cache(struct wsgi_request *wsgi_req, struct uwsgi_
if (wsgi_req->socket->can_offload && !ur->custom && !urcc->no_offload) {
if (!uwsgi_offload_request_memory_do(wsgi_req, value, valsize)) {
wsgi_req->via = UWSGI_VIA_OFFLOAD;
wsgi_req->status = 202;
return UWSGI_ROUTE_BREAK;
}
}
+146
View File
@@ -0,0 +1,146 @@
#include <uwsgi.h>
#ifdef UWSGI_ROUTING
extern struct uwsgi_server uwsgi;
/*
by Unbit
syntax:
route = /^foobar1(.*)/ hash:key=foo$1poo,algo=murmur2,var=MYHASH,items=node1;node2;node3
*/
struct uwsgi_router_hash_conf {
char *key;
size_t key_len;
char *var;
size_t var_len;
char *algo;
char *items;
size_t items_len;
};
static int uwsgi_routing_func_hash(struct wsgi_request *wsgi_req, struct uwsgi_route *ur){
struct uwsgi_router_hash_conf *urhc = (struct uwsgi_router_hash_conf *) ur->data2;
struct uwsgi_hash_algo *uha = uwsgi_hash_algo_get(urhc->algo);
if (!uha) {
uwsgi_log("[uwsgi-hash-router] unable to find hash algo \"%s\"\n", urhc->algo);
return UWSGI_ROUTE_BREAK;
}
char **subject = (char **) (((char *)(wsgi_req))+ur->subject);
uint16_t *subject_len = (uint16_t *) (((char *)(wsgi_req))+ur->subject_len);
struct uwsgi_buffer *ub = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urhc->key, urhc->key_len);
if (!ub) return UWSGI_ROUTE_BREAK;
uint32_t h = uha->func(ub->buf, ub->pos);
uwsgi_buffer_destroy(ub);
// now count the number of items
uint32_t items = 1;
size_t i, ilen = urhc->items_len;
for(i=0;i<ilen;i++) {
if (urhc->items[i] == ';') items++;
}
// skip last semicolon
if (urhc->items[ilen-1] == ';') items--;
uint32_t hashed_result = h % items;
uint32_t found = 0;
char *value = urhc->items;
uint16_t vallen = 0;
for(i=0;i<ilen;i++) {
if (!value) {
value = urhc->items + i;
}
if (urhc->items[i] == ';') {
if (found == hashed_result) {
vallen = (urhc->items+i) - value;
break;
}
value = NULL;
found++;
}
}
if (vallen == 0) {
// first item
if (hashed_result == 0) {
value = urhc->items;
vallen = urhc->items_len;
}
// last item
else {
vallen = (urhc->items + urhc->items_len) - value;
}
}
if (!vallen) {
uwsgi_log("[uwsgi-hash-router] BUG !!! unable to hash items\n");
return UWSGI_ROUTE_BREAK;
}
if (!uwsgi_req_append(wsgi_req, urhc->var, urhc->var_len, value, vallen)) {
uwsgi_log("[uwsgi-hash-router] unable to append hash var to the request\n");
return UWSGI_ROUTE_BREAK;
}
return UWSGI_ROUTE_NEXT;
}
static int uwsgi_router_hash(struct uwsgi_route *ur, char *args) {
ur->func = uwsgi_routing_func_hash;
ur->data = args;
ur->data_len = strlen(args);
struct uwsgi_router_hash_conf *urhc = uwsgi_calloc(sizeof(struct uwsgi_router_hash_conf));
if (uwsgi_kvlist_parse(ur->data, ur->data_len, ',', '=',
"key", &urhc->key,
"var", &urhc->var,
"algo", &urhc->algo,
"items", &urhc->items,
NULL)) {
uwsgi_log("invalid route syntax: %s\n", args);
exit(1);
}
if (!urhc->key || !urhc->var || !urhc->items) {
uwsgi_log("invalid route syntax: you need to specify a hash key, a var and a set of items\n");
exit(1);
}
urhc->key_len = strlen(urhc->key);
urhc->var_len = strlen(urhc->var);
urhc->items_len = strlen(urhc->items);
if (!urhc->algo) urhc->algo = "djb33x";
ur->data2 = urhc;
return 0;
}
static void router_hash_register() {
uwsgi_register_router("hash", uwsgi_router_hash);
}
struct uwsgi_plugin router_hash_plugin = {
.name = "router_hash",
.on_load = router_hash_register,
};
#else
struct uwsgi_plugin router_hash_plugin = {
.name = "router_hash",
};
#endif
+6
View File
@@ -0,0 +1,6 @@
NAME='router_hash'
CFLAGS = []
LDFLAGS = []
LIBS = []
GCC_LIST = ['router_hash']
+174 -33
View File
@@ -1,20 +1,24 @@
#include <uwsgi.h>
#ifdef UWSGI_ROUTING
#define MEMCACHED_BUFSIZE 8192
extern struct uwsgi_server uwsgi;
/*
memcached internal router
memcached internal router and transformation
route = /^foobar1(.*)/ memcached:addr=127.0.0.1:11211,key=foo$1poo,type=body
route = /^foobar1(.*)/ memcached:addr=127.0.0.1:11211,key=foo$1poo
route = /^foobar1(.*)/ memcachedstore:addr=127.0.0.1:11211,key=foo$1poo
*/
struct uwsgi_router_memcached_conf {
char *addr;
size_t addr_len;
char *key;
size_t key_len;
@@ -22,12 +26,19 @@ struct uwsgi_router_memcached_conf {
char *content_type;
size_t content_type_len;
// 0 -> full, 1 -> body
char *type;
int type_num;
char *no_offload;
char *expires;
};
// this is allocated for each transformation
struct uwsgi_transformation_memcached_conf {
struct uwsgi_buffer *addr;
struct uwsgi_buffer *key;
char *expires;
};
static size_t memcached_firstline_parse(char *buf, size_t len) {
// check for "VALUE x 0 0"
if (len < 11) return 0;
@@ -48,6 +59,86 @@ static size_t memcached_firstline_parse(char *buf, size_t len) {
}
}
// store an item in memcached
static void memcached_store(char *addr, struct uwsgi_buffer *key, struct uwsgi_buffer *value, char *expires) {
int timeout = uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT];
int fd = uwsgi_connect(addr, 0, 1);
if (fd < 0) return;
// wait for connection
int ret = uwsgi.wait_write_hook(fd, timeout);
if (ret <= 0) goto end;
// build the request
struct uwsgi_buffer *ub = uwsgi_buffer_new(uwsgi.page_size);
if (uwsgi_buffer_append(ub, "set ", 4)) goto end2;
if (uwsgi_buffer_append(ub, key->buf, key->pos)) goto end2;
if (uwsgi_buffer_append(ub, " 0 " , 3)) goto end2;
if (uwsgi_buffer_append(ub, expires, strlen(expires))) goto end2;
if (uwsgi_buffer_append(ub, " " , 1)) goto end2;
if (uwsgi_buffer_num64(ub, value->pos)) goto end2;
if (uwsgi_buffer_append(ub, "\r\n" , 2)) goto end2;
if (uwsgi_write_true_nb(fd, ub->buf, ub->pos, timeout)) goto end2;
if (uwsgi_write_true_nb(fd, value->buf, value->pos, timeout)) goto end2;
if (uwsgi_write_true_nb(fd, "\r\n", 2, timeout)) goto end2;
// we are not interested in command result... (ugly but it works)
end2:
uwsgi_buffer_destroy(ub);
end:
close(fd);
}
static int transform_memcached(struct wsgi_request *wsgi_req, struct uwsgi_transformation *ut) {
struct uwsgi_transformation_memcached_conf *utmc = (struct uwsgi_transformation_memcached_conf *) ut->data;
struct uwsgi_buffer *ub = ut->chunk;
// store only successfull response
if (wsgi_req->write_errors == 0 && wsgi_req->status == 200 && ub->pos > 0) {
memcached_store(utmc->addr->buf, utmc->key, ub, utmc->expires);
}
// free resources
uwsgi_buffer_destroy(utmc->key);
uwsgi_buffer_destroy(utmc->addr);
free(utmc);
return 0;
}
// be tolerant on errors
static int uwsgi_routing_func_memcached_store(struct wsgi_request *wsgi_req, struct uwsgi_route *ur){
struct uwsgi_router_memcached_conf *urmc = (struct uwsgi_router_memcached_conf *) ur->data2;
struct uwsgi_transformation_memcached_conf *utmc = uwsgi_calloc(sizeof(struct uwsgi_transformation_memcached_conf));
// build key and name
char **subject = (char **) (((char *)(wsgi_req))+ur->subject);
uint16_t *subject_len = (uint16_t *) (((char *)(wsgi_req))+ur->subject_len);
utmc->key = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urmc->key, urmc->key_len);
if (!utmc->key) goto error;
utmc->addr = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urmc->addr, urmc->addr_len);
if (!utmc->addr) goto error;
utmc->expires = urmc->expires;
uwsgi_add_transformation(wsgi_req, transform_memcached, utmc);
return UWSGI_ROUTE_NEXT;
error:
if (utmc->key) uwsgi_buffer_destroy(utmc->key);
if (utmc->addr) uwsgi_buffer_destroy(utmc->addr);
free(utmc);
return UWSGI_ROUTE_NEXT;
}
static int uwsgi_routing_func_memcached(struct wsgi_request *wsgi_req, struct uwsgi_route *ur){
// this is the buffer for the memcached response
char buf[MEMCACHED_BUFSIZE];
@@ -62,13 +153,24 @@ static int uwsgi_routing_func_memcached(struct wsgi_request *wsgi_req, struct uw
struct uwsgi_buffer *ub_key = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urmc->key, urmc->key_len);
if (!ub_key) return UWSGI_ROUTE_BREAK;
int fd = uwsgi_connect(urmc->addr, 0, 1);
if (fd < 0) { uwsgi_buffer_destroy(ub_key) ; goto end; }
struct uwsgi_buffer *ub_addr = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urmc->addr, urmc->addr_len);
if (!ub_addr) {
uwsgi_buffer_destroy(ub_key);
return UWSGI_ROUTE_BREAK;
}
int fd = uwsgi_connect(ub_addr->buf, 0, 1);
if (fd < 0) {
uwsgi_buffer_destroy(ub_key);
uwsgi_buffer_destroy(ub_addr);
goto end;
}
// wait for connection;
int ret = uwsgi.wait_write_hook(fd, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]);
if (ret <= 0) {
uwsgi_buffer_destroy(ub_key) ;
uwsgi_buffer_destroy(ub_addr);
close(fd);
goto end;
}
@@ -77,11 +179,13 @@ static int uwsgi_routing_func_memcached(struct wsgi_request *wsgi_req, struct uw
char *cmd = uwsgi_concat3n("get ", 4, ub_key->buf, ub_key->pos, "\r\n", 2);
if (uwsgi_write_true_nb(fd, cmd, 6+ub_key->pos, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT])) {
uwsgi_buffer_destroy(ub_key);
uwsgi_buffer_destroy(ub_addr);
free(cmd);
close(fd);
goto end;
}
uwsgi_buffer_destroy(ub_key);
uwsgi_buffer_destroy(ub_addr);
free(cmd);
// ok, start reading the response...
@@ -130,26 +234,34 @@ read:
goto end;
}
if (urmc->type_num == 1) {
if (uwsgi_response_prepare_headers(wsgi_req, "200 OK", 6)) { close(fd); goto end; }
if (uwsgi_response_add_content_type(wsgi_req, urmc->content_type, urmc->content_type_len)) { close(fd); goto end; }
if (uwsgi_response_add_content_length(wsgi_req, response_size)) { close(fd); goto end; }
}
// from now on, every error will trigger a BREAK...
// send headers
if (uwsgi_response_prepare_headers(wsgi_req, "200 OK", 6)) goto error;
if (uwsgi_response_add_content_type(wsgi_req, urmc->content_type, urmc->content_type_len)) goto error;
if (uwsgi_response_add_content_length(wsgi_req, response_size)) goto error;
// the first chunk could already contains part of the body
size_t remains = pos-(found+2);
if (remains >= response_size) {
uwsgi_response_write_body_do(wsgi_req, buf+found+2, response_size);
close(fd);
goto end;
goto done;
}
// send what we have
if (uwsgi_response_write_body_do(wsgi_req, buf+found+2, remains)) {
close(fd);
goto end;
}
if (uwsgi_response_write_body_do(wsgi_req, buf+found+2, remains)) goto error;
// and now start reading til the output is consumed
response_size -= remains;
// try to offload via the pipe engine
if (wsgi_req->socket->can_offload && !ur->custom && !urmc->no_offload) {
if (!uwsgi_offload_request_pipe_do(wsgi_req, fd, response_size)) {
wsgi_req->via = UWSGI_VIA_OFFLOAD;
return UWSGI_ROUTE_BREAK;
}
}
while(response_size > 0) {
ssize_t len = read(fd, buf, UMIN(MEMCACHED_BUFSIZE, response_size));
if (len > 0) goto write;
@@ -165,23 +277,22 @@ wait2:
}
goto error;
write:
if (uwsgi_response_write_body_do(wsgi_req, buf, len)) {
goto error;
}
if (uwsgi_response_write_body_do(wsgi_req, buf, len)) goto error;
response_size -= len;
}
done:
close(fd);
if (ur->custom)
return UWSGI_ROUTE_NEXT;
return UWSGI_ROUTE_BREAK;
error:
close(fd);
return UWSGI_ROUTE_BREAK;
end:
if (ur->custom)
return UWSGI_ROUTE_NEXT;
return UWSGI_ROUTE_BREAK;
return UWSGI_ROUTE_NEXT;
}
static int uwsgi_router_memcached(struct uwsgi_route *ur, char *args) {
@@ -193,27 +304,23 @@ static int uwsgi_router_memcached(struct uwsgi_route *ur, char *args) {
"addr", &urmc->addr,
"key", &urmc->key,
"content_type", &urmc->content_type,
"type", &urmc->type, NULL)) {
"no_offload", &urmc->no_offload,
NULL)) {
uwsgi_log("invalid route syntax: %s\n", args);
exit(1);
}
if (!urmc->key || !urmc->addr) {
uwsgi_log("invalid route syntax: you need to specify a memcached address and key pattern\n");
exit(1);
return -1;
}
urmc->key_len = strlen(urmc->key);
urmc->addr_len = strlen(urmc->addr);
if (!urmc->type) urmc->type = "full";
if (!urmc->content_type) urmc->content_type = "text/html";
urmc->content_type_len = strlen(urmc->content_type);
if (!strcmp(urmc->type, "body")) {
urmc->type_num = 1;
}
ur->data2 = urmc;
return 0;
}
@@ -224,12 +331,46 @@ static int uwsgi_router_memcached_continue(struct uwsgi_route *ur, char *args) {
return 0;
}
static int uwsgi_router_memcached_store(struct uwsgi_route *ur, char *args) {
ur->func = uwsgi_routing_func_memcached_store;
ur->data = args;
ur->data_len = strlen(args);
struct uwsgi_router_memcached_conf *urmc = uwsgi_calloc(sizeof(struct uwsgi_router_memcached_conf));
if (uwsgi_kvlist_parse(ur->data, ur->data_len, ',', '=',
"addr", &urmc->addr,
"key", &urmc->key,
"expires", &urmc->expires, NULL)) {
uwsgi_log("invalid memcachedstore route syntax: %s\n", args);
return -1;
}
if (!urmc->key || !urmc->addr) {
uwsgi_log("invalid memcachedstore route syntax: you need to specify an address and a key\n");
return -1;
}
urmc->key_len = strlen(urmc->key);
urmc->addr_len = strlen(urmc->addr);
if (!urmc->expires) urmc->expires = "0";
ur->data2 = urmc;
return 0;
}
static void router_memcached_register() {
uwsgi_register_router("memcached", uwsgi_router_memcached);
uwsgi_register_router("memcached-continue", uwsgi_router_memcached_continue);
uwsgi_register_router("memcachedstore", uwsgi_router_memcached_store);
uwsgi_register_router("memcached-store", uwsgi_router_memcached_store);
}
#endif
struct uwsgi_plugin router_memcached_plugin = {
.name = "router_memcached",
#ifdef UWSGI_ROUTING
.on_load = router_memcached_register,
#endif
};
+374
View File
@@ -0,0 +1,374 @@
#include <uwsgi.h>
#ifdef UWSGI_ROUTING
#define REDIS_BUFSIZE 8192
extern struct uwsgi_server uwsgi;
/*
redis internal router and transformation
route = /^foobar1(.*)/ redis:addr=127.0.0.1:11211,key=foo$1poo
route = /^foobar1(.*)/ redisstore:addr=127.0.0.1:11211,key=foo$1poo
*/
struct uwsgi_router_redis_conf {
char *addr;
size_t addr_len;
char *key;
size_t key_len;
char *content_type;
size_t content_type_len;
char *no_offload;
char *expires;
};
// this is allocated for each transformation
struct uwsgi_transformation_redis_conf {
struct uwsgi_buffer *addr;
struct uwsgi_buffer *key;
char *expires;
};
static size_t redis_firstline_parse(char *buf, size_t len) {
// check for "$0"
if (len < 2) return 0;
if (buf[0] != '$') return 0;
if (buf[1] == '-') return 0;
return uwsgi_str_num(buf + 1, len - 1);
}
// store an item in redis
static void redis_store(char *addr, struct uwsgi_buffer *key, struct uwsgi_buffer *value, char *expires) {
int timeout = uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT];
int fd = uwsgi_connect(addr, 0, 1);
if (fd < 0) return;
// wait for connection
int ret = uwsgi.wait_write_hook(fd, timeout);
if (ret <= 0) goto end;
// build the request
struct uwsgi_buffer *ub = uwsgi_buffer_new(uwsgi.page_size);
if (uwsgi_buffer_append(ub, "*3\r\n$3\r\nSET\r\n$", 14)) goto end2;
if (uwsgi_buffer_num64(ub, key->pos)) goto end2;
if (uwsgi_buffer_append(ub, "\r\n" , 2)) goto end2;
if (uwsgi_buffer_append(ub, key->buf, key->pos)) goto end2;
if (uwsgi_buffer_append(ub, "\r\n$" , 3)) goto end2;
if (uwsgi_buffer_num64(ub, value->pos)) goto end2;
if (uwsgi_buffer_append(ub, "\r\n" , 2)) goto end2;
if (uwsgi_write_true_nb(fd, ub->buf, ub->pos, timeout)) goto end2;
if (uwsgi_write_true_nb(fd, value->buf, value->pos, timeout)) goto end2;
ub->pos = 0;
if (strcmp(expires, "0")) {
if (uwsgi_buffer_append(ub, "\r\n*3\r\n$6\r\nEXPIRE\r\n$" , 19)) goto end2;
if (uwsgi_buffer_num64(ub, key->pos)) goto end2;
if (uwsgi_buffer_append(ub, "\r\n" , 2)) goto end2;
if (uwsgi_buffer_append(ub, key->buf, key->pos)) goto end2;
if (uwsgi_buffer_append(ub, "\r\n$" , 3)) goto end2;
if (uwsgi_buffer_num64(ub, strlen(expires))) goto end2;
if (uwsgi_buffer_append(ub, "\r\n" , 2)) goto end2;
if (uwsgi_buffer_append(ub, expires, strlen(expires))) goto end2;
}
if (uwsgi_buffer_append(ub, "\r\n" , 2)) goto end2;
if (uwsgi_write_true_nb(fd, ub->buf, ub->pos, timeout)) goto end2;
// we are not interested in command result... (ugly but it works)
end2:
uwsgi_buffer_destroy(ub);
end:
close(fd);
}
static int transform_redis(struct wsgi_request *wsgi_req, struct uwsgi_transformation *ut) {
struct uwsgi_transformation_redis_conf *utrc = (struct uwsgi_transformation_redis_conf *) ut->data;
struct uwsgi_buffer *ub = ut->chunk;
// store only successfull response
if (wsgi_req->write_errors == 0 && wsgi_req->status == 200 && ub->pos > 0) {
redis_store(utrc->addr->buf, utrc->key, ub, utrc->expires);
}
// free resources
uwsgi_buffer_destroy(utrc->key);
uwsgi_buffer_destroy(utrc->addr);
free(utrc);
return 0;
}
// be tolerant on errors
static int uwsgi_routing_func_redis_store(struct wsgi_request *wsgi_req, struct uwsgi_route *ur){
struct uwsgi_router_redis_conf *urrc = (struct uwsgi_router_redis_conf *) ur->data2;
struct uwsgi_transformation_redis_conf *utrc = uwsgi_calloc(sizeof(struct uwsgi_transformation_redis_conf));
// build key and name
char **subject = (char **) (((char *)(wsgi_req))+ur->subject);
uint16_t *subject_len = (uint16_t *) (((char *)(wsgi_req))+ur->subject_len);
utrc->key = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urrc->key, urrc->key_len);
if (!utrc->key) goto error;
utrc->addr = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urrc->addr, urrc->addr_len);
if (!utrc->addr) goto error;
utrc->expires = urrc->expires;
uwsgi_add_transformation(wsgi_req, transform_redis, utrc);
return UWSGI_ROUTE_NEXT;
error:
if (utrc->key) uwsgi_buffer_destroy(utrc->key);
if (utrc->addr) uwsgi_buffer_destroy(utrc->addr);
free(utrc);
return UWSGI_ROUTE_NEXT;
}
static int uwsgi_routing_func_redis(struct wsgi_request *wsgi_req, struct uwsgi_route *ur){
// this is the buffer for the redis response
char buf[REDIS_BUFSIZE];
size_t i;
char last_char = 0;
struct uwsgi_router_redis_conf *urrc = (struct uwsgi_router_redis_conf *) ur->data2;
char **subject = (char **) (((char *)(wsgi_req))+ur->subject);
uint16_t *subject_len = (uint16_t *) (((char *)(wsgi_req))+ur->subject_len);
struct uwsgi_buffer *ub_key = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urrc->key, urrc->key_len);
if (!ub_key) return UWSGI_ROUTE_BREAK;
struct uwsgi_buffer *ub_addr = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urrc->addr, urrc->addr_len);
if (!ub_addr) {
uwsgi_buffer_destroy(ub_key);
return UWSGI_ROUTE_BREAK;
}
int fd = uwsgi_connect(ub_addr->buf, 0, 1);
if (fd < 0) {
uwsgi_buffer_destroy(ub_key);
uwsgi_buffer_destroy(ub_addr);
goto end;
}
// wait for connection;
int ret = uwsgi.wait_write_hook(fd, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]);
if (ret <= 0) {
uwsgi_buffer_destroy(ub_key) ;
uwsgi_buffer_destroy(ub_addr);
close(fd);
goto end;
}
// build the request and send it
char *cmd = uwsgi_concat3n("get ", 4, ub_key->buf, ub_key->pos, "\r\n", 2);
if (uwsgi_write_true_nb(fd, cmd, 6+ub_key->pos, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT])) {
uwsgi_buffer_destroy(ub_key);
uwsgi_buffer_destroy(ub_addr);
free(cmd);
close(fd);
goto end;
}
uwsgi_buffer_destroy(ub_key);
uwsgi_buffer_destroy(ub_addr);
free(cmd);
// ok, start reading the response...
// first we need to get a full line;
size_t found = 0;
size_t pos = 0;
for(;;) {
ssize_t len = read(fd, buf + pos, REDIS_BUFSIZE - pos);
if (len > 0) {
pos += len;
goto read;
}
if (len < 0) {
if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) goto wait;
}
close(fd);
goto end;
wait:
ret = uwsgi.wait_read_hook(fd, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]);
// when we have a chunk try to read the first line
if (ret > 0) {
len = read(fd, buf + pos, REDIS_BUFSIZE - pos);
if (len > 0) {
pos += len;
goto read;
}
}
close(fd);
goto end;
read:
for(i=0;i<pos;i++) {
if (last_char == '\r' && buf[i] == '\n') {
found = i-1;
break;
}
last_char = buf[i];
}
if (found) break;
}
// ok parse the first line
size_t response_size = redis_firstline_parse(buf, found);
if (response_size == 0) {
close(fd);
goto end;
}
// from now on, every error will trigger a BREAK...
// send headers
if (uwsgi_response_prepare_headers(wsgi_req, "200 OK", 6)) goto error;
if (uwsgi_response_add_content_type(wsgi_req, urrc->content_type, urrc->content_type_len)) goto error;
if (uwsgi_response_add_content_length(wsgi_req, response_size)) goto error;
// the first chunk could already contains part of the body
size_t remains = pos-(found+2);
if (remains >= response_size) {
uwsgi_response_write_body_do(wsgi_req, buf+found+2, response_size);
goto done;
}
// send what we have
if (uwsgi_response_write_body_do(wsgi_req, buf+found+2, remains)) goto error;
// and now start reading til the output is consumed
response_size -= remains;
// try to offload via the pipe engine
if (wsgi_req->socket->can_offload && !ur->custom && !urrc->no_offload) {
if (!uwsgi_offload_request_pipe_do(wsgi_req, fd, response_size)) {
wsgi_req->via = UWSGI_VIA_OFFLOAD;
return UWSGI_ROUTE_BREAK;
}
}
while(response_size > 0) {
ssize_t len = read(fd, buf, UMIN(REDIS_BUFSIZE, response_size));
if (len > 0) goto write;
if (len < 0) {
if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) goto wait2;
}
goto error;
wait2:
ret = uwsgi.wait_read_hook(fd, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]);
if (ret > 0) {
len = read(fd, buf, UMIN(REDIS_BUFSIZE, response_size));
if (len > 0) goto write;
}
goto error;
write:
if (uwsgi_response_write_body_do(wsgi_req, buf, len)) goto error;
response_size -= len;
}
done:
close(fd);
if (ur->custom)
return UWSGI_ROUTE_NEXT;
return UWSGI_ROUTE_BREAK;
error:
close(fd);
return UWSGI_ROUTE_BREAK;
end:
return UWSGI_ROUTE_NEXT;
}
static int uwsgi_router_redis(struct uwsgi_route *ur, char *args) {
ur->func = uwsgi_routing_func_redis;
ur->data = args;
ur->data_len = strlen(args);
struct uwsgi_router_redis_conf *urrc = uwsgi_calloc(sizeof(struct uwsgi_router_redis_conf));
if (uwsgi_kvlist_parse(ur->data, ur->data_len, ',', '=',
"addr", &urrc->addr,
"key", &urrc->key,
"content_type", &urrc->content_type,
"no_offload", &urrc->no_offload,
NULL)) {
uwsgi_log("invalid route syntax: %s\n", args);
exit(1);
}
if (!urrc->key || !urrc->addr) {
uwsgi_log("invalid route syntax: you need to specify a redis address and key pattern\n");
return -1;
}
urrc->key_len = strlen(urrc->key);
urrc->addr_len = strlen(urrc->addr);
if (!urrc->content_type) urrc->content_type = "text/html";
urrc->content_type_len = strlen(urrc->content_type);
ur->data2 = urrc;
return 0;
}
static int uwsgi_router_redis_continue(struct uwsgi_route *ur, char *args) {
uwsgi_router_redis(ur, args);
ur->custom = 1;
return 0;
}
static int uwsgi_router_redis_store(struct uwsgi_route *ur, char *args) {
ur->func = uwsgi_routing_func_redis_store;
ur->data = args;
ur->data_len = strlen(args);
struct uwsgi_router_redis_conf *urrc = uwsgi_calloc(sizeof(struct uwsgi_router_redis_conf));
if (uwsgi_kvlist_parse(ur->data, ur->data_len, ',', '=',
"addr", &urrc->addr,
"key", &urrc->key,
"expires", &urrc->expires, NULL)) {
uwsgi_log("invalid redisstore route syntax: %s\n", args);
return -1;
}
if (!urrc->key || !urrc->addr) {
uwsgi_log("invalid redisstore route syntax: you need to specify an address and a key\n");
return -1;
}
urrc->key_len = strlen(urrc->key);
urrc->addr_len = strlen(urrc->addr);
if (!urrc->expires) urrc->expires = "0";
ur->data2 = urrc;
return 0;
}
static void router_redis_register() {
uwsgi_register_router("redis", uwsgi_router_redis);
uwsgi_register_router("redis-continue", uwsgi_router_redis_continue);
uwsgi_register_router("redisstore", uwsgi_router_redis_store);
uwsgi_register_router("redis-store", uwsgi_router_redis_store);
}
#endif
struct uwsgi_plugin router_redis_plugin = {
.name = "router_redis",
#ifdef UWSGI_ROUTING
.on_load = router_redis_register,
#endif
};
+6
View File
@@ -0,0 +1,6 @@
NAME='router_redis'
CFLAGS = []
LDFLAGS = []
LIBS = []
GCC_LIST = ['router_redis']
+138
View File
@@ -0,0 +1,138 @@
#include <uwsgi.h>
#if defined(UWSGI_ROUTING)
/*
offload transformation
each chunk is buffered to the transformation custom buffer
if the buffer grows higher than the specified limit, the response will be buffered to disk
on final the memory (or the tmpfile) will be offloaded
even if we are talking about buffering, this is a streaming transformation
as we need to hold control on memory usage
"offload" MUST BE the last transformation in the chain. Hellish things will happen otherwise...
*/
static int transform_offload(struct wsgi_request *wsgi_req, struct uwsgi_transformation *ut) {
if (ut->is_final) {
struct uwsgi_transformation *orig_ut = (struct uwsgi_transformation *) ut->data;
// sendfile offload
if (orig_ut->fd > -1) {
if (!uwsgi_offload_request_sendfile_do(wsgi_req, orig_ut->fd, orig_ut->len)) {
// the fd will be closed by the offload engine
orig_ut->fd = -1;
wsgi_req->via = UWSGI_VIA_OFFLOAD;
wsgi_req->response_size += orig_ut->len;
return 0;
}
// fallback to non-offloaded write
if (uwsgi_simple_sendfile(wsgi_req, orig_ut->fd, 0, orig_ut->len)) {
return -1;
}
wsgi_req->response_size += orig_ut->len;
return 0;
}
// memory offload
if (orig_ut->ub) {
if (!uwsgi_offload_request_memory_do(wsgi_req, orig_ut->ub->buf, orig_ut->ub->pos)) {
// memory will be freed by the offload engine
orig_ut->ub->buf = NULL;
wsgi_req->via = UWSGI_VIA_OFFLOAD;
wsgi_req->response_size += orig_ut->ub->pos;
return 0;
}
// fallback to non-offloaded write
if (uwsgi_simple_write(wsgi_req, orig_ut->ub->buf, orig_ut->ub->pos)) {
return -1;
}
wsgi_req->response_size += orig_ut->ub->pos;
return -1;
}
return 0;
}
// check if we need to start buffering to disk
if (ut->fd == -1 && ut->len + ut->chunk->pos > ut->custom64) {
ut->fd = uwsgi_tmpfd();
if (ut->fd < 0) return -1;
// save already buffered data
if (ut->ub) {
ssize_t wlen = write(ut->fd, ut->ub->buf, ut->ub->pos);
if (wlen != (ssize_t) ut->ub->pos) {
uwsgi_error("transform_offload/write()");
return -1;
}
}
}
// if fd > -1, append to file
if (ut->fd > -1) {
ssize_t wlen = write(ut->fd, ut->chunk->buf, ut->chunk->pos);
if (wlen != (ssize_t) ut->chunk->pos) {
uwsgi_error("transform_offload/write()");
return -1;
}
ut->len += wlen;
goto done;
}
// buffer to memory
if (!ut->ub) {
ut->ub = uwsgi_buffer_new(ut->chunk->pos);
}
// append the chunk to the custom buffer
if (uwsgi_buffer_append(ut->ub, ut->chunk->buf, ut->chunk->pos)) return -1;
ut->len += ut->chunk->pos;
done:
// reset the chunk !!!
ut->chunk->pos = 0;
return 0;
}
static int uwsgi_routing_func_offload(struct wsgi_request *wsgi_req, struct uwsgi_route *ur) {
if (!wsgi_req->socket->can_offload) {
uwsgi_log("unable to use the offload transformation without offload threads !!!\n");
return UWSGI_ROUTE_BREAK;
}
struct uwsgi_transformation *ut = uwsgi_add_transformation(wsgi_req, transform_offload, NULL);
ut->can_stream = 1;
ut->custom64 = ur->custom;
// add a "final" transformation to add the trailing chunk
ut = uwsgi_add_transformation(wsgi_req, transform_offload, ut);
ut->is_final = 1;
return UWSGI_ROUTE_NEXT;
}
static int uwsgi_router_offload(struct uwsgi_route *ur, char *args) {
ur->func = uwsgi_routing_func_offload;
if (args[0] == 0) {
// 1 MB limit
ur->custom = 1024 * 1024;
}
else {
ur->custom = strtoul(args, NULL, 10);
}
return 0;
}
static void router_offload_register(void) {
uwsgi_register_router("offload", uwsgi_router_offload);
}
struct uwsgi_plugin transformation_offload_plugin = {
.name = "transformation_offload",
.on_load = router_offload_register,
};
#else
struct uwsgi_plugin transformation_offload_plugin = {
.name = "transformation_offload",
};
#endif
@@ -0,0 +1,6 @@
NAME='transformation_offload'
CFLAGS = []
LDFLAGS = []
LIBS = []
GCC_LIST = ['offload']
+13 -1
View File
@@ -155,6 +155,11 @@ ssize_t uwsgi_proto_fastcgi_read_body(struct wsgi_request *wsgi_req, char *buf,
memcpy(buf, wsgi_req->proto_parser_remains_buf, remains);
wsgi_req->proto_parser_remains -= remains;
wsgi_req->proto_parser_remains_buf += remains;
// we consumed all of the body, we can safely move the memory
if (wsgi_req->proto_parser_remains == 0 && wsgi_req->proto_parser_move) {
memmove(wsgi_req->proto_parser_buf, wsgi_req->proto_parser_buf + wsgi_req->proto_parser_move, wsgi_req->proto_parser_pos);
wsgi_req->proto_parser_move = 0;
}
return remains;
}
@@ -176,7 +181,14 @@ ssize_t uwsgi_proto_fastcgi_read_body(struct wsgi_request *wsgi_req, char *buf,
// copy remaining
wsgi_req->proto_parser_remains = fcgi_len - remains;
wsgi_req->proto_parser_remains_buf = wsgi_req->proto_parser_buf + sizeof(struct fcgi_record) + remains;
memmove(wsgi_req->proto_parser_buf, wsgi_req->proto_parser_buf + fcgi_all_len, wsgi_req->proto_parser_pos - fcgi_all_len);
// we consumed all of the body, we can safely move the memory
if (wsgi_req->proto_parser_remains == 0) {
memmove(wsgi_req->proto_parser_buf, wsgi_req->proto_parser_buf + fcgi_all_len, wsgi_req->proto_parser_pos - fcgi_all_len);
}
else {
// postpone memory move
wsgi_req->proto_parser_move = fcgi_all_len;
}
wsgi_req->proto_parser_pos -= fcgi_all_len;
return remains;
}
+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 = '2013-05-26'
s.date = '2013-06-05'
s.summary = "uWSGI"
s.description = "The uWSGI server for Ruby/Rack"
s.authors = ["Unbit"]
+16 -1
View File
@@ -278,7 +278,9 @@ extern int pivot_root(const char *new_root, const char *put_old);
#endif
#if defined(__HAIKU__) || defined(__CYGWIN__)
#ifndef WAIT_ANY
#define WAIT_ANY (-1)
#endif
#define PRIO_MAX 20
#endif
@@ -1220,6 +1222,10 @@ struct uwsgi_transformation {
uint8_t flushed;
void *data;
uint64_t round;
int fd;
struct uwsgi_buffer *ub;
uint64_t len;
uint64_t custom64;
struct uwsgi_transformation *next;
};
@@ -1398,6 +1404,7 @@ struct wsgi_request {
int headers_hvec;
uint64_t proto_parser_pos;
uint64_t proto_parser_move;
int64_t proto_parser_status;
void *proto_parser_buf;
uint64_t proto_parser_buf_size;
@@ -1945,6 +1952,7 @@ struct uwsgi_server {
struct uwsgi_offload_engine *offload_engine_sendfile;
struct uwsgi_offload_engine *offload_engine_transfer;
struct uwsgi_offload_engine *offload_engine_memory;
struct uwsgi_offload_engine *offload_engine_pipe;
int offload_threads;
int offload_threads_events;
struct uwsgi_thread **offload_thread;
@@ -1995,6 +2003,7 @@ struct uwsgi_server {
struct uwsgi_string_list *touch_logrotate;
struct uwsgi_string_list *touch_logreopen;
struct uwsgi_string_list *touch_exec;
struct uwsgi_string_list *touch_signal;
int propagate_touch;
@@ -3227,6 +3236,7 @@ struct uwsgi_gateway_socket *uwsgi_new_gateway_socket(char *, char *);
struct uwsgi_gateway_socket *uwsgi_new_gateway_socket_from_fd(int, char *);
void escape_shell_arg(char *, size_t, char *);
void escape_json(char *, size_t, char *);
void *uwsgi_malloc_shared(size_t);
void *uwsgi_calloc_shared(size_t);
@@ -3355,7 +3365,8 @@ void uwsgi_opt_set_cap(char *, char *, void *);
void uwsgi_opt_set_unshare(char *, char *, void *);
#endif
char *uwsgi_tmpname(char *, char *);
int uwsgi_tmpfd();
FILE *uwsgi_tmpfile();
#ifdef UWSGI_ROUTING
struct uwsgi_router *uwsgi_register_router(char *, int (*)(struct uwsgi_route *, char *));
@@ -3770,6 +3781,10 @@ struct uwsgi_thread *uwsgi_offload_thread_start(void);
int uwsgi_offload_request_sendfile_do(struct wsgi_request *, int, size_t);
int uwsgi_offload_request_net_do(struct wsgi_request *, char *, struct uwsgi_buffer *);
int uwsgi_offload_request_memory_do(struct wsgi_request *, char *, size_t);
int uwsgi_offload_request_pipe_do(struct wsgi_request *, int, size_t);
int uwsgi_simple_sendfile(struct wsgi_request *, int, size_t, size_t);
int uwsgi_simple_write(struct wsgi_request *, char *, size_t);
void uwsgi_subscription_set_algo(char *);
+1 -1
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '1.9.11'
uwsgi_version = '1.9.12'
import os
import re