diff --git a/spooler.c b/spooler.c index dd6b0660..5070ed8c 100644 --- a/spooler.c +++ b/spooler.c @@ -3,12 +3,21 @@ #include +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); } } } diff --git a/testapp.py b/testapp.py index 09fe0593..52a716ff 100644 --- a/testapp.py +++ b/testapp.py @@ -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' diff --git a/uwsgi.c b/uwsgi.c index 2fefe20a..b669d30c 100644 --- a/uwsgi.c +++ b/uwsgi.c @@ -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); diff --git a/uwsgi.h b/uwsgi.h index 90b48eba..a90d0076 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -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); diff --git a/uwsgi_pymodule.c b/uwsgi_pymodule.c index e2a3ffda..de366ec9 100644 --- a/uwsgi_pymodule.c +++ b/uwsgi_pymodule.c @@ -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 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;