finally a worker (almost) solid spooler. Ready for beta1

This commit is contained in:
roberto@voldemort
2009-12-19 20:29:55 +01:00
parent ec077bc843
commit 18efb257eb
5 changed files with 210 additions and 31 deletions
+68 -21
View File
@@ -3,12 +3,21 @@
#include <dirent.h>
extern char *spool_dir;
int spool_request(char *spooldir, char *filename, int rn, char *buffer, int size) {
struct uwsgi_packet_header {
uint8_t modifier1;
uint16_t datasize;
uint8_t modifier2;
};
int spool_request(char *filename, int rn, char *buffer, int size) {
char hostname[256+1];
struct timeval tv;
int fd;
struct uwsgi_packet_header uh ;
if (gethostname(hostname,256)) {
perror("gethostname()");
@@ -19,7 +28,7 @@ int spool_request(char *spooldir, char *filename, int rn, char *buffer, int size
hostname[256] = 0 ;
if (snprintf(filename,1024,"%s/uwsgi_spoolfile_on_%s_%d_%d_%llu_%llu", spooldir, hostname, getpid(), rn, (unsigned long long) tv.tv_sec, (unsigned long long) tv.tv_usec ) <= 0) {
if (snprintf(filename,1024,"%s/uwsgi_spoolfile_on_%s_%d_%d_%llu_%llu", spool_dir, hostname, getpid(), rn, (unsigned long long) tv.tv_sec, (unsigned long long) tv.tv_usec ) <= 0) {
return 0;
}
@@ -35,13 +44,25 @@ int spool_request(char *spooldir, char *filename, int rn, char *buffer, int size
return 0;
}
fprintf(stderr,"writing %d bytes to spool file.\n",size);
uh.modifier1 = 17 ;
uh.modifier2 = 0 ;
uh.datasize = (uint16_t) size ;
#ifdef __BIG_ENDIAN__
uh.datasize= uwsgi_swap16(uh.datasize);
#endif
if (write(fd, &uh, 4) != 4) {
goto clear ;
}
if (write(fd, buffer, size) != size) {
goto clear ;
}
close(fd);
fprintf(stderr,"written %d bytes to spool file %s.\n",size + 4, filename);
return 1;
@@ -53,13 +74,16 @@ clear:
return 0;
}
void spooler(char *spooldir, PyObject *uwsgi_module) {
void spooler(PyObject *uwsgi_module) {
DIR *sdir ;
struct dirent *dp;
PyObject *uwsgi_module_dict, *spooler_callable, *spool_result, *spool_tuple, *spool_env ;
int spool_fd ;
uint16_t uwstrlen ;
int rlen;
int rlen = 0;
int datasize ;
struct uwsgi_packet_header uh ;
char *key;
char *val;
@@ -79,7 +103,18 @@ void spooler(char *spooldir, PyObject *uwsgi_module) {
}
if (chdir(spooldir)) {
spool_env = PyDict_New();
if (!spool_env) {
fprintf(stderr,"could not create spooler env.\n");
exit(1);
}
if (PyTuple_SetItem(spool_tuple, 0, spool_env)) {
PyErr_Print();
exit(1);
}
if (chdir(spool_dir)) {
perror("chdir()");
exit(1);
}
@@ -110,25 +145,36 @@ void spooler(char *spooldir, PyObject *uwsgi_module) {
}
spool_fd = open(dp->d_name, O_RDONLY) ;
if (spool_fd < 0) {
perror("open()");
continue;
}
if (flock(spool_fd, LOCK_EX)) {
perror("flock()");
close(spool_fd);
continue;
}
if (spool_fd < 0) {
perror("open()");
continue;
}
spool_env = PyDict_New();
if (!spool_env) {
PyErr_Print();
if (read(spool_fd, &uh, 4) != 4) {
perror("read()");
close(spool_fd);
continue;
}
while( (rlen = read(spool_fd, &uwstrlen, 2) ) == 2) {
#ifdef __BIG_ENDIAN__
uh.datasize= uwsgi_swap16(uh.datasize);
#endif
datasize = 0 ;
while( datasize < uh.datasize) {
rlen = read(spool_fd, &uwstrlen, 2) ;
if (rlen != 2) {
perror("read()");
goto next_spool ;
}
datasize += rlen ;
key = NULL;
val = NULL;
if (uwstrlen > 0) {
@@ -143,6 +189,7 @@ void spooler(char *spooldir, PyObject *uwsgi_module) {
free(key);
goto next_spool;
}
datasize += rlen ;
key[rlen] = 0 ;
@@ -152,6 +199,7 @@ void spooler(char *spooldir, PyObject *uwsgi_module) {
free(key);
goto next_spool;
}
datasize += rlen ;
if (uwstrlen > 0) {
val = malloc(uwstrlen+1);
@@ -167,6 +215,7 @@ void spooler(char *spooldir, PyObject *uwsgi_module) {
free(key);
goto next_spool;
}
datasize += rlen ;
val[rlen] = 0 ;
/* ready to add item to the dict */
}
@@ -187,10 +236,6 @@ void spooler(char *spooldir, PyObject *uwsgi_module) {
}
if (PyTuple_SetItem(spool_tuple, 0, spool_env)) {
PyErr_Print();
goto retry_later;
}
spool_result = PyEval_CallObject(spooler_callable, spool_tuple);
if (!spool_result) {
PyErr_Print();
@@ -199,11 +244,14 @@ void spooler(char *spooldir, PyObject *uwsgi_module) {
}
if (PyInt_Check(spool_result)) {
if (PyInt_AsLong(spool_result) == 17) {
Py_DECREF(spool_result);
fprintf(stderr,"retry this task later...\n");
goto retry_later;
}
}
Py_DECREF(spool_result);
fprintf(stderr,"done with task/spool %s\n", dp->d_name);
next_spool:
@@ -213,9 +261,8 @@ next_spool:
exit(1);
}
retry_later:
Py_DECREF(spool_env);
PyDict_Clear(spool_env);
close(spool_fd);
Py_DECREF(spooler_callable);
}
}
}
+3 -1
View File
@@ -31,10 +31,12 @@ def myspooler(env):
print env
for i in range(1,100):
uwsgi.sharedarea_inclong(100)
time.sleep(1)
#time.sleep(1)
uwsgi.spooler = myspooler
print "SPOOLER: ", uwsgi.send_to_spooler({'TESTKEY':'TESTVALUE', 'APPNAME':'uWSGI'})
def helloworld():
return 'Hello World'
+11 -6
View File
@@ -561,6 +561,8 @@ int single_app_mode = 0;
int memory_debug = 0 ;
#endif
char *spool_dir = NULL ;
int main(int argc, char *argv[], char *envp[]) {
struct timeval check_interval = {.tv_sec = 1, .tv_usec = 0 };
@@ -586,7 +588,6 @@ int main(int argc, char *argv[], char *envp[]) {
#endif
#ifndef ROCK_SOLID
char *spool_dir = NULL ;
char spool_filename[1024];
pid_t spooler_pid = 0 ;
#endif
@@ -1113,7 +1114,7 @@ int main(int argc, char *argv[], char *envp[]) {
#ifndef ROCK_SOLID
if (spool_dir != NULL) {
spooler_pid = spooler_start(spool_dir, serverfd, uwsgi_module);
spooler_pid = spooler_start(serverfd, uwsgi_module);
}
#endif
@@ -1226,7 +1227,7 @@ int main(int argc, char *argv[], char *envp[]) {
/* reload the spooler */
if (spool_dir && spooler_pid > 0) {
if (diedpid == spooler_pid) {
spooler_pid = spooler_start(spool_dir,serverfd, uwsgi_module);
spooler_pid = spooler_start(serverfd, uwsgi_module);
continue;
}
}
@@ -1472,7 +1473,7 @@ int main(int argc, char *argv[], char *envp[]) {
}
fprintf(stderr,"managing spool request...\n");
i = spool_request(spool_dir, spool_filename, requests+1, buffer,wsgi_req.size) ;
i = spool_request(spool_filename, requests+1, buffer,wsgi_req.size) ;
wsgi_req.modifier = 255 ;
wsgi_req.size = 0 ;
if (i > 0) {
@@ -2637,6 +2638,10 @@ void init_uwsgi_embedded_module() {
init_uwsgi_module_advanced(new_uwsgi_module);
if (spool_dir != NULL) {
init_uwsgi_module_spooler(new_uwsgi_module);
}
if (sharedareasize > 0 && sharedarea) {
init_uwsgi_module_sharedarea(new_uwsgi_module);
@@ -2645,7 +2650,7 @@ void init_uwsgi_embedded_module() {
#endif
#ifndef ROCK_SOLID
pid_t spooler_start(char *spool_dir, int serverfd, PyObject *uwsgi_module) {
pid_t spooler_start(int serverfd, PyObject *uwsgi_module) {
pid_t pid ;
pid = fork();
@@ -2655,7 +2660,7 @@ pid_t spooler_start(char *spool_dir, int serverfd, PyObject *uwsgi_module) {
}
else if (pid == 0) {
close(serverfd);
spooler(spool_dir, uwsgi_module);
spooler(uwsgi_module);
}
else if (pid > 0) {
fprintf(stderr,"spawned the uWSGI spooler on dir %s with pid %d\n", spool_dir, pid);
+4 -3
View File
@@ -207,11 +207,12 @@ void uwsgi_wsgi_config(void);
void init_uwsgi_module_sharedarea(PyObject *);
void init_uwsgi_module_advanced(PyObject *);
void init_uwsgi_module_spooler(PyObject *);
#ifndef ROCK_SOLID
int spool_request(char *, char *, int, char *, int);
void spooler(char *, PyObject *);
pid_t spooler_start(char *,int, PyObject *);
int spool_request(char *, int, char *, int);
void spooler(PyObject *);
pid_t spooler_start(int, PyObject *);
#endif
void set_harakiri(int);
+124
View File
@@ -4,6 +4,9 @@ extern char *sharedarea ;
extern void *sharedareamutex ;
extern int sharedareasize ;
char *spool_buffer = NULL ;
extern int buffer_size ;
#ifdef __APPLE__
#define LOCK_SHAREDAREA OSSpinLockLock((OSSpinLock *) sharedareamutex);
#define UNLOCK_SHAREDAREA OSSpinLockUnlock((OSSpinLock *) sharedareamutex);
@@ -253,6 +256,98 @@ PyObject *py_uwsgi_sharedarea_read(PyObject *self, PyObject *args) {
return PyString_FromStringAndSize(sharedarea+pos, len);
}
PyObject *py_uwsgi_send_spool(PyObject *self, PyObject *args) {
PyObject *spool_dict, *spool_vars ;
PyObject *zero, *key, *val;
extern int requests ;
uint16_t keysize, valsize ;
char *cur_buf ;
int i ;
char spool_filename[1024];
spool_dict = PyTuple_GetItem(args, 0);
if (!PyDict_Check(spool_dict)) {
Py_INCREF(Py_None);
return Py_None;
}
spool_vars = PyDict_Items(spool_dict);
if (!spool_vars) {
Py_INCREF(Py_None);
return Py_None;
}
cur_buf = spool_buffer ;
for(i=0;i<PyList_Size(spool_vars);i++) {
zero = PyList_GetItem(spool_vars, i);
if (zero) {
if (PyTuple_Check(zero)) {
key = PyTuple_GetItem(zero, 0);
val = PyTuple_GetItem(zero, 1);
if (PyString_Check(key) && PyString_Check(val)) {
keysize = PyString_Size(key) ;
valsize= PyString_Size(val) ;
if (cur_buf + keysize + 2 + valsize + 2 <= spool_buffer + buffer_size) {
#ifdef __BIG_ENDIAN__
keysize = uwsgi_swap16(keysize);
#endif
memcpy(cur_buf, &keysize, 2);
cur_buf+=2;
#ifdef __BIG_ENDIAN__
keysize = uwsgi_swap16(keysize);
#endif
memcpy(cur_buf, PyString_AsString(key), keysize);
cur_buf += keysize;
#ifdef __BIG_ENDIAN__
valsize = uwsgi_swap16(valsize);
#endif
memcpy(cur_buf, &valsize, 2);
cur_buf+=2;
#ifdef __BIG_ENDIAN__
valsize = uwsgi_swap16(valsize);
#endif
memcpy(cur_buf, PyString_AsString(val), valsize);
cur_buf += valsize;
}
else {
Py_DECREF(zero);
Py_INCREF(Py_None);
return Py_None;
}
}
else {
Py_DECREF(zero);
Py_INCREF(Py_None);
return Py_None;
}
}
else {
Py_DECREF(zero);
Py_INCREF(Py_None);
return Py_None;
}
}
else {
Py_INCREF(Py_None);
return Py_None;
}
}
i = spool_request(spool_filename, requests+1, spool_buffer, cur_buf - spool_buffer) ;
if (i > 0) {
return Py_True;
}
Py_DECREF(spool_vars);
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_send_message(PyObject *self, PyObject *args) {
PyObject *arg_host, *arg_port, *arg_modifier1, *arg_modifier2, *arg_message, *arg_timeout;
@@ -278,6 +373,11 @@ PyObject *py_uwsgi_send_message(PyObject *self, PyObject *args) {
}
static PyMethodDef uwsgi_spooler_methods[] = {
{"send_to_spooler", py_uwsgi_send_spool, METH_VARARGS, ""},
{NULL, NULL},
};
static PyMethodDef uwsgi_advanced_methods[] = {
{"send_uwsgi_message", py_uwsgi_send_message, METH_VARARGS, ""},
{NULL, NULL},
@@ -297,6 +397,30 @@ static PyMethodDef uwsgi_sa_methods[] = {
void init_uwsgi_module_spooler(PyObject *current_uwsgi_module) {
PyMethodDef *uwsgi_function;
PyObject *uwsgi_module_dict;
uwsgi_module_dict = PyModule_GetDict(current_uwsgi_module);
if (!uwsgi_module_dict) {
fprintf(stderr,"could not get uwsgi module __dict__\n");
exit(1);
}
spool_buffer = malloc(buffer_size);
if (!spool_buffer) {
perror("malloc()");
exit(1);
}
for (uwsgi_function = uwsgi_spooler_methods; uwsgi_function->ml_name != NULL; uwsgi_function++) {
PyObject *func = PyCFunction_New(uwsgi_function, NULL);
PyDict_SetItemString(uwsgi_module_dict, uwsgi_function->ml_name, func);
Py_DECREF(func);
}
}
void init_uwsgi_module_advanced(PyObject *current_uwsgi_module) {
PyMethodDef *uwsgi_function;
PyObject *uwsgi_module_dict;