From 867478769228d9b8d2d3d1d431bf8bfde66a018d Mon Sep 17 00:00:00 2001 From: "roberto@sirius.local" Date: Tue, 2 Mar 2010 11:50:18 +0100 Subject: [PATCH] leak free ( i hope ) erlang node --- erlang.c | 184 ++++++++++++++++++++++++++++++++++++++++++------------- uwsgi.c | 6 +- uwsgi.h | 4 +- 3 files changed, 149 insertions(+), 45 deletions(-) diff --git a/erlang.c b/erlang.c index 59bc7f57..e23efa5e 100644 --- a/erlang.c +++ b/erlang.c @@ -4,6 +4,8 @@ extern struct uwsgi_server uwsgi; +static void erlang_log(void); + PyObject *py_erlang_connect(PyObject *self, PyObject *args) { char *erlang_node ; @@ -23,15 +25,28 @@ PyObject *py_erlang_connect(PyObject *self, PyObject *args) { PyObject *py_erlang_recv_message(PyObject *self, PyObject *args) { int erfd ; ErlMessage emsg; + PyObject *pyer; + unsigned char erlang_buffer[8192]; if (!PyArg_ParseTuple(args, "i:erlang_recv_message", &erfd)) { return NULL ; } - if (erl_receive_msg(erfd, (unsigned char *) uwsgi.buffer, uwsgi.buffer_size, &emsg) == ERL_MSG) { - return eterm_to_py(emsg.msg); +cycle: + memset(&emsg, 0, sizeof(ErlMessage)); + if (erl_receive_msg(erfd, erlang_buffer, 8192, &emsg) == ERL_MSG) { + if (emsg.type == ERL_TICK) { goto cycle;} + pyer = eterm_to_py(emsg.msg); + if (emsg.msg) {erl_free_compound(emsg.msg);} + if (emsg.to) {erl_free_compound(emsg.to);} + if (emsg.from) { erl_free_compound(emsg.from);} + if (!pyer) { + goto clear; + } + return pyer ; } +clear: Py_INCREF(Py_None); return Py_None; @@ -40,6 +55,7 @@ PyObject *py_erlang_recv_message(PyObject *self, PyObject *args) { PyObject *py_erlang_send_message(PyObject *self, PyObject *args) { ETERM *pymessage ; + ETERM *pid; PyObject *erfd, *ermessage, *erdest, *zero ; int er_number, er_serial, er_creation ; @@ -73,39 +89,49 @@ PyObject *py_erlang_send_message(PyObject *self, PyObject *args) { if (PyString_Check(erdest)) { if (!erl_reg_send(PyInt_AsLong(erfd), PyString_AsString(erdest), pymessage)) { erl_err_msg("erl_reg_send()"); - goto clear; + goto clear2; } } else if (PyDict_Check(erdest)) { - fprintf(stderr,"ready to send\n"); zero = PyDict_GetItemString(erdest,"node"); - if (!zero) { goto clear; } if (!PyString_Check(zero)) { goto clear; } + if (!zero) { goto clear2; } if (!PyString_Check(zero)) { goto clear2; } er_node = PyString_AsString(zero); zero = PyDict_GetItemString(erdest,"number"); - if (!zero) { goto clear; } if (!PyInt_Check(zero)) { goto clear; } + if (!zero) { goto clear2; } if (!PyInt_Check(zero)) { goto clear2; } er_number = PyInt_AsLong(zero); zero = PyDict_GetItemString(erdest,"serial"); - if (!zero) { goto clear; } if (!PyInt_Check(zero)) { goto clear; } + if (!zero) { goto clear2; } if (!PyInt_Check(zero)) { goto clear2; } er_serial = PyInt_AsLong(zero); zero = PyDict_GetItemString(erdest,"creation"); - if (!zero) { goto clear; } if (!PyInt_Check(zero)) { goto clear; } + if (!zero) { goto clear2; } if (!PyInt_Check(zero)) { goto clear2; } er_creation = PyInt_AsLong(zero); - if (!erl_send(PyInt_AsLong(erfd), erl_mk_pid((const char *) er_node, er_number,er_serial,er_creation), pymessage)) { + pid = erl_mk_pid((const char *) er_node, er_number,er_serial,er_creation) ; + + if (!pid) { goto clear2; } + + if (!erl_send(PyInt_AsLong(erfd), pid, pymessage)) { erl_err_msg("erl_send()"); - goto clear; + erl_free_term(pid); + goto clear2; } + + erl_free_term(pid); } else { goto clear; } + erl_free_compound(pymessage); + Py_INCREF(Py_True); return Py_True; +clear2: + erl_free_compound(pymessage); clear: PyErr_Print(); Py_INCREF(Py_None); @@ -135,7 +161,7 @@ static PyMethodDef uwsgi_erlang_methods[] = { }; -int init_erlang(char *nodename) { +int init_erlang(char *nodename, char *cookie) { struct sockaddr_in e_addr; @@ -143,6 +169,10 @@ int init_erlang(char *nodename) { char *node ; int efd ; int rlen ; + char *cookiefile ; + char *cookiehome ; + char cookievalue[128]; + int cookiefd ; PyMethodDef *uwsgi_function; @@ -154,6 +184,40 @@ int init_erlang(char *nodename) { return -1; } + if (cookie == NULL) { + // get the cookie from the home + cookiehome = getenv("HOME"); + if (!cookiehome) { + fprintf(stderr,"unable to get erlang cookie from your home.\n"); + return -1; + } + cookiefile = malloc(strlen(cookiehome)+1+strlen(".erlang.cookie")+1); + if (!cookiefile) { + perror("malloc()"); + } + cookiefile[0] = 0 ; + strcat(cookiefile, cookiehome); + strcat(cookiefile, "/.erlang.cookie"); + + cookiefd = open(cookiefile, O_RDONLY); + if (cookiefd < 0) { + perror("open()"); + free(cookiefile); + return -1 ; + } + + memset(cookievalue, 0, 128); + if (read(cookiefd, cookievalue, 127) < 1) { + fprintf(stderr,"invalid cookie found in %s\n", cookiefile); + close(cookiefd); + free(cookiefile); + return -1 ; + } + cookie = cookievalue ; + close(cookiefd); + free(cookiefile); + } + node = malloc((ip-nodename)+1); if (node == NULL) { perror("malloc()"); @@ -164,7 +228,7 @@ int init_erlang(char *nodename) { erl_init(NULL, 0); - if (erl_connect_xinit(ip+1, node, nodename, NULL, "RVHWRLDVUWOTBIRRALYZ", 0) == -1) { + if (erl_connect_xinit(ip+1, node, nodename, NULL, cookie, 0) == -1) { fprintf(stderr,"*** unable to initialize erlang c-node ***\n"); return -1; } @@ -237,11 +301,9 @@ ETERM *py_to_eterm(PyObject *pobj) { } if (PyString_Check(pobj)) { - fprintf(stderr,"creating atom from string\n"); eobj = erl_mk_atom(PyString_AsString(pobj)); } else if (PyInt_Check(pobj)) { - fprintf(stderr,"creating ERLANG int\n"); eobj = erl_mk_int(PyInt_AsLong(pobj)); } else if (PyList_Check(pobj)) { @@ -251,7 +313,6 @@ ETERM *py_to_eterm(PyObject *pobj) { eobj2 = py_to_eterm(pobj2); eobj = erl_cons(eobj2, eobj); } - fprintf(stderr,"created list\n"); } else if (PyDict_Check(pobj)) { // a pid @@ -273,12 +334,10 @@ ETERM *py_to_eterm(PyObject *pobj) { } else if (PyTuple_Check(pobj)) { count = PyTuple_Size(pobj); - fprintf(stderr,"PYTUPLE !!!\n"); eobj3 = malloc(sizeof(ETERM *)*count) ; for(i=0;iob_type->tp_name); + } clear: if (eobj == NULL) { @@ -335,7 +397,6 @@ PyObject *eterm_to_py(ETERM *obj) { fprintf(stderr,"FOUND A BINARY %.*s\n", ERL_BIN_SIZE(obj), ERL_BIN_PTR(obj)); break; case ERL_PID: - fprintf(stderr,"FOUND A PID %s %d %d %d\n", ERL_PID_NODE(obj), ERL_PID_NUMBER(obj), ERL_PID_SERIAL(obj), ERL_PID_CREATION(obj)); eobj = PyDict_New(); if (PyDict_SetItemString(eobj, "node", PyString_FromString(ERL_PID_NODE(obj)) )) { PyErr_Print(); break;} if (PyDict_SetItemString(eobj, "number", PyInt_FromLong(ERL_PID_NUMBER(obj)) )) { PyErr_Print(); break;} @@ -361,7 +422,17 @@ void erlang_loop() { ErlMessage em; ETERM *eresponse; - int rlen; + PyObject *callable = PyDict_GetItemString(uwsgi.embedded_dict, "erlang_func"); + if (!callable) { + PyErr_Print(); + fprintf(stderr,"- you have not defined a uwsgi.erlang_func callable, Erlang message manager will be disabled -\n"); + } + + PyObject *pargs = PyTuple_New(1); + if (!pargs) { + PyErr_Print(); + fprintf(stderr,"- error preparing arg tuple for uwsgi.erlang_func callable, Erlang message manager will be disabled -\n"); + } while(uwsgi.workers[uwsgi.mywid].manage_next_request) { @@ -370,44 +441,71 @@ void erlang_loop() { uwsgi.poll.fd = erl_accept(uwsgi.erlangfd, &econn); - fprintf(stderr, "ERL_ACCEPT: %d\n", uwsgi.poll.fd); if (uwsgi.poll.fd >=0) { UWSGI_SET_ERLANGING ; for(;;) { - rlen = erl_receive_msg(uwsgi.poll.fd, (unsigned char *) uwsgi.buffer, uwsgi.buffer_size, &em); - fprintf(stderr,"ERL: %d\n", rlen); - if (rlen == ERL_MSG) { - PyObject *zero = eterm_to_py(em.msg); - PyObject *callable = PyDict_GetItemString(uwsgi.embedded_dict, "erlang_func"); - if (!callable) { - fprintf(stderr,"AIAAA\n"); - PyErr_Print(); - } - PyObject *pargs = PyTuple_New(1); - if (!pargs) { - fprintf(stderr,"oops1\n"); - PyErr_Print(); - } - if (PyTuple_SetItem(pargs,0,zero)) { - fprintf(stderr,"oops2\n"); - PyErr_Print(); - } - PyObject *erlang_result = PyEval_CallObject (callable, pargs); - eresponse = py_to_eterm(erlang_result); + if (erl_receive_msg(uwsgi.poll.fd, (unsigned char *) uwsgi.buffer, uwsgi.buffer_size, &em) == ERL_MSG) { + if (em.type == ERL_TICK) continue; - rlen = erl_send(uwsgi.poll.fd, em.from, eresponse); - fprintf(stderr,"ERL_SEND: %d\n", rlen); + PyObject *zero = eterm_to_py(em.msg); + if (em.msg) {erl_free_compound(em.msg);} + if (em.to) {erl_free_compound(em.to);} + + if (!zero) { + PyErr_Print(); + continue; + } + + if (PyTuple_SetItem(pargs,0,zero)) { + PyErr_Print(); + continue; } + + PyObject *erlang_result = PyEval_CallObject (callable, pargs); + + //Py_DECREF(zero); + + if (erlang_result) { + eresponse = py_to_eterm(erlang_result); + if (eresponse) { + erl_send(uwsgi.poll.fd, em.from, eresponse); + erl_free_compound(eresponse); + } + Py_DECREF(erlang_result); + } + + if (em.from) { erl_free_compound(em.from);} + + uwsgi.workers[0].requests++; + uwsgi.workers[uwsgi.mywid].requests++; + if (uwsgi.options[UWSGI_OPTION_LOGGING]) + erlang_log(); } + else { + break; + } } erl_close_connection(uwsgi.poll.fd); - fprintf(stderr,"CONNECTION CLOSED\n"); UWSGI_UNSET_ERLANGING ; + } } } +static void erlang_log() { + if (uwsgi.options[UWSGI_OPTION_MEMORY_DEBUG]) { + get_memusage(); + } + else { + uwsgi.workers[uwsgi.mywid].rss_size = 0 ; + uwsgi.workers[uwsgi.mywid].vsz_size = 0 ; + } + fprintf(stderr,"[Erlang worker %d pid %d] request %llu done {rss: %llu vsz: %llu}\n", uwsgi.mywid, uwsgi.mypid, uwsgi.workers[uwsgi.mywid].requests, + uwsgi.workers[uwsgi.mywid].rss_size, + uwsgi.workers[uwsgi.mywid].vsz_size); +} + #else #warning "*** Erlang support is disabled ***" #endif diff --git a/uwsgi.c b/uwsgi.c index fe07420f..ae1e5dc7 100644 --- a/uwsgi.c +++ b/uwsgi.c @@ -611,6 +611,7 @@ int main (int argc, char *argv[], char *envp[]) { {"udp", required_argument, 0, LONG_ARGS_UDP}, {"check-interval", required_argument, 0, LONG_ARGS_CHECK_INTERVAL}, {"erlang", required_argument, 0, LONG_ARGS_ERLANG}, + {"erlang-cookie", required_argument, 0, LONG_ARGS_ERLANG_COOKIE}, {0, 0, 0, 0} }; #endif @@ -1096,7 +1097,7 @@ int main (int argc, char *argv[], char *envp[]) { #ifdef UWSGI_ERLANG if (uwsgi.erlang_node) { uwsgi.erlang_nodes = 1; - uwsgi.erlangfd = init_erlang(uwsgi.erlang_node); + uwsgi.erlangfd = init_erlang(uwsgi.erlang_node, uwsgi.erlang_cookie); } #endif @@ -2280,6 +2281,9 @@ void manage_opt(int i, char *optarg) { case LONG_ARGS_ERLANG: uwsgi.erlang_node = optarg; break; + case LONG_ARGS_ERLANG_COOKIE: + uwsgi.erlang_cookie = optarg; + break; #endif case LONG_ARGS_PYTHONPATH: if (uwsgi.python_path_cnt < 63) { diff --git a/uwsgi.h b/uwsgi.h index 1326b1a6..546dbf99 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -99,6 +99,7 @@ PyAPI_FUNC (PyObject *) PyMarshal_ReadObjectFromString (char *, Py_ssize_t); #define LONG_ARGS_LIMIT_AS 17009 #define LONG_ARGS_UDP 17010 #define LONG_ARGS_ERLANG 17011 +#define LONG_ARGS_ERLANG_COOKIE 17012 #define UWSGI_CLEAR_STATUS uwsgi.workers[uwsgi.mywid].status = 0 @@ -282,6 +283,7 @@ struct __attribute__ ((packed)) wsgi_request { char *udp_socket; #ifdef UWSGI_ERLANG char *erlang_node; + char *erlang_cookie; #endif int (**hooks)(struct uwsgi_server *, struct wsgi_request*) ; @@ -500,7 +502,7 @@ int uwsgi_request_ping(struct uwsgi_server*, struct wsgi_request*); #include #include -int init_erlang(char *); +int init_erlang(char *, char *); void erlang_loop(void); PyObject *eterm_to_py(ETERM *); ETERM *py_to_eterm(PyObject *);