leak free ( i hope ) erlang node

This commit is contained in:
roberto@sirius.local
2010-03-02 11:50:18 +01:00
parent 3654016c53
commit 8674787692
3 changed files with 149 additions and 45 deletions
+141 -43
View File
@@ -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;i<count;i++) {
pobj2 = PyTuple_GetItem(pobj, i);
if (!pobj2) {
fprintf(stderr,"oops\n");
break;
}
eobj3[i] = py_to_eterm(pobj2);
@@ -286,6 +345,9 @@ ETERM *py_to_eterm(PyObject *pobj) {
eobj = erl_mk_tuple(eobj3, count);
free(eobj3);
}
else {
fprintf(stderr, "UNMANAGED PYTHON TYPE: %s\n", pobj->ob_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
+5 -1
View File
@@ -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) {
+3 -1
View File
@@ -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 <erl_interface.h>
#include <ei.h>
int init_erlang(char *);
int init_erlang(char *, char *);
void erlang_loop(void);
PyObject *eterm_to_py(ETERM *);
ETERM *py_to_eterm(PyObject *);