extended the uwsgi api to allow easy implementation of websockets servers

This commit is contained in:
roberto@mrspurr
2010-12-06 13:48:23 +01:00
parent 68c6dd13b8
commit 9aaf36141e
4 changed files with 200 additions and 20 deletions
+94 -19
View File
@@ -176,10 +176,13 @@ static void *http_request(void *u_h_r) {
int state = uwsgi_http_method;
int http_body_len = 0;
int http_upgrade = 0;
struct pollfd http_poll[2];
size_t len;
int i, j;
int i, j, rlen;
char HTTP_header_key[1024];
@@ -279,6 +282,13 @@ static void *http_request(void *u_h_r) {
*ptr++ = 0;
http_body_len = atoi(tmp_buf);
}
else if (!strcmp("CONNECTION", HTTP_header_key)) {
if (ptr+1 > watermark2) { close(uwsgi_fd); goto clear;}
*ptr++ = 0;
if (!strcmp(tmp_buf, "Upgrade")) {
http_upgrade = 1;
}
}
ptr = tmp_buf;
state = uwsgi_http_header_key;
} else if (state == uwsgi_http_protocol_r) {
@@ -322,31 +332,96 @@ static void *http_request(void *u_h_r) {
uwsgi_error("write()");
}
if (http_body_len > 0) {
if (http_body_len >= (int) len - (i + 1)) {
if (http_upgrade) {
// send already available data
if ( (len - (i + 1)) > 0) {
if (write(uwsgi_fd, buf + i + 1, len - (i + 1)) < 0) {
uwsgi_error("write()");
}
http_body_len -= len - (i + 1);
} else {
if (write(uwsgi_fd, buf + i, http_body_len) < 0) {
uwsgi_error("write()");
}
http_body_len = 0;
uwsgi_error("write()");
close(uwsgi_fd);
goto clear;
}
}
while (http_body_len > 0) {
int to_read = 4096;
if (http_body_len < to_read) {
to_read = http_body_len;
http_poll[0].fd = clientfd;
http_poll[0].events = POLLIN;
http_poll[1].fd = uwsgi_fd;
http_poll[1].events = POLLIN;
for(;;) {
rlen = poll(http_poll, 2, -1);
if (rlen < 0) {
uwsgi_error("poll()");
close(uwsgi_fd);
goto clear;
}
len = read(clientfd, uwsgipkt, to_read);
if (write(uwsgi_fd, uwsgipkt, len) < 0) {
uwsgi_error("write()");
else if (rlen > 0) {
if (http_poll[0].revents & POLLIN) {
len = read(clientfd, uwsgipkt, 4096);
if (len > 0) {
if (write(uwsgi_fd, uwsgipkt, len) < 0) {
uwsgi_error("write()");
close(uwsgi_fd);
goto clear;
}
}
else {
// client disconnected
close(uwsgi_fd);
goto clear;
}
}
else if (http_poll[1].revents & POLLIN) {
len = read(uwsgi_fd, uwsgipkt, 4096);
if (len > 0) {
if (write(clientfd, uwsgipkt, len) < 0) {
uwsgi_error("write()");
close(uwsgi_fd);
goto clear;
}
}
else {
// client disconnected
close(uwsgi_fd);
goto clear;
}
}
}
else {
// timeout
close(uwsgi_fd);
goto clear;
}
}
}
else {
if (http_body_len > 0) {
if (http_body_len >= (int) len - (i + 1)) {
if (write(uwsgi_fd, buf + i + 1, len - (i + 1)) < 0) {
uwsgi_error("write()");
}
http_body_len -= len - (i + 1);
} else {
if (write(uwsgi_fd, buf + i, http_body_len) < 0) {
uwsgi_error("write()");
}
http_body_len = 0;
}
while (http_body_len > 0) {
int to_read = 4096;
if (http_body_len < to_read) {
to_read = http_body_len;
}
len = read(clientfd, uwsgipkt, to_read);
if (write(uwsgi_fd, uwsgipkt, len) < 0) {
uwsgi_error("write()");
}
http_body_len -= len;
}
http_body_len -= len;
}
}
while ((len = read(uwsgi_fd, uwsgipkt, 4096)) > 0) {
if (write(clientfd, uwsgipkt, len) < 0) {
uwsgi_error("write()");
+54
View File
@@ -166,6 +166,59 @@ PyObject *py_uwsgi_close(PyObject * self, PyObject * args) {
}
PyObject *py_uwsgi_recv_block(PyObject * self, PyObject * args) {
char buf[4096];
char *bufptr;
ssize_t rlen = 0, len ;
int fd, size, remains, ret, timeout = -1;
if (!PyArg_ParseTuple(args, "ii|i:recv_block", &fd, &size, &timeout)) {
return NULL;
}
if (fd < 0) goto clear;
UWSGI_RELEASE_GIL
// security check
if (size > 4096) size = 4096;
remains = size;
bufptr = buf;
while(remains > 0) {
uwsgi_log("%d %d %d\n", remains, size, timeout);
ret = uwsgi_waitfd(fd, timeout);
if (ret > 0) {
len = read(fd, bufptr, UMIN(remains, size)) ;
if (len > 0) {
bufptr+=len;
rlen += len;
remains -= len;
}
else {
break;
}
}
else {
uwsgi_log("error waiting for block data\n");
break;
}
}
UWSGI_GET_GIL
if ( rlen == size) {
return PyString_FromStringAndSize(buf, rlen);
}
clear:
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_recv(PyObject * self, PyObject * args) {
int fd, max_size = 4096;
@@ -1639,6 +1692,7 @@ static PyMethodDef uwsgi_advanced_methods[] = {
{"is_connected", py_uwsgi_is_connected, METH_VARARGS, ""},
{"send", py_uwsgi_send, METH_VARARGS, ""},
{"recv", py_uwsgi_recv, METH_VARARGS, ""},
{"recv_block", py_uwsgi_recv_block, METH_VARARGS, ""},
{"close", py_uwsgi_close, METH_VARARGS, ""},
{"parsefile", py_uwsgi_parse_file, METH_VARARGS, ""},
+4 -1
View File
@@ -1141,7 +1141,10 @@ int uwsgi_waitfd(int fd, int timeout) {
if (!timeout) timeout = uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT];
ret = poll(upoll, 1, timeout*1000);
timeout = timeout*1000;
if (timeout < 0) timeout = -1;
ret = poll(upoll, 1, timeout);
if (ret < 0) {
uwsgi_error("poll()");
+48
View File
@@ -0,0 +1,48 @@
import uwsgi
def application(e, s):
print e
client = e['wsgi.input'].fileno()
print client
data = uwsgi.recv_block(client, 8)
print "data", data, len(data)
key1 = e['HTTP_SEC_WEBSOCKET_KEY1']
key2 = e['HTTP_SEC_WEBSOCKET_KEY2']
total1 = ''
div1 = 0
for c in key1:
if c in '0'..'9':
total1 += c
for c in key1:
if c == ' ':
div1 += 1
if div1 == 0:
raise StopIteration
total1 = int(total1) / div1
total2 = ''
div2 = 0
for c in key2:
if c in '0'..'9':
total2 += c
for c in key2:
if c == ' ':
div2 += 1
if div2 == 0:
raise StopIteration
total2 = int(total2) / div1