mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-10-04 16:08:31 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0f397eeb55 | ||
|
|
bf929398a9 | ||
|
|
4f8a0f753e | ||
|
|
7db5da9566 | ||
|
|
4d1389805c | ||
|
|
ee87e2b946 | ||
|
|
6356d78768 | ||
|
|
afdf7e740e | ||
|
|
51ce3ddf78 | ||
|
|
192186056e | ||
|
|
336a7febd7 | ||
|
|
63e73d1b73 | ||
|
|
a02ab9a723 | ||
|
|
95f15ab9dd | ||
|
|
5d0ea0066a | ||
|
|
9ff6b2da9b | ||
|
|
6ac7e27847 | ||
|
|
4ed1c5e747 | ||
|
|
bae3bc096f | ||
|
|
9d73683de2 | ||
|
|
f5fec0ab71 | ||
|
|
9909ab10c4 | ||
|
|
65cff01fee | ||
|
|
10190a3012 | ||
|
|
79ab696f4e |
+1
-1
@@ -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
|
||||
|
||||
@@ -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
@@ -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
@@ -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
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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
@@ -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) {
|
||||
|
||||
@@ -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},
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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, ...);
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
@@ -0,0 +1,6 @@
|
||||
NAME='router_hash'
|
||||
|
||||
CFLAGS = []
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
GCC_LIST = ['router_hash']
|
||||
@@ -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
|
||||
};
|
||||
|
||||
@@ -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
|
||||
};
|
||||
@@ -0,0 +1,6 @@
|
||||
NAME='router_redis'
|
||||
|
||||
CFLAGS = []
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
GCC_LIST = ['router_redis']
|
||||
@@ -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
@@ -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
@@ -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"]
|
||||
|
||||
@@ -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
@@ -1,6 +1,6 @@
|
||||
# uWSGI build system
|
||||
|
||||
uwsgi_version = '1.9.11'
|
||||
uwsgi_version = '1.9.12'
|
||||
|
||||
import os
|
||||
import re
|
||||
|
||||
Reference in New Issue
Block a user