From 1d5b05265d15fae8d77d8482d3914aeeea39c771 Mon Sep 17 00:00:00 2001 From: "roberto@debian32" Date: Sun, 23 Oct 2011 19:18:31 +0200 Subject: [PATCH] initial farm implementation --- master_utils.c | 9 +++ mule.c | 106 ++++++++++++++++++++++++++++++-- plugins/python/uwsgi_pymodule.c | 31 ++++++++++ signal.c | 14 +++++ uwsgi.c | 63 +++++++++++++++++++ uwsgi.h | 25 ++++++++ 6 files changed, 244 insertions(+), 4 deletions(-) diff --git a/master_utils.c b/master_utils.c index 01f8d9a0..496f2457 100644 --- a/master_utils.c +++ b/master_utils.c @@ -71,6 +71,15 @@ void uwsgi_fixup_fds(int wid, int muleid) { } } + for(i=0;imules; + + while(umf) { + if (umf->mule->id == muleid) { + return 1; + } + umf = umf->next; + } + + return 0; +} + +int farm_has_signaled(int fd) { + + int i; + for(i=0;imule->id == uwsgi.muleid && uwsgi.farms[i].signal_pipe[1] == fd) { + return 1; + } + umf = umf->next; + } + } + + return 0; +} + +int farm_has_msg(int fd) { + + int i; + for(i=0;imule->id == uwsgi.muleid && uwsgi.farms[i].queue_pipe[1] == fd) { + return 1; + } + umf = umf->next; + } + } + + return 0; +} + + +void uwsgi_mule_add_farm_to_queue(int queue) { + + int i; + for(i=0;inext; + } + + uwsgi_mf = uwsgi_malloc(sizeof(struct uwsgi_mule_farm)); + old_umf->next = uwsgi_mf; + } + + uwsgi_mf->mule = um; + uwsgi_mf->next = NULL; + + return uwsgi_mf; +} + diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index 31b910df..cce27fcf 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -1085,6 +1085,36 @@ PyObject *py_uwsgi_setprocname(PyObject * self, PyObject * args) { return Py_None; } +PyObject *py_uwsgi_farm_msg(PyObject * self, PyObject * args) { + + char *message = NULL; + Py_ssize_t message_len = 0; + char *farm_name = NULL; + ssize_t len; + int i; + + if (!PyArg_ParseTuple(args, "ss#:farm_msg", &farm_name, &message, &message_len)) { + return NULL; + } + + for(i=0;ireceiver, "farm", 4)) { + i = atoi(use->receiver+4); + if (i > uwsgi.farms_cnt || i <= 0) { + uwsgi_log("invalid signal target: %s\n", use->receiver); + } + else { + if (write(uwsgi.farms[i-1].signal_pipe[0], &sig, 1) != 1) { + uwsgi_error("write()"); + uwsgi_log("could not deliver signal %d to farm %d (%s)\n", sig, i, uwsgi.farms[i-1].name); + } + } + } + else { // unregistered signal, sending it to all the workers uwsgi_log("^^^ UNSUPPORTED SIGNAL TARGET: %s ^^^\n", use->receiver); diff --git a/uwsgi.c b/uwsgi.c index c788f457..666d0c22 100644 --- a/uwsgi.c +++ b/uwsgi.c @@ -126,6 +126,7 @@ static struct option long_base_options[] = { {"spooler-chdir", required_argument, 0, LONG_ARGS_SPOOLER_CHDIR}, #endif {"mule", optional_argument, 0, LONG_ARGS_MULE}, + {"farm", required_argument, 0, LONG_ARGS_FARM}, {"disable-logging", no_argument, 0, 'L'}, {"pidfile", required_argument, 0, LONG_ARGS_PIDFILE}, @@ -2253,10 +2254,67 @@ skipzero: exit(1); } + uwsgi.mules[i].id = i+1; + snprintf(uwsgi.mules[i].name, 0xff, "uWSGI mule %d", i+1); } } + if (uwsgi.farms_cnt > 0) { + uwsgi.farms = (struct uwsgi_farm *) mmap(NULL, sizeof(struct uwsgi_farm) * uwsgi.farms_cnt, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANON, -1, 0); + if (!uwsgi.farms) { + uwsgi_error("mmap()"); + exit(1); + } + memset(uwsgi.farms, 0, sizeof(struct uwsgi_farm) * uwsgi.farms_cnt); + + struct uwsgi_string_list *farm_name = uwsgi.farms_list; + for(i=0;ivalue); + + char *mules_list = strchr(farm_value, ':'); + if (!mules_list) { + uwsgi_log("invalid farm value (%s) must be in the form name:mule[,muleN].\n", farm_value); + exit(1); + } + + mules_list[0] = 0; + mules_list++; + + strncpy(uwsgi.farms[i].name, farm_value, 0xff); + + // create the socket pipe + if (socketpair(AF_UNIX, SOCK_STREAM, 0, uwsgi.farms[i].signal_pipe)) { + uwsgi_error("socketpair()\n"); + } + + if (socketpair(AF_UNIX, SOCK_DGRAM, 0, uwsgi.farms[i].queue_pipe)) { + uwsgi_error("socketpair()"); + exit(1); + } + + + char *p = strtok(mules_list, ","); + while(p != NULL) { + struct uwsgi_mule *um = get_mule_by_id( atoi( p ) ); + if (!um) { + uwsgi_log("invalid mule id: %s\n", p); + exit(1); + } + + uwsgi_mule_farm_new(&uwsgi.farms[i].mules, um); + + p = strtok(NULL, ","); + } + uwsgi_log("created farm %d name: %s mules:%s\n", i+1, uwsgi.farms[i].name, strchr(farm_name->value, ':')+1); + + farm_name = farm_name->next; + + } + + } + /* uwsgi.shared->hooks[0] = uwsgi_request_wsgi; @@ -3377,6 +3435,11 @@ static int manage_base_opt(int i, char *optarg) { uwsgi.mules_cnt++; uwsgi_string_new_list(&uwsgi.mules_patches, optarg); return 1; + case LONG_ARGS_FARM: + uwsgi.master_process = 1; + uwsgi.farms_cnt++; + uwsgi_string_new_list(&uwsgi.farms_list, optarg); + return 1; case LONG_ARGS_SOCKET_PROTOCOL: // TODO map each socket to a specific protocol return 1; diff --git a/uwsgi.h b/uwsgi.h index 297916ac..264ae9ee 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -532,6 +532,7 @@ struct uwsgi_opt { #define LONG_ARGS_PROCNAME_APPEND 17154 #define LONG_ARGS_PROCNAME 17155 #define LONG_ARGS_PROCNAME_MASTER 17156 +#define LONG_ARGS_FARM 17157 #define UWSGI_OK 0 @@ -1294,12 +1295,15 @@ struct uwsgi_server { /* the list of mules */ struct uwsgi_string_list *mules_patches; struct uwsgi_mule *mules; + struct uwsgi_string_list *farms_list; + struct uwsgi_farm *farms; pid_t mypid; int mywid; int muleid; int mules_cnt; + int farms_cnt; rlim_t max_fd; @@ -1731,6 +1735,7 @@ struct uwsgi_worker { char name[0xff]; }; + struct uwsgi_mule { int id; pid_t pid; @@ -1753,6 +1758,23 @@ struct uwsgi_mule { char name[0xff]; }; +struct uwsgi_mule_farm { + struct uwsgi_mule *mule; + struct uwsgi_mule_farm *next; +}; + +struct uwsgi_farm { + int id; + char name[0xff]; + + int signal_pipe[2]; + int queue_pipe[2]; + + struct uwsgi_mule_farm *mules; + +}; + + char *uwsgi_get_cwd(void); @@ -2395,6 +2417,9 @@ void http_url_decode(char *, uint16_t *, char *); pid_t uwsgi_fork(char *); +struct uwsgi_mule *get_mule_by_id(int); +struct uwsgi_mule_farm *uwsgi_mule_farm_new(struct uwsgi_mule_farm **, struct uwsgi_mule *); + #ifdef UWSGI_CAP void uwsgi_build_cap(char *); #endif