Compare commits

...
29 Commits
Author SHA1 Message Date
roberto@debian32 bb1ef28d69 updated decorators with @farm([name]) 2011-10-23 19:59:53 +02:00
roberto@debian32 1d5b05265d initial farm implementation 2011-10-23 19:18:31 +02:00
roberto@debian32 d9fef3c30e a little hack for psgi loader waiting for a truly agnostic init_uwsgi_app 2011-10-23 12:14:53 +02:00
roberto@debian32 e5dfb9efc3 improved get_memusage 2011-10-23 12:04:02 +02:00
roberto@debian32 02bab15e86 perl threading improvements 2011-10-23 11:49:25 +02:00
roberto@debian32 6f4ddead09 Added tag 1.0-rc2 for changeset 7a4e021ea7ba 2011-10-23 08:37:03 +02:00
roberto@gambit b361248ff6 fixed add_exported_opts 2011-10-23 08:35:15 +02:00
roberto@debian32 e9c42ceaa5 support for multiple values in options 2011-10-22 18:32:32 +02:00
roberto@debian32 5d7029744a PATH_INFO is correctly escaped in http modes 2011-10-22 09:35:01 +02:00
roberto@gambit 1084a3afa8 fixed mule patch loading 2011-10-20 10:38:59 +02:00
roberto@gambit ed254bffee fixed mule process name 2011-10-20 10:20:44 +02:00
roberto@djangocloud 733784d7dd another procname fix... 2011-10-19 16:13:10 +10:00
roberto@djangocloud cb5b26bcb4 another bunch of procname fixes 2011-10-19 14:26:07 +10:00
roberto@djangocloud 0065152bb1 daemons throttling subsystem 2011-10-18 18:53:22 +10:00
roberto@djangocloud 03266203fd another fix for environ 2011-10-18 18:30:13 +10:00
roberto@djangocloud 4a5f68f9d8 fixed environ 2011-10-18 18:12:39 +10:00
roberto@djangocloud 83e31a1c27 another setprocname fix 2011-10-18 16:20:33 +10:00
roberto@djangocloud 0efed5ef9c various emperor procname fixes 2011-10-18 16:15:06 +10:00
roberto@djangocloud cdfef1281d fixed another spooler compilation bug 2011-10-18 15:12:59 +10:00
roberto@debian32 e00907e142 --spooler-harakiri 2011-10-17 19:07:12 +02:00
roberto@app1 6dcbcf39e3 fixed compilation when spooler is disabled 2011-10-17 03:16:55 -05:00
roberto@debian32 044579a8e0 more verbose annihilation report 2011-10-16 11:16:13 +02:00
roberto@debian32 d2b52b4663 added uwsgi.setprocname() 2011-10-16 10:37:26 +02:00
roberto@debian32 7f75b2c02c added --procname and --procname-master 2011-10-16 10:18:45 +02:00
roberto@freebsd64 24a3f671e5 freebsd setproctitle 2011-10-16 10:06:08 +02:00
roberto@debian32 875100e3c8 check for mules before nuke 2011-10-16 09:58:37 +02:00
roberto@mrspurr 34bc7c0cfd check for procname 2011-10-16 09:54:35 +02:00
roberto@debian32 b39a10ff7c procname initial management 2011-10-16 09:45:41 +02:00
roberto@debian32 02331904de uWSGI 1.0-rc1 2011-10-15 08:47:00 +02:00
23 changed files with 1005 additions and 202 deletions
+2
View File
@@ -32,3 +32,5 @@ f4d3c4dcd7a63e21fcc9d8d567374c69935178bd 0.9.8.3
fb168b0b86169219aa9b8e40f0caa6297cf34dbc 0.9.9-rc1
6aa667612019fcff8e68af9469f73aff6fa378a8 0.9.9-rc2
376ffd8f190623e5e2521fe9fa64bf03c31e71be 0.9.9
7048ae11cfc8453c6bb599af3c0dfacd7ebae43c 1.0-rc1
7a4e021ea7ba2953e474326612c3bb82d6817f0e 1.0-rc2
+13 -13
View File
@@ -179,7 +179,7 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
struct uwsgi_instance *c_ui = ui;
struct uwsgi_instance *n_ui = NULL;
pid_t pid;
char **argv;
char **vassal_argv;
char *uef;
char **uenvs;
int counter;
@@ -340,9 +340,9 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
uct = uct->next;
}
argv = uwsgi_malloc(sizeof(char *) * counter);
vassal_argv = uwsgi_malloc(sizeof(char *) * counter);
// set args
argv[0] = uwsgi.binary_path;
vassal_argv[0] = uwsgi.binary_path;
if (uwsgi.emperor_broodlord) {
colon = strchr(name, ':');
@@ -351,29 +351,29 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
}
}
if (!strcmp(name + (strlen(name) - 4), ".xml"))
argv[1] = "--xml";
vassal_argv[1] = "--xml";
if (!strcmp(name + (strlen(name) - 4), ".ini"))
argv[1] = "--ini";
vassal_argv[1] = "--ini";
if (!strcmp(name + (strlen(name) - 4), ".yml"))
argv[1] = "--yaml";
vassal_argv[1] = "--yaml";
if (!strcmp(name + (strlen(name) - 5), ".yaml"))
argv[1] = "--yaml";
vassal_argv[1] = "--yaml";
if (!strcmp(name + (strlen(name) - 3), ".js"))
argv[1] = "--json";
vassal_argv[1] = "--json";
if (colon) {
colon[0] = ':';
}
argv[2] = name;
vassal_argv[2] = name;
counter = 3;
uct = uwsgi.vassals_templates;
while(uct) {
argv[counter] = "--inherit";
argv[counter+1] = uct->filename;
vassal_argv[counter] = "--inherit";
vassal_argv[counter+1] = uct->filename;
counter+=2;
uct = uct->next;
}
argv[counter] = NULL;
vassal_argv[counter] = NULL;
// close all of the unneded fd
for(i=3;i<sysconf(_SC_OPEN_MAX);i++) {
@@ -392,7 +392,7 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
}
// start !!!
if (execvp(argv[0], argv)) {
if (execvp(vassal_argv[0], vassal_argv)) {
uwsgi_error("execvp()");
}
uwsgi_log("is the uwsgi binary in your system PATH ?\n");
+5 -5
View File
@@ -47,7 +47,7 @@ struct uwsgi_gateway *register_gateway(char *name, void (*loop)(void)) {
}
}
gw_pid = fork();
gw_pid = uwsgi_fork(name);
if (gw_pid < 0) {
uwsgi_error("fork()");
return NULL;
@@ -90,7 +90,7 @@ struct uwsgi_gateway *register_gateway(char *name, void (*loop)(void)) {
ug->loop = loop;
ug->num = num;
uwsgi_log( "spawned uWSGI %s %d (pid: %d)\n", ug->name, ug->num, (int) ug->pid);
uwsgi_log( "spawned %s %d (pid: %d)\n", ug->name, ug->num, (int) ug->pid);
uwsgi.gateways_cnt++;
@@ -103,7 +103,7 @@ void gateway_respawn(int id) {
pid_t gw_pid;
struct uwsgi_gateway *ug = &uwsgi.gateways[id];
gw_pid = fork();
gw_pid = uwsgi_fork(ug->name);
if (gw_pid < 0) {
uwsgi_error("fork()");
return;
@@ -124,10 +124,10 @@ void gateway_respawn(int id) {
ug->pid = gw_pid;
ug->respawns++;
if (ug->respawns == 1) {
uwsgi_log( "spawned uWSGI %s %d (pid: %d)\n", ug->name, ug->num, (int) gw_pid);
uwsgi_log( "spawned %s %d (pid: %d)\n", ug->name, ug->num, (int) gw_pid);
}
else {
uwsgi_log( "respawned uWSGI %s %d (pid: %d)\n", ug->name, ug->num, (int) gw_pid);
uwsgi_log( "respawned %s %d (pid: %d)\n", ug->name, ug->num, (int) gw_pid);
}
}
+18 -17
View File
@@ -129,22 +129,22 @@ void log_request(struct wsgi_request *wsgi_req) {
rlen = writev(2, logvec, logvecpos+1);
}
void get_memusage() {
void get_memusage(uint64_t *rss, uint64_t *vsz) {
#ifdef UNBIT
uwsgi.workers[uwsgi.mywid].vsz_size = syscall(356);
*vsz = syscall(356);
#elif defined(__linux__)
FILE *procfile;
int i;
procfile = fopen("/proc/self/stat", "r");
if (procfile) {
i = fscanf(procfile, "%*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %llu %lld", (unsigned long long *) &uwsgi.workers[uwsgi.mywid].vsz_size, (unsigned long long *) &uwsgi.workers[uwsgi.mywid].rss_size);
i = fscanf(procfile, "%*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %*s %llu %lld", (unsigned long long *) vsz, (unsigned long long *) rss);
if (i != 2) {
uwsgi_log( "warning: invalid record in /proc/self/stat\n");
}
fclose(procfile);
}
uwsgi.workers[uwsgi.mywid].rss_size = uwsgi.workers[uwsgi.mywid].rss_size * uwsgi.page_size;
*rss = *rss * uwsgi.page_size;
#elif defined (__sun__)
psinfo_t info;
int procfd;
@@ -152,8 +152,8 @@ void get_memusage() {
procfd = open("/proc/self/psinfo", O_RDONLY);
if (procfd >= 0) {
if ( read(procfd, (char *) &info, sizeof(info)) > 0) {
uwsgi.workers[uwsgi.mywid].rss_size = (uint64_t) info.pr_rssize * 1024;
uwsgi.workers[uwsgi.mywid].vsz_size = (uint64_t) info.pr_size * 1024;
*rss = (uint64_t) info.pr_rssize * 1024;
*vsz = (uint64_t) info.pr_size * 1024;
}
close(procfd);
}
@@ -164,8 +164,8 @@ void get_memusage() {
mach_msg_type_number_t t_size = sizeof(struct task_basic_info);
if (task_info(mach_task_self(), TASK_BASIC_INFO, (task_info_t) & t_info, &t_size) == KERN_SUCCESS) {
uwsgi.workers[uwsgi.mywid].rss_size = t_info.resident_size;
uwsgi.workers[uwsgi.mywid].vsz_size = t_info.virtual_size;
*rss = t_info.resident_size;
*vsz = t_info.virtual_size;
}
#elif defined(__FreeBSD__) || defined(__NetBSD__) || defined(__DragonFly__) || defined(__OpenBSD__)
kvm_t *kv;
@@ -184,8 +184,8 @@ void get_memusage() {
struct kinfo_proc *kproc;
kproc = kvm_getprocs(kv, KERN_PROC_PID, uwsgi.mypid, &cnt);
if (kproc && cnt > 0) {
uwsgi.workers[uwsgi.mywid].vsz_size = kproc->ki_size;
uwsgi.workers[uwsgi.mywid].rss_size = kproc->ki_rssize * uwsgi.page_size;
*vsz = kproc->ki_size;
*rss = kproc->ki_rssize * uwsgi.page_size;
}
#elif defined(__NetBSD__) || defined(__OpenBSD__)
struct kinfo_proc2 *kproc2;
@@ -193,11 +193,11 @@ void get_memusage() {
kproc2 = kvm_getproc2(kv, KERN_PROC_PID, uwsgi.mypid, sizeof(struct kinfo_proc2), &cnt);
if (kproc2 && cnt > 0) {
#ifdef __OpenBSD__
uwsgi.workers[uwsgi.mywid].vsz_size = (kproc2->p_vm_dsize + kproc2->p_vm_ssize + kproc2->p_vm_tsize) * uwsgi.page_size;
*vsz = (kproc2->p_vm_dsize + kproc2->p_vm_ssize + kproc2->p_vm_tsize) * uwsgi.page_size;
#else
uwsgi.workers[uwsgi.mywid].vsz_size = kproc2->p_vm_msize * uwsgi.page_size;
*vsz = kproc2->p_vm_msize * uwsgi.page_size;
#endif
uwsgi.workers[uwsgi.mywid].rss_size = kproc2->p_vm_rssize * uwsgi.page_size;
*rss = kproc2->p_vm_rssize * uwsgi.page_size;
}
#endif
@@ -207,13 +207,14 @@ void get_memusage() {
area_info ai;
int32 cookie;
uwsgi.workers[uwsgi.mywid].vsz_size = 0;
uwsgi.workers[uwsgi.mywid].rss_size = 0;
*vsz = 0;
*rss = 0;
while(get_next_area_info(0, &cookie, &ai) == B_OK) {
uwsgi.workers[uwsgi.mywid].vsz_size += ai.ram_size;
*vsz += ai.ram_size;
if ( (ai.protection & B_WRITE_AREA) != 0) {
uwsgi.workers[uwsgi.mywid].rss_size += ai.ram_size;
*rss += ai.ram_size;
}
}
#endif
}
+57 -4
View File
@@ -294,6 +294,18 @@ int master_loop(char **argv, char **environ) {
struct uwsgi_rb_timer *min_timeout;
struct rb_root *rb_timers = uwsgi_init_rb_timer();
if (uwsgi.procname_master) {
uwsgi_set_processname(uwsgi.procname_master);
}
else if (uwsgi.procname) {
uwsgi_set_processname(uwsgi.procname);
}
else if (uwsgi.auto_procname) {
uwsgi_set_processname("uWSGI master");
}
uwsgi.current_time = time(NULL);
uwsgi_unix_signal(SIGHUP, grace_them_all);
@@ -661,8 +673,8 @@ healthy:
if ((uwsgi.cheap || ready_to_die >= uwsgi.numproc) && uwsgi.to_hell) {
// call a series of waitpid to ensure all processes (gateways and daemons) are dead
for(i=0;i<(uwsgi.gateways_cnt+ushared->daemons_cnt);i++) {
// call a series of waitpid to ensure all processes (gateways, mules and daemons) are dead
for(i=0;i<(uwsgi.gateways_cnt+ushared->daemons_cnt+uwsgi.mules_cnt);i++) {
diedpid = waitpid(WAIT_ANY, &waitpid_status, WNOHANG);
}
@@ -670,8 +682,8 @@ healthy:
exit(0);
}
if ( (uwsgi.cheap || ready_to_reload >= uwsgi.numproc) && uwsgi.to_heaven) {
// call a series of waitpid to ensure all processes (gateways and daemons) are dead
for(i=0;i<(uwsgi.gateways_cnt+ushared->daemons_cnt);i++) {
// call a series of waitpid to ensure all processes (gateways, mules and daemons) are dead
for(i=0;i<(uwsgi.gateways_cnt+ushared->daemons_cnt+uwsgi.mules_cnt);i++) {
diedpid = waitpid(WAIT_ANY, &waitpid_status, WNOHANG);
}
@@ -1364,6 +1376,15 @@ healthy:
}
}
}
#ifdef UWSGI_SPOOLER
if (uwsgi.shared->spooler_pid > 0 && uwsgi.shared->spooler_harakiri > 0) {
if (uwsgi.shared->spooler_harakiri < (time_t) uwsgi.current_time) {
uwsgi_log("*** HARAKIRI ON THE SPOOLER (pid: %d) ***\n", uwsgi.shared->spooler_pid);
kill(uwsgi.shared->spooler_pid, SIGKILL);
uwsgi.shared->spooler_harakiri = 0;
}
}
#endif
#ifdef UWSGI_UDP
// check for cluster nodes
@@ -1488,6 +1509,37 @@ healthy:
uwsgi.mywid = find_worker_id(diedpid);
if (uwsgi.mywid <= 0) {
// check spooler, mules, gateways and daemons
#ifdef UWSGI_SPOOLER
if (uwsgi.spool_dir && uwsgi.shared->spooler_pid > 0) {
if (diedpid == uwsgi.shared->spooler_pid) {
uwsgi_log("spooler (pid: %d) annihilated\n", (int) diedpid);
goto next;
}
}
#endif
for(i=0;i<uwsgi.mules_cnt;i++) {
if (uwsgi.mules[i].pid == diedpid) {
uwsgi_log("mule %d (pid: %d) annihilated\n", i+1, (int) diedpid);
goto next;
}
}
for(i=0;i<uwsgi.gateways_cnt;i++) {
if (uwsgi.gateways[i].pid == diedpid) {
uwsgi_log("gateway %d (pid: %d) annihilated\n", i+1, (int) diedpid);
goto next;
}
}
for(i=0;i<uwsgi.shared->daemons_cnt;i++) {
if (uwsgi.shared->daemons[i].pid == diedpid) {
uwsgi_log("daemon %d (pid: %d) annihilated\n", i+1, (int) diedpid);
goto next;
}
}
if (WIFEXITED(waitpid_status)) {
uwsgi_log("subprocess %d exited with code %d\n", (int) diedpid, WEXITSTATUS(waitpid_status));
}
@@ -1497,6 +1549,7 @@ healthy:
else if (WIFSTOPPED(waitpid_status)) {
uwsgi_log("subprocess %d stopped\n", (int) diedpid);
}
next:
continue;
}
else {
+10 -1
View File
@@ -71,6 +71,15 @@ void uwsgi_fixup_fds(int wid, int muleid) {
}
}
for(i=0;i<uwsgi.farms_cnt;i++) {
if (uwsgi.farms[i].signal_pipe[0] != -1) close(uwsgi.farms[i].signal_pipe[0]);
if (muleid == 0) {
if (uwsgi.farms[i].signal_pipe[1] != -1) close(uwsgi.farms[i].signal_pipe[1]);
if (uwsgi.farms[i].queue_pipe[1] != -1) close(uwsgi.farms[i].queue_pipe[1]);
}
}
}
@@ -81,7 +90,7 @@ int uwsgi_respawn_worker(int wid) {
int respawns = uwsgi.workers[wid].respawn_count;
int i;
pid_t pid = fork();
pid_t pid = uwsgi_fork(uwsgi.workers[wid].name);
if (pid == 0) {
signal(SIGWINCH, worker_wakeup);
+104 -6
View File
@@ -16,7 +16,7 @@ void uwsgi_mule(int id) {
int i;
pid_t pid = fork();
pid_t pid = uwsgi_fork(uwsgi.mules[id-1].name);
if (pid == 0) {
uwsgi.muleid = id;
// avoid race conditions
@@ -44,10 +44,10 @@ void uwsgi_mule(int id) {
if (uwsgi.mules[id-1].patch) {
uwsgi_log("loading patch %s\n", uwsgi.mules[id-1].patch);
for (i = 0; i < 0xFF; i++) {
if (uwsgi.p[i]->mule) {
if (uwsgi.p[i]->mule(uwsgi.mules[id-1].patch) == 1) {
uwsgi_log("loaded mule patch %s\n", uwsgi.mules[id-1].patch);
break;
}
}
@@ -63,6 +63,64 @@ void uwsgi_mule(int id) {
}
}
int uwsgi_farm_has_mule(struct uwsgi_farm *farm, int muleid) {
struct uwsgi_mule_farm *umf = farm->mules;
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;i<uwsgi.farms_cnt;i++) {
struct uwsgi_mule_farm *umf = uwsgi.farms[i].mules;
while(umf) {
if (umf->mule->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;i<uwsgi.farms_cnt;i++) {
struct uwsgi_mule_farm *umf = uwsgi.farms[i].mules;
while(umf) {
if (umf->mule->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;i<uwsgi.farms_cnt;i++) {
if (uwsgi_farm_has_mule(&uwsgi.farms[i], uwsgi.muleid)) {
event_queue_add_fd_read(queue, uwsgi.farms[i].signal_pipe[1]);
event_queue_add_fd_read(queue, uwsgi.farms[i].queue_pipe[1]);
}
}
}
void uwsgi_mule_handler() {
ssize_t len;
@@ -79,13 +137,15 @@ void uwsgi_mule_handler() {
event_queue_add_fd_read(mule_queue, uwsgi.my_signal_socket);
event_queue_add_fd_read(mule_queue, uwsgi.mules[uwsgi.muleid-1].queue_pipe[1]);
uwsgi_mule_add_farm_to_queue(mule_queue);
for(;;) {
rlen = event_queue_wait(mule_queue, -1, &interesting_fd);
if (rlen <= 0) {
continue;
}
if (interesting_fd == uwsgi.signal_socket || interesting_fd == uwsgi.my_signal_socket) {
if (interesting_fd == uwsgi.signal_socket || interesting_fd == uwsgi.my_signal_socket || farm_has_signaled(interesting_fd)) {
len = read(interesting_fd, &uwsgi_signal, 1);
if (len <= 0) {
uwsgi_log_verbose("uWSGI mule %d braying: my master died, i will follow him...\n", uwsgi.muleid);
@@ -98,8 +158,8 @@ void uwsgi_mule_handler() {
uwsgi_log_verbose("error managing signal %d on mule %d\n", uwsgi_signal, uwsgi.mywid);
}
}
else if (interesting_fd == uwsgi.mules[uwsgi.muleid-1].queue_pipe[1]) {
len = read(uwsgi.mules[uwsgi.muleid-1].queue_pipe[1], message, 65536);
else if (interesting_fd == uwsgi.mules[uwsgi.muleid-1].queue_pipe[1] || farm_has_msg(interesting_fd)) {
len = read(interesting_fd, message, 65536);
if (len < 0) {
uwsgi_error("read()");
}
@@ -107,9 +167,47 @@ void uwsgi_mule_handler() {
exit(1);
}
else {
uwsgi_log("*** mule %d received a %d bytes message\n", uwsgi.muleid, len);
uwsgi_log("*** mule %d received a %d bytes message ***\n", uwsgi.muleid, len);
}
}
}
}
struct uwsgi_mule *get_mule_by_id(int id) {
int i;
for(i=0;i<uwsgi.mules_cnt;i++) {
if (uwsgi.mules[i].id == id) {
return &uwsgi.mules[i];
}
}
return NULL;
}
struct uwsgi_mule_farm *uwsgi_mule_farm_new(struct uwsgi_mule_farm **umf, struct uwsgi_mule *um) {
struct uwsgi_mule_farm *uwsgi_mf = *umf, *old_umf;
if (!uwsgi_mf) {
*umf = uwsgi_malloc(sizeof(struct uwsgi_mule_farm));
uwsgi_mf = *umf;
}
else {
while(uwsgi_mf) {
old_umf = uwsgi_mf;
uwsgi_mf = uwsgi_mf->next;
}
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;
}
+1 -1
View File
@@ -422,7 +422,7 @@ int erlang_init() {
uwsgi_log("Erlang C-Node %s registered on port %d\n", ei_thisnodename(&uerl.cnode), ntohs(sin.sin_port));
if (register_fat_gateway("erlang", erlang_loop) == NULL) {
if (register_fat_gateway("uWSGI erlang c-node", erlang_loop) == NULL) {
uwsgi_log("unable to register the erlang gateway\n");
exit(1);
}
+2 -2
View File
@@ -701,13 +701,13 @@ int fastrouter_init() {
if (!ufr.nevents) ufr.nevents = 64;
if (ufr.code_string_code && ufr.code_string_function) {
if (register_fat_gateway("fastrouter", fastrouter_loop) == NULL) {
if (register_fat_gateway("uWSGI fastrouter", fastrouter_loop) == NULL) {
uwsgi_log("unable to register the fastrouter gateway\n");
exit(1);
}
}
else {
if (register_gateway("fastrouter", fastrouter_loop) == NULL) {
if (register_gateway("uWSGI fastrouter", fastrouter_loop) == NULL) {
uwsgi_log("unable to register the fastrouter gateway\n");
exit(1);
}
+11 -3
View File
@@ -136,6 +136,8 @@ struct http_session {
char uss[MAX_HTTP_VEC*2];
char buffer[UMAX16];
char path_info[UMAX16];
uint16_t path_info_len;
struct uwsgi_subscriber_name *un;
@@ -325,13 +327,19 @@ int http_parse(struct http_session *h_session) {
base = ptr;
while(ptr < watermark) {
if (*ptr == '?' && !query_string) {
h_session->uh.pktsize += http_add_uwsgi_var(h_session->iov, h_session->uss+c, h_session->uss+c+2, "PATH_INFO", 9, base, ptr-base, &c);
// PATH_INFO must be url-decoded !!!
h_session->path_info_len = ptr-base;
http_url_decode(base, &h_session->path_info_len, h_session->path_info);
h_session->uh.pktsize += http_add_uwsgi_var(h_session->iov, h_session->uss+c, h_session->uss+c+2, "PATH_INFO", 9, h_session->path_info, h_session->path_info_len, &c);
query_string = ptr+1;
}
else if (*ptr == ' ') {
h_session->uh.pktsize += http_add_uwsgi_var(h_session->iov, h_session->uss+c, h_session->uss+c+2, "REQUEST_URI", 11, base, ptr-base, &c);
if (!query_string) {
h_session->uh.pktsize += http_add_uwsgi_var(h_session->iov, h_session->uss+c, h_session->uss+c+2, "PATH_INFO", 9, base, ptr-base, &c);
// PATH_INFO must be url-decoded !!!
h_session->path_info_len = ptr-base;
http_url_decode(base, &h_session->path_info_len, h_session->path_info);
h_session->uh.pktsize += http_add_uwsgi_var(h_session->iov, h_session->uss+c, h_session->uss+c+2, "PATH_INFO", 9, h_session->path_info, h_session->path_info_len, &c);
h_session->uh.pktsize += http_add_uwsgi_var(h_session->iov, h_session->uss+c, h_session->uss+c+2, "QUERY_STRING", 12, "", 0, &c);
}
else {
@@ -883,7 +891,7 @@ int http_init() {
uhttp.server = bind_to_tcp(uhttp.socket_name, uwsgi.listen_queue, strchr(uhttp.socket_name,':'));
if (register_gateway("http", http_loop) == NULL) {
if (register_gateway("uWSGI http", http_loop) == NULL) {
uwsgi_log("unable to register the http gateway\n");
exit(1);
}
+1 -1
View File
@@ -34,7 +34,7 @@ struct uwsgi_perl {
void init_perl_embedded_module(void);
void uwsgi_psgi_app(void);
int psgi_response(struct wsgi_request *, PerlInterpreter *, AV*);
int psgi_response(struct wsgi_request *, AV*);
#define psgi_xs(func) newXS("uwsgi::" #func, XS_##func, "uwsgi")
#define psgi_check_args(x) if (items < x) Perl_croak(aTHX_ "Usage: uwsgi::%s takes %d arguments", __FUNCTION__ + 3, x)
+26 -14
View File
@@ -41,19 +41,11 @@ XS(XS_stream)
AV *response = (AV* ) SvREFCNT_inc(SvRV(ST(0))) ;
if (av_len(response) == 2) {
#ifdef my_perl
while (psgi_response(wsgi_req, my_perl, response) != UWSGI_OK);
#else
while (psgi_response(wsgi_req, uperl.main, response) != UWSGI_OK);
#endif
while (psgi_response(wsgi_req, response) != UWSGI_OK);
}
else if (av_len(response) == 1) {
#ifdef my_perl
while (psgi_response(wsgi_req, my_perl, response) != UWSGI_OK);
#else
while (psgi_response(wsgi_req, uperl.main, response) != UWSGI_OK);
#endif
while (psgi_response(wsgi_req, response) != UWSGI_OK);
SvREFCNT_dec(response);
ST(0) = sv_bless(newRV(sv_newmortal()), uperl.streaming_stash);
XSRETURN(1);
@@ -236,6 +228,7 @@ xs_init(pTHX)
void uwsgi_psgi_app() {
struct stat stat_psgi;
int id = uwsgi_apps_cnt;
if (uperl.psgi) {
@@ -291,7 +284,6 @@ void uwsgi_psgi_app() {
uperl.psgibuffer[stat_psgi.st_size] = 0;
if (uwsgi.threads < 2) {
uperl.psgi_main = perl_eval_pv(uwsgi_concat4("#line 1 ", uperl.psgi, "\n", uperl.psgibuffer), 0);
if (!uperl.psgi_main) {
uwsgi_log("unable to find PSGI function entry point.\n");
@@ -303,11 +295,31 @@ void uwsgi_psgi_app() {
exit(1);
}
if (uwsgi.threads < 2) {
free(uperl.psgibuffer);
close(uperl.fd);
}
}
struct uwsgi_app *wi = &uwsgi_apps[id];
memset(wi, 0, sizeof(struct uwsgi_app));
wi->modifier1 = 5;
wi->mountpoint = "";
wi->mountpoint_len = 0;
uwsgi_log("PSGI app %d (%s) loaded at %p\n", id, uperl.psgi, uperl.psgi_main);
uwsgi_apps_cnt++;
// check if we need to emulate fork() COW
int i;
if (uwsgi.mywid == 0) {
for(i=1;i<=uwsgi.numproc;i++) {
memcpy(&uwsgi.workers[i].apps[id], &uwsgi.workers[0].apps[id], sizeof(struct uwsgi_app));
uwsgi.workers[i].apps_cnt = uwsgi_apps_cnt;
}
}
uwsgi_log("PSGI app (%s) loaded at %p\n", uperl.psgi, uperl.psgi_main);
}
+41 -62
View File
@@ -23,7 +23,7 @@ SV *uwsgi_perl_obj_new(char *class, size_t class_len) {
SV *newobj;
// set current context ?
//dTHX;
dTHX;
dSP;
ENTER;
@@ -49,7 +49,7 @@ SV *uwsgi_perl_call_stream(SV *func) {
SV *ret = NULL;
// set current context ?
//dTHX;
dTHX;
dSP;
ENTER;
@@ -79,7 +79,7 @@ int uwsgi_perl_obj_can(SV *obj, char *method, size_t len) {
int ret;
// set current context ? needed for threading
//dTHX;
dTHX;
dSP;
ENTER;
@@ -107,6 +107,7 @@ SV *uwsgi_perl_obj_call(SV *obj, char *method) {
SV *ret = NULL;
dTHX;
dSP;
ENTER;
@@ -140,6 +141,7 @@ AV *psgi_call(struct wsgi_request *wsgi_req, SV *psgi_func, SV *env) {
AV *ret = NULL;
dTHX;
dSP;
ENTER;
@@ -149,6 +151,7 @@ AV *psgi_call(struct wsgi_request *wsgi_req, SV *psgi_func, SV *env) {
PUTBACK;
call_sv(psgi_func, G_SCALAR | G_EVAL);
SPAGAIN;
if(SvTRUE(ERRSV)) {
@@ -326,51 +329,19 @@ int uwsgi_perl_init(){
}
#ifdef my_perl
#ifdef USE_ITHREADS
void uwsgi_perl_enable_threads() {
int i;
if (pthread_key_create(&uperl.u_interpreter, NULL)) {
uwsgi_error("pthread_key_create()");
exit(1);
}
uperl.interp = malloc( sizeof(PerlInterpreter*) * uwsgi.threads );
if (!uperl.interp) {
uwsgi_error("malloc()");
exit(1);
}
uperl.interp = uwsgi_malloc( sizeof(PerlInterpreter*) * uwsgi.threads );
uperl.interp[0] = uperl.main;
for(i=1;i<uwsgi.threads;i++) {
uperl.interp[i] = perl_clone(uperl.main, CLONEf_KEEP_PTR_TABLE);
if (!uperl.interp[i]) {
uwsgi_log("unable to create new perl interpreter\n");
exit(1);
}
}
uperl.psgi_func = uwsgi_malloc( sizeof(SV*) * uwsgi.threads );
uperl.psgi_func = malloc( sizeof(SV*) * uwsgi.threads );
if (!uperl.psgi_func) {
uwsgi_error("malloc()");
exit(1);
}
dTHXa(uperl.main);
PERL_SET_CONTEXT(uperl.main);
uperl.psgi_main = perl_eval_pv(uperl.psgibuffer, 0);
if (!uperl.psgi_main) {
uwsgi_log("unable to find PSGI function entry point.\n");
exit(1);
}
if(SvTRUE(ERRSV)) {
uwsgi_log("%s\n", SvPV_nolen(ERRSV));
exit(1);
}
}
#endif
@@ -381,12 +352,10 @@ int uwsgi_perl_request(struct wsgi_request *wsgi_req) {
SV **harakiri;
SV *psgi_func = uperl.psgi_main;
// ugly hack
register PerlInterpreter *my_perl = uperl.main;
#ifdef UWSGI_ASYNC
if (wsgi_req->async_status == UWSGI_AGAIN) {
return psgi_response(wsgi_req, my_perl, wsgi_req->async_placeholder);
return psgi_response(wsgi_req, wsgi_req->async_placeholder);
}
#endif
@@ -401,19 +370,14 @@ int uwsgi_perl_request(struct wsgi_request *wsgi_req) {
return -1;
}
#ifdef UWSGI_THREADING
if (uwsgi.threads > 1 && wsgi_req->async_id > 0) {
if (uwsgi.threads > 1) {
psgi_func = uperl.psgi_func[wsgi_req->async_id];
my_perl = pthread_getspecific(uperl.u_interpreter);
}
#endif
ENTER;
SAVETMPS;
wsgi_req->async_environ = build_psgi_env(wsgi_req);
if (!wsgi_req->async_environ) goto clear;
@@ -431,7 +395,7 @@ int uwsgi_perl_request(struct wsgi_request *wsgi_req) {
goto clear2;
}
while (psgi_response(wsgi_req, my_perl, wsgi_req->async_result) != UWSGI_OK) {
while (psgi_response(wsgi_req, wsgi_req->async_result) != UWSGI_OK) {
#ifdef UWSGI_ASYNC
if (uwsgi.async > 1) {
FREETMPS;
@@ -471,26 +435,41 @@ void uwsgi_perl_after_request(struct wsgi_request *wsgi_req) {
}
#ifdef USE_ITHREADS
void uwsgi_perl_init_thread(int core_id) {
pthread_setspecific(uperl.u_interpreter, uperl.interp[core_id]);
if (core_id > 0) {
uperl.interp[core_id] = perl_clone(uperl.main, CLONEf_KEEP_PTR_TABLE);
if (!uperl.interp[core_id]) {
uwsgi_log("unable to create new perl interpreter\n");
exit(1);
}
}
dTHXa(uperl.interp[core_id]);
PERL_SET_CONTEXT(uperl.interp[core_id]);
uperl.psgi_func[core_id] = perl_eval_pv(uperl.psgibuffer, 0);
if (!uperl.psgi_func[core_id]) {
uwsgi_log("unable to find PSGI function entry point.\n");
exit(1);
}
pthread_setspecific(uperl.u_interpreter, uperl.interp[core_id]);
if(SvTRUE(ERRSV)) {
uwsgi_log("%s\n", SvPV_nolen(ERRSV));
exit(1);
}
uperl.psgi_func[0] = uperl.psgi_main;
if (core_id > 0) {
uperl.psgi_func[core_id] = perl_eval_pv(uwsgi_concat4("#line 1 ", uperl.psgi, "\n", uperl.psgibuffer), 0);
if (!uperl.psgi_func[core_id]) {
uwsgi_log("unable to find PSGI function entry point.\n");
exit(1);
}
if(SvTRUE(ERRSV)) {
uwsgi_log("%s\n", SvPV_nolen(ERRSV));
exit(1);
}
}
}
#endif
int uwsgi_perl_manage_options(int i, char *optarg) {
@@ -531,11 +510,11 @@ struct uwsgi_plugin psgi_plugin = {
.init_apps = uwsgi_psgi_app,
//.magic = uwsgi_perl_magic,
//.help = uwsgi_perl_help,
#ifdef my_perl
#ifdef USE_ITHREADS
.enable_threads = uwsgi_perl_enable_threads,
.init_thread = uwsgi_perl_init_thread,
#endif
.manage_opt = uwsgi_perl_manage_options,
.init_thread = uwsgi_perl_init_thread,
.request = uwsgi_perl_request,
.after_request = uwsgi_perl_after_request,
+1 -1
View File
@@ -4,7 +4,7 @@ extern struct uwsgi_server uwsgi;
extern struct http_status_codes hsc[];
int psgi_response(struct wsgi_request *wsgi_req, PerlInterpreter *my_perl, AV *response) {
int psgi_response(struct wsgi_request *wsgi_req, AV *response) {
SV **status_code, **hitem ;
AV *headers, *body =NULL;
+138 -3
View File
@@ -1072,6 +1072,83 @@ PyObject *py_uwsgi_embedded_data(PyObject * self, PyObject * args) {
}
PyObject *py_uwsgi_setprocname(PyObject * self, PyObject * args) {
char *name = NULL;
if (!PyArg_ParseTuple(args, "s:setprocname", &name)) {
return NULL;
}
uwsgi_set_processname(name);
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_in_farm(PyObject * self, PyObject * args) {
char *farm_name = NULL;
int i;
if (!PyArg_ParseTuple(args, "|s:in_farm", &farm_name)) {
return NULL;
}
if (uwsgi.muleid == 0) goto none;
for(i=0;i<uwsgi.farms_cnt;i++) {
if (!farm_name) {
if (uwsgi_farm_has_mule(&uwsgi.farms[i], uwsgi.muleid)) {
Py_INCREF(Py_True);
return Py_True;
}
}
else {
if (!strcmp(farm_name, uwsgi.farms[i].name)) {
if (uwsgi_farm_has_mule(&uwsgi.farms[i], uwsgi.muleid)) {
Py_INCREF(Py_True);
return Py_True;
}
}
}
}
none:
Py_INCREF(Py_None);
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;i<uwsgi.farms_cnt;i++) {
if (!strcmp(farm_name, uwsgi.farms[i].name)) {
len = write(uwsgi.farms[i].queue_pipe[0], message, message_len);
if (len <= 0) {
uwsgi_error("write()");
}
break;
}
}
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_mule_msg(PyObject * self, PyObject * args) {
char *message = NULL;
@@ -1118,6 +1195,57 @@ PyObject *py_uwsgi_mule_get_msg(PyObject * self, PyObject * args) {
return PyString_FromStringAndSize(message, len);
}
PyObject *py_uwsgi_farm_get_msg(PyObject * self, PyObject * args) {
ssize_t len = 0;
// this buffer will be configurable
char message[65536];
int i, count = 0, pos = 0, ret;
struct pollfd *farmpoll;
if (uwsgi.muleid == 0) {
return PyErr_Format(PyExc_ValueError, "you can receive mule messages only in a mule !!!");
}
UWSGI_RELEASE_GIL;
for(i=0;i<uwsgi.farms_cnt;i++) {
if (uwsgi_farm_has_mule(&uwsgi.farms[i], uwsgi.muleid)) count++;
}
farmpoll = uwsgi_malloc( sizeof(struct pollfd) * count);
for(i=0;i<uwsgi.farms_cnt;i++) {
if (uwsgi_farm_has_mule(&uwsgi.farms[i], uwsgi.muleid)) {
farmpoll[pos].fd = uwsgi.farms[i].queue_pipe[1];
farmpoll[pos].events = POLLIN;
pos++;
}
}
ret = poll(farmpoll, count, -1);
if (ret <= 0) {
uwsgi_error("poll()");
free(farmpoll);
Py_INCREF(Py_None);
return Py_None;
}
for(i=0;i<count;i++) {
if (farmpoll[i].revents & POLLIN) {
len = read(farmpoll[i].fd, message, 65536);
break;
}
}
UWSGI_GET_GIL;
if (len <= 0) {
uwsgi_error("read()");
free(farmpoll);
Py_INCREF(Py_None);
return Py_None;
}
free(farmpoll);
return PyString_FromStringAndSize(message, len);
}
PyObject *py_uwsgi_extract(PyObject * self, PyObject * args) {
char *name;
@@ -2479,12 +2607,13 @@ PyObject *py_uwsgi_logsize(PyObject * self, PyObject * args) {
PyObject *py_uwsgi_mem(PyObject * self, PyObject * args) {
uint64_t rss=0, vsz = 0;
PyObject *ml = PyTuple_New(2);
get_memusage();
get_memusage(&rss, &vsz);
PyTuple_SetItem(ml, 0, PyLong_FromUnsignedLongLong(uwsgi.workers[uwsgi.mywid].rss_size));
PyTuple_SetItem(ml, 1, PyLong_FromUnsignedLongLong(uwsgi.workers[uwsgi.mywid].vsz_size));
PyTuple_SetItem(ml, 0, PyLong_FromUnsignedLongLong(rss));
PyTuple_SetItem(ml, 1, PyLong_FromUnsignedLongLong(vsz));
return ml;
@@ -2636,6 +2765,7 @@ PyObject *py_uwsgi_grunt(PyObject * self, PyObject * args) {
goto clear;
}
// use a normal fork here
grunt_pid = fork();
if (grunt_pid < 0) {
uwsgi_error("fork()");
@@ -2809,6 +2939,8 @@ static PyMethodDef uwsgi_advanced_methods[] = {
{"unlock", py_uwsgi_unlock, METH_VARARGS, ""},
{"cl", py_uwsgi_cl, METH_VARARGS, ""},
{"setprocname", py_uwsgi_setprocname, METH_VARARGS, ""},
{"listen_queue", py_uwsgi_listen_queue, METH_VARARGS, ""},
{"attach_daemon", py_uwsgi_attach_daemon, METH_VARARGS, ""},
@@ -2868,7 +3000,10 @@ static PyMethodDef uwsgi_advanced_methods[] = {
{"extract", py_uwsgi_extract, METH_VARARGS, ""},
{"mule_msg", py_uwsgi_mule_msg, METH_VARARGS, ""},
{"farm_msg", py_uwsgi_farm_msg, METH_VARARGS, ""},
{"mule_get_msg", py_uwsgi_mule_get_msg, METH_VARARGS, ""},
{"farm_get_msg", py_uwsgi_farm_get_msg, METH_VARARGS, ""},
{"in_farm", py_uwsgi_in_farm, METH_VARARGS, ""},
//{"call_hook", py_uwsgi_call_hook, METH_VARARGS, ""},
{NULL, NULL},
+24 -6
View File
@@ -124,9 +124,18 @@ static int http_parse(struct wsgi_request *wsgi_req, char *watermark) {
base = ptr;
while (ptr < watermark) {
if (*ptr == '?' && !query_string) {
wsgi_req->uh.pktsize += proto_base_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, base, ptr - base);
wsgi_req->path_info = (wsgi_req->buffer + wsgi_req->uh.pktsize) - (ptr - base);
wsgi_req->path_info_len = ptr - base;
if (watermark + (ptr - base) < (char *)(wsgi_req->proto_parser_buf + uwsgi.buffer_size)) {
wsgi_req->path_info = watermark;
wsgi_req->path_info_len = ptr - base;
http_url_decode(base, &wsgi_req->path_info_len, wsgi_req->path_info);
wsgi_req->uh.pktsize += proto_base_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, wsgi_req->path_info, wsgi_req->path_info_len);
}
else {
uwsgi_log("not enough space in wsgi_req http proto_parser_buf to encode PATH_INFO, consider tuning it with --buffer-size\n");
wsgi_req->uh.pktsize += proto_base_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, base, ptr - base);
wsgi_req->path_info = (wsgi_req->buffer + wsgi_req->uh.pktsize) - (ptr - base);
wsgi_req->path_info_len = ptr - base;
}
wsgi_req->path_info_pos = 3;
query_string = ptr + 1;
}
@@ -136,9 +145,18 @@ static int http_parse(struct wsgi_request *wsgi_req, char *watermark) {
wsgi_req->uri_len = ptr - base;
if (!query_string) {
wsgi_req->uh.pktsize += proto_base_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, base, ptr - base);
wsgi_req->path_info = (wsgi_req->buffer + wsgi_req->uh.pktsize) - (ptr - base);
wsgi_req->path_info_len = ptr - base;
if (watermark + (ptr - base) < (char *)(wsgi_req->proto_parser_buf + uwsgi.buffer_size)) {
wsgi_req->path_info = watermark;
wsgi_req->path_info_len = ptr - base;
http_url_decode(base, &wsgi_req->path_info_len, wsgi_req->path_info);
wsgi_req->uh.pktsize += proto_base_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, wsgi_req->path_info, wsgi_req->path_info_len);
}
else {
uwsgi_log("not enough space in wsgi_req http proto_parser_buf to encode PATH_INFO, consider tuning it with --buffer-size\n");
wsgi_req->uh.pktsize += proto_base_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, base, ptr - base);
wsgi_req->path_info = (wsgi_req->buffer + wsgi_req->uh.pktsize) - (ptr - base);
wsgi_req->path_info_len = ptr - base;
}
wsgi_req->path_info_pos = 5;
wsgi_req->uh.pktsize += proto_base_add_uwsgi_var(wsgi_req, "QUERY_STRING", 12, "", 0);
}
+31 -1
View File
@@ -30,6 +30,13 @@ int uwsgi_signal_handler(uint8_t sig) {
set_mule_harakiri(uwsgi.shared->options[UWSGI_OPTION_MULE_HARAKIRI]);
}
}
#ifdef UWSGI_SPOOLER
else if (uwsgi.spool_dir && uwsgi.shared->spooler_pid > 0 && (getpid() == uwsgi.shared->spooler_pid)) {
if(uwsgi.shared->options[UWSGI_OPTION_SPOOLER_HARAKIRI] > 0) {
set_spooler_harakiri(uwsgi.shared->options[UWSGI_OPTION_SPOOLER_HARAKIRI]);
}
}
#endif
int ret = uwsgi.p[use->modifier1]->signal_handler(sig, use->handler);
@@ -39,12 +46,19 @@ int uwsgi_signal_handler(uint8_t sig) {
set_harakiri(0);
}
}
if (uwsgi.muleid > 0) {
else if (uwsgi.muleid > 0) {
uwsgi.mules[uwsgi.muleid-1].sig = 0;
if(uwsgi.mules[uwsgi.muleid-1].harakiri > 0) {
set_mule_harakiri(0);
}
}
#ifdef UWSGI_SPOOLER
else if (uwsgi.spool_dir && uwsgi.shared->spooler_pid > 0 && (getpid() == uwsgi.shared->spooler_pid)) {
if(uwsgi.shared->options[UWSGI_OPTION_SPOOLER_HARAKIRI] > 0) {
set_spooler_harakiri(0);
}
}
#endif
return ret;
}
@@ -221,6 +235,7 @@ void uwsgi_route_signal(uint8_t sig) {
else if (!strcmp(use->receiver, "subscribed")) {
}
// route to spooler
#ifdef UWSGI_SPOOLER
else if (!strcmp(use->receiver, "spooler")) {
if (ushared->worker_signal_pipe[0] != -1) {
if (write(ushared->spooler_signal_pipe[0], &sig, 1) != 1) {
@@ -229,6 +244,7 @@ void uwsgi_route_signal(uint8_t sig) {
}
}
}
#endif
else if (!strcmp(use->receiver, "mules")) {
for(i=0;i<uwsgi.mules_cnt;i++) {
if (write(uwsgi.mules[i].signal_pipe[0], &sig, 1) != 1) {
@@ -255,6 +271,20 @@ void uwsgi_route_signal(uint8_t sig) {
}
}
}
else if (!strncmp(use->receiver, "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);
+7 -1
View File
@@ -16,7 +16,7 @@ pid_t spooler_start() {
int i;
pid_t pid = fork();
pid_t pid = uwsgi_fork("uWSGI spooler");
if (pid < 0) {
uwsgi_error("fork()");
exit(1);
@@ -413,7 +413,13 @@ void spooler_manage_task(char *dir, char *task) {
for(i=0;i<0xff;i++) {
if (uwsgi.p[i]->spooler) {
time_t now = time(NULL);
if(uwsgi.shared->options[UWSGI_OPTION_SPOOLER_HARAKIRI] > 0) {
set_spooler_harakiri(uwsgi.shared->options[UWSGI_OPTION_SPOOLER_HARAKIRI]);
}
ret = uwsgi.p[i]->spooler(task, spool_buf, uh.pktsize, body, body_len);
if(uwsgi.shared->options[UWSGI_OPTION_SPOOLER_HARAKIRI] > 0) {
set_spooler_harakiri(0);
}
if (body) {
free(body);
}
+220 -20
View File
@@ -2,6 +2,7 @@
extern struct uwsgi_server uwsgi;
extern char **environ;
/* statistically ordered */
struct http_status_codes hsc[] = {
@@ -121,6 +122,20 @@ void set_mule_harakiri(int sec) {
}
}
#ifdef UWSGI_SPOOLER
void set_spooler_harakiri(int sec) {
if (sec == 0) {
uwsgi.shared->spooler_harakiri = 0;
}
else {
uwsgi.shared->spooler_harakiri = time(NULL) + sec;
}
if (!uwsgi.master_process) {
alarm(sec);
}
}
#endif
void daemonize(char *logfile) {
pid_t pid;
@@ -738,7 +753,7 @@ void uwsgi_close_request(struct wsgi_request *wsgi_req) {
int waitpid_status;
int tmp_id;
uint64_t tmp_rt;
uint64_t tmp_rt, rss = 0, vsz = 0;
gettimeofday(&wsgi_req->end_of_request, NULL);
@@ -748,8 +763,11 @@ void uwsgi_close_request(struct wsgi_request *wsgi_req) {
uwsgi.workers[uwsgi.mywid].avg_response_time = (uwsgi.workers[uwsgi.mywid].avg_response_time+tmp_rt)/2;
// get memory usage
if (uwsgi.shared->options[UWSGI_OPTION_MEMORY_DEBUG] == 1 || uwsgi.force_get_memusage )
get_memusage();
if (uwsgi.shared->options[UWSGI_OPTION_MEMORY_DEBUG] == 1 || uwsgi.force_get_memusage ) {
get_memusage(&rss, &vsz);
uwsgi.workers[uwsgi.mywid].vsz_size = vsz;
uwsgi.workers[uwsgi.mywid].rss_size = rss;
}
// close the connection with the webserver
@@ -798,11 +816,11 @@ void uwsgi_close_request(struct wsgi_request *wsgi_req) {
goodbye_cruel_world();
}
if (uwsgi.reload_on_as && (rlim_t) uwsgi.workers[uwsgi.mywid].vsz_size >= uwsgi.reload_on_as) {
if (uwsgi.reload_on_as && (rlim_t) vsz >= uwsgi.reload_on_as) {
goodbye_cruel_world();
}
if (uwsgi.reload_on_rss && (rlim_t) uwsgi.workers[uwsgi.mywid].rss_size >= uwsgi.reload_on_rss) {
if (uwsgi.reload_on_rss && (rlim_t) rss >= uwsgi.reload_on_rss) {
goodbye_cruel_world();
}
@@ -1697,24 +1715,45 @@ int uwsgi_read_whole_body(struct wsgi_request *wsgi_req, char *buf, size_t len)
return 0;
}
void add_exported_option(char *key, char *value, int configured) {
void _add_exported_option(char *key, char *value, int configured) {
if (!uwsgi.exported_opts) {
uwsgi.exported_opts = uwsgi_malloc(sizeof(struct uwsgi_opt *));
}
else {
uwsgi.exported_opts = realloc(uwsgi.exported_opts, sizeof(struct uwsgi_opt *) * (uwsgi.exported_opts_cnt + 1));
if (!uwsgi.exported_opts) {
uwsgi_error("realloc()");
exit(1);
}
uwsgi.exported_opts = uwsgi_malloc(sizeof(struct uwsgi_opt *));
}
else {
uwsgi.exported_opts = realloc(uwsgi.exported_opts, sizeof(struct uwsgi_opt *) * (uwsgi.exported_opts_cnt + 1));
if (!uwsgi.exported_opts) {
uwsgi_error("realloc()");
exit(1);
}
}
uwsgi.exported_opts[uwsgi.exported_opts_cnt] = uwsgi_malloc(sizeof(struct uwsgi_opt));
uwsgi.exported_opts[uwsgi.exported_opts_cnt]->key = key;
uwsgi.exported_opts[uwsgi.exported_opts_cnt]->value = value;
uwsgi.exported_opts[uwsgi.exported_opts_cnt]->configured = configured;
uwsgi.exported_opts_cnt++;
}
void add_exported_option(char *key, char *value, int configured) {
char *v = value;
if (value && strchr(value, ';')) {
v = uwsgi_str(value);
}
uwsgi.exported_opts[uwsgi.exported_opts_cnt] = uwsgi_malloc(sizeof(struct uwsgi_opt));
uwsgi.exported_opts[uwsgi.exported_opts_cnt]->key = key;
uwsgi.exported_opts[uwsgi.exported_opts_cnt]->value = value;
uwsgi.exported_opts[uwsgi.exported_opts_cnt]->configured = configured;
uwsgi.exported_opts_cnt++;
if (v == NULL) {
_add_exported_option(key, v, configured);
return;
}
char *p = strtok(v, ";");
while(p != NULL) {
_add_exported_option(key, p, configured);
p = strtok(NULL, ";");
}
}
@@ -2225,8 +2264,13 @@ void spawn_daemon(struct uwsgi_daemon *ud) {
char *a;
int cnt = 1;
int devnull = -1;
int throttle = 0;
pid_t pid = fork();
if (uwsgi.current_time - ud->last_spawn <= 3) {
throttle = ud->respawns - (uwsgi.current_time-ud->last_spawn);
}
pid_t pid = uwsgi_fork("uWSGI external daemon");
if (pid < 0) {
uwsgi_error("fork()");
return;
@@ -2289,6 +2333,12 @@ void spawn_daemon(struct uwsgi_daemon *ud) {
argv[cnt] = NULL;
if (throttle) {
uwsgi_log_verbose("throttling %s for %d seconds\n", argv[0], throttle);
sleep(throttle);
}
uwsgi_log_verbose("running %s\n", argv[0]);
if (execvp(argv[0], argv)) {
uwsgi_error("execvp()");
}
@@ -2417,6 +2467,37 @@ int uwsgi_list_has_str(char *list, char *str) {
return 0;
}
char hex2num(char *str) {
char val = 0;
val <<= 4;
if (str[0] >= '0' && str[0] <= '9') {
val += str[0] & 0x0F;
}
else if ( str[0] >= 'A' && str[0] <= 'F') {
val += (str[0] & 0x0F) + 9;
}
else {
return 0;
}
val <<= 4;
if (str[1] >= '0' && str[1] <= '9') {
val += str[1] & 0x0F;
}
else if ( str[1] >= 'A' && str[1] <= 'F') {
val += (str[1] & 0x0F) + 9;
}
else {
return 0;
}
return val;
}
int uwsgi_str2_num(char *str) {
int num = 0;
@@ -3187,3 +3268,122 @@ void uwsgi_apply_config_pass(char symbol, char*(*hook)(char *) ) {
}
}
void uwsgi_set_processname(char *name) {
#if defined(__linux__) || defined(__sun__)
size_t amount = 0;
// prepare for strncat
*uwsgi.orig_argv[0] = 0;
if (uwsgi.procname_prefix) {
amount += strlen(uwsgi.procname_prefix);
if ((int)amount > uwsgi.max_procname-1) return;
strncat(uwsgi.orig_argv[0], uwsgi.procname_prefix, uwsgi.max_procname-(amount+1));
}
amount += strlen(name);
if ((int)amount > uwsgi.max_procname-1) return;
strncat(uwsgi.orig_argv[0], name, (uwsgi.max_procname-amount+1));
if (uwsgi.procname_append) {
amount += strlen(uwsgi.procname_append);
if ((int)amount > uwsgi.max_procname-1) return;
strncat(uwsgi.orig_argv[0], uwsgi.procname_append, uwsgi.max_procname-(amount+1));
}
memset(uwsgi.orig_argv[0]+amount+1, ' ', uwsgi.max_procname-(amount-1));
#elif defined(__FreeBSD__)
if (uwsgi.procname_prefix) {
if (!uwsgi.procname_append) {
setproctitle("-%s%s", uwsgi.procname_prefix, name);
}
else {
setproctitle("-%s%s%s", uwsgi.procname_prefix, name, uwsgi.procname_append);
}
}
else if (uwsgi.procname_append) {
if (!uwsgi.procname_prefix) {
setproctitle("-%s%s", name, uwsgi.procname_append);
}
else {
setproctitle("-%s%s%s", uwsgi.procname_prefix, name, uwsgi.procname_append);
}
}
else {
setproctitle("-%s", name);
}
#endif
}
// this is a wrapper for fork restoring original argv
pid_t uwsgi_fork(char *name) {
pid_t pid = fork();
if (pid == 0) {
#if defined(__linux__) || defined(__sun__)
int i;
for(i=0;i<uwsgi.argc;i++) {
strcpy(uwsgi.orig_argv[i],uwsgi.argv[i]);
}
#endif
if (uwsgi.auto_procname && name) {
if (uwsgi.procname) {
uwsgi_set_processname(uwsgi.procname);
}
else {
uwsgi_set_processname(name);
}
}
}
return pid;
}
void http_url_decode(char *buf, uint16_t *len, char *dst) {
uint16_t i;
int percent = 0;
char value[2];
size_t new_len = 0;
char *ptr = dst;
value[0] = '0';
value[1] = '0';
for(i=0;i<*len;i++) {
if (buf[i] == '%') {
if (percent == 0) {
percent = 1;
}
else {
*ptr++= '%';
new_len++;
percent = 0;
}
}
else {
if (percent == 1) {
value[0] = buf[i];
percent = 2;
}
else if (percent == 2) {
value[1] = buf[i];
*ptr++= hex2num(value);
percent = 0;
new_len++;
}
else {
*ptr++= buf[i];
new_len++;
}
}
}
*len = new_len;
}
+207 -40
View File
@@ -78,6 +78,11 @@ static struct option long_base_options[] = {
#ifdef UWSGI_THREADING
{"enable-threads", no_argument, 0, 'T'},
#endif
{"auto-procname", no_argument, &uwsgi.auto_procname, 1},
{"procname-prefix", required_argument, 0, LONG_ARGS_PROCNAME_PREFIX},
{"procname-append", required_argument, 0, LONG_ARGS_PROCNAME_APPEND},
{"procname", required_argument, 0, LONG_ARGS_PROCNAME},
{"procname-master", required_argument, 0, LONG_ARGS_PROCNAME_MASTER},
{"single-interpreter", no_argument, 0, 'i'},
{"master", no_argument, 0, 'M'},
{"emperor", required_argument, 0, LONG_ARGS_EMPEROR},
@@ -121,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},
@@ -510,7 +516,7 @@ void kill_them_all(int signum) {
#ifdef UWSGI_SPOOLER
if (uwsgi.spool_dir && uwsgi.shared->spooler_pid > 0) {
kill(uwsgi.shared->spooler_pid, SIGKILL);
uwsgi_log("killed the spooler with pid %d\n", uwsgi.shared->spooler_pid);
uwsgi_log("killing the spooler with pid %d\n", uwsgi.shared->spooler_pid);
}
#endif
@@ -518,7 +524,7 @@ void kill_them_all(int signum) {
if (uwsgi.emperor_pid >= 0) {
kill(uwsgi.emperor_pid, SIGKILL);
waitpid(uwsgi.emperor_pid, &i, 0);
uwsgi_log("killed the emperor with pid %d\n", uwsgi.emperor_pid);
uwsgi_log("killing the emperor with pid %d\n", uwsgi.emperor_pid);
}
@@ -560,14 +566,14 @@ void grace_them_all(int signum) {
#ifdef UWSGI_SPOOLER
if (uwsgi.spool_dir && uwsgi.shared->spooler_pid > 0) {
kill(uwsgi.shared->spooler_pid, SIGKILL);
uwsgi_log("killed the spooler with pid %d\n", uwsgi.shared->spooler_pid);
uwsgi_log("killing the spooler with pid %d\n", uwsgi.shared->spooler_pid);
}
#endif
if (uwsgi.emperor_pid >= 0) {
kill(uwsgi.emperor_pid, SIGKILL);
waitpid(uwsgi.emperor_pid, &i, 0);
uwsgi_log("killed the emperor with pid %d\n", uwsgi.emperor_pid);
uwsgi_log("killing the emperor with pid %d\n", uwsgi.emperor_pid);
}
for (i = 0; i < uwsgi.shared->daemons_cnt; i++) {
@@ -653,6 +659,7 @@ void reap_them_all(int signum) {
}
for (i = 0; i < uwsgi.mules_cnt; i++) {
if (!uwsgi.mules) break;
if (uwsgi.mules[i].pid > 0)
kill(uwsgi.mules[i].pid, SIGKILL);
}
@@ -660,7 +667,7 @@ void reap_them_all(int signum) {
if (uwsgi.emperor_pid >= 0) {
kill(uwsgi.emperor_pid, SIGKILL);
waitpid(uwsgi.emperor_pid, &i, 0);
uwsgi_log("killed the emperor with pid %d\n", uwsgi.emperor_pid);
uwsgi_log("killing the emperor with pid %d\n", uwsgi.emperor_pid);
}
if (!uwsgi.workers)
@@ -847,6 +854,56 @@ void signal_pidfile(int sig, char *filename) {
}
}
void fixup_argv_and_environ(int argc, char **argv, char **environ) {
uwsgi.orig_argv = argv;
uwsgi.argv = argv;
uwsgi.argc = argc;
uwsgi.environ = environ;
#if defined(__linux__) || defined(__sun__)
int i;
int env_count = 0;
uwsgi.argv = uwsgi_malloc( sizeof(char *) * argc);
for(i=0;i<argc;i++) {
if (i==0 || argv[0] + uwsgi.max_procname + 1 == argv[i]) {
uwsgi.max_procname += strlen(argv[i]) + 1;
}
uwsgi.argv[i] = strdup(argv[i]);
}
uwsgi.max_procname++;
for (i = 0; environ[i] != NULL; i++) {
// useless
//if ((environ[0] + uwsgi.max_procname + 1) == environ[i]) {
uwsgi.max_procname += strlen(environ[i]) + 1;
//}
env_count++;
}
uwsgi.environ = uwsgi_malloc(sizeof(char *) * env_count);
for(i=0;i<env_count;i++) {
uwsgi.environ[i] = strdup(environ[i]);
#ifdef UWSGI_DEBUG
uwsgi_log("ENVIRON: %s\n", uwsgi.environ[i]);
#endif
environ[i] = uwsgi.environ[i];
}
#ifdef UWSGI_DEBUG
uwsgi_log("max space for custom process name = %d\n", uwsgi.max_procname);
#endif
//environ = uwsgi.environ;
#endif
}
#ifdef UWSGI_AS_SHARED_LIBRARY
int uwsgi_init(int argc, char *argv[], char *envp[]) {
@@ -886,6 +943,7 @@ int main(int argc, char *argv[], char *envp[]) {
signal(SIGTERM, SIG_IGN);
signal(SIGPIPE, SIG_IGN);
//initialize masterpid with a default value
masterpid = getpid();
@@ -1024,6 +1082,9 @@ int main(int argc, char *argv[], char *envp[]) {
uwsgi.page_size = getpagesize();
uwsgi.binary_path = uwsgi_get_binary_path(argv[0]);
// ok we can now safely play with argv and environ
fixup_argv_and_environ(argc, argv, environ);
//initialize embedded plugins
UWSGI_LOAD_EMBEDDED_PLUGINS
// now a bit of magic, if the executable basename contains a 'uwsgi_' string,
@@ -1067,7 +1128,7 @@ int main(int argc, char *argv[], char *envp[]) {
#endif
uwsgi.option_index = -1;
while ((i = getopt_long(argc, argv, short_options, uwsgi.long_options, &uwsgi.option_index)) != -1) {
while ((i = getopt_long(uwsgi.argc, uwsgi.argv, short_options, uwsgi.long_options, &uwsgi.option_index)) != -1) {
if (uwsgi.option_index > -1) {
optname = (char *) uwsgi.long_options[uwsgi.option_index].name;
@@ -1091,8 +1152,8 @@ int main(int argc, char *argv[], char *envp[]) {
#endif
if (optind < argc) {
for(i=optind;i<argc;i++) {
char *lazy = argv[i];
for(i=optind;i<uwsgi.argc;i++) {
char *lazy = uwsgi.argv[i];
if (lazy[0] != '[') {
if (0) {
}
@@ -1159,7 +1220,6 @@ int main(int argc, char *argv[], char *envp[]) {
}
}
struct uwsgi_string_list *config_file;
#ifdef UWSGI_XML
if (uwsgi.xml_config != NULL) {
@@ -1169,6 +1229,7 @@ int main(int argc, char *argv[], char *envp[]) {
}
#endif
#ifdef UWSGI_INI
struct uwsgi_string_list *config_file;
if (uwsgi.ini != NULL) {
config_file = uwsgi.ini;
while(config_file) {
@@ -1337,8 +1398,8 @@ int main(int argc, char *argv[], char *envp[]) {
//call after_opt hooks
if (uwsgi.binary_path == argv[0]) {
uwsgi.binary_path = uwsgi_str(argv[0]);
if (uwsgi.binary_path == uwsgi.argv[0]) {
uwsgi.binary_path = uwsgi_str(uwsgi.argv[0]);
}
if (!uwsgi.no_initial_output) {
@@ -1433,7 +1494,7 @@ int main(int argc, char *argv[], char *envp[]) {
exit(1);
}
uwsgi.emperor_pid = fork();
uwsgi.emperor_pid = uwsgi_fork("uWSGI Emperor");
if (uwsgi.emperor_pid < 0) {
uwsgi_error("pid()");
exit(1);
@@ -1454,7 +1515,7 @@ int main(int argc, char *argv[], char *envp[]) {
// call jail systems
for (i = 0; i < uwsgi.gp_cnt; i++) {
if (uwsgi.gp[i]->jail) {
uwsgi.gp[i]->jail(uwsgi_start, argv);
uwsgi.gp[i]->jail(uwsgi_start, uwsgi.argv);
}
}
@@ -1462,12 +1523,12 @@ int main(int argc, char *argv[], char *envp[]) {
// TODO pluginize basic Linux namespace support
#ifdef __linux__
if (uwsgi.ns) {
linux_namespace_start((void *) argv);
linux_namespace_start((void *) uwsgi.argv);
// never here
}
else {
#endif
uwsgi_start((void *) argv);
uwsgi_start((void *) uwsgi.argv);
#ifdef __linux__
}
#endif
@@ -1479,8 +1540,6 @@ int main(int argc, char *argv[], char *envp[]) {
int uwsgi_start(void *v_argv) {
char **argv = v_argv;
#ifdef UWSGI_DEBUG
int so_bufsize;
socklen_t so_bufsize_len;
@@ -1590,7 +1649,7 @@ int uwsgi_start(void *v_argv) {
exit(1);
}
uwsgi.emperor_pid = fork();
uwsgi.emperor_pid = uwsgi_fork("uWSGI Emperor");
if (uwsgi.emperor_pid < 0) {
uwsgi_error("pid()");
exit(1);
@@ -2154,6 +2213,7 @@ skipzero:
for(i=1;i<=uwsgi.numproc;i++) {
uwsgi.workers[i].signal_pipe[0] = - 1;
uwsgi.workers[i].signal_pipe[1] = - 1;
snprintf(uwsgi.workers[i].name, 0xff, "uWSGI worker %d", i);
}
if (uwsgi.master_process) {
@@ -2193,9 +2253,68 @@ skipzero:
uwsgi_error("socketpair()");
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;i<uwsgi.farms_cnt;i++) {
char *farm_value = uwsgi_str(farm_name->value);
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;
@@ -2370,15 +2489,14 @@ skipzero:
}
}
if (getpid() == masterpid && uwsgi.master_process == 1) {
#ifdef UWSGI_AS_SHARED_LIBRARY
int ml_ret = master_loop(argv, environ);
int ml_ret = master_loop(uwsgi.argv, uwsgi.environ);
if (ml_ret == -1) {
return 0;
}
#else
(void) master_loop(argv, environ);
(void) master_loop(uwsgi.argv, uwsgi.environ);
#endif
//from now on the process is a real worker
}
@@ -2799,6 +2917,22 @@ static int manage_base_opt(int i, char *optarg) {
case 0:
return 1;
case LONG_ARGS_PROCNAME_PREFIX:
uwsgi.auto_procname = 1;
uwsgi.procname_prefix = optarg;
return 1;
case LONG_ARGS_PROCNAME_APPEND:
uwsgi.auto_procname = 1;
uwsgi.procname_append = optarg;
return 1;
case LONG_ARGS_PROCNAME:
uwsgi.auto_procname = 1;
uwsgi.procname = optarg;
return 1;
case LONG_ARGS_PROCNAME_MASTER:
uwsgi.auto_procname = 1;
uwsgi.procname_master = optarg;
return 1;
#ifdef UWSGI_UDP
case LONG_ARGS_CLUSTER_RELOAD:
send_udp_message(98, optarg, "", 0);
@@ -3301,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;
@@ -3521,12 +3660,14 @@ static int manage_base_opt(int i, char *optarg) {
case 't':
uwsgi.shared->options[UWSGI_OPTION_HARAKIRI] = atoi(optarg);
return 1;
#ifdef UWSGI_SPOOLER
case LONG_ARGS_SPOOLER_HARAKIRI:
uwsgi.shared->options[UWSGI_OPTION_SPOOLER_HARAKIRI] = atoi(optarg);
return 1;
case LONG_ARGS_SPOOLER_CHDIR:
uwsgi.spooler_chdir = optarg;
return 1;
#endif
case LONG_ARGS_MULE_HARAKIRI:
uwsgi.shared->options[UWSGI_OPTION_MULE_HARAKIRI] = atoi(optarg);
return 1;
@@ -3601,32 +3742,58 @@ static int manage_base_opt(int i, char *optarg) {
return 0;
}
void manage_opt(int i, char *optarg) {
int _manage_opt(int i, char *p) {
int j;
if (manage_base_opt(i, optarg)) {
if (manage_base_opt(i, p)) {
return 1;
}
for (j = 0; j < 0xFF; j++) {
if (uwsgi.p[j]->manage_opt) {
if (uwsgi.p[j]->manage_opt(i, p)) {
return 1;
}
}
}
for (j = 0; j < uwsgi.gp_cnt; j++) {
if (uwsgi.gp[j]->manage_opt) {
if (uwsgi.gp[j]->manage_opt(i, p)) {
return 1;
}
}
}
return 0;
}
void manage_opt(int i, char *optarg) {
char *value = optarg;
if (value && strchr(optarg, ';')) {
value = uwsgi_str(optarg);
}
if (!value) {
if (!_manage_opt(i, NULL)) {
exit(1);
}
return;
}
for (j = 0; j < 0xFF; j++) {
if (uwsgi.p[j]->manage_opt) {
if (uwsgi.p[j]->manage_opt(i, optarg)) {
return;
}
}
char *p = strtok(value, ";");
while(p != NULL) {
if (_manage_opt(i, p)) goto next;
exit(1);
next:
p = strtok(NULL, ";");
}
for (j = 0; j < uwsgi.gp_cnt; j++) {
if (uwsgi.gp[j]->manage_opt) {
if (uwsgi.gp[j]->manage_opt(i, optarg)) {
return;
}
}
}
//never here
exit(1);
return;
}
+52 -1
View File
@@ -528,6 +528,11 @@ struct uwsgi_opt {
#define LONG_ARGS_EMPEROR_STATS 17150
#define LONG_ARGS_SPOOLER_CHDIR 17151
#define LONG_ARGS_LOCKS 17152
#define LONG_ARGS_PROCNAME_PREFIX 17153
#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
@@ -966,6 +971,17 @@ struct uwsgi_server {
char hostname[256];
int hostname_len;
char **orig_argv;
char **argv;
int argc;
int max_procname;
int auto_procname;
char **environ;
char *procname_prefix;
char *procname_append;
char *procname_master;
char *procname;
// quiet startup
int no_initial_output;
@@ -1279,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;
@@ -1594,6 +1613,7 @@ struct uwsgi_shared {
#ifdef UWSGI_SPOOLER
pid_t spooler_pid;
int spooler_frequency;
time_t spooler_harakiri;
#endif
#ifdef UWSGI_SNMP
@@ -1712,8 +1732,10 @@ struct uwsgi_worker {
uint64_t avg_response_time;
char name[0xff];
};
struct uwsgi_mule {
int id;
pid_t pid;
@@ -1736,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);
@@ -1763,7 +1802,7 @@ void daemonize(char *);
void logto(char *);
void log_request(struct wsgi_request *);
void get_memusage(void);
void get_memusage(uint64_t *, uint64_t *);
void harakiri(void);
void stats(int);
@@ -1787,6 +1826,7 @@ pid_t spooler_start(void);
void set_harakiri(int);
void set_mule_harakiri(int);
void set_spooler_harakiri(int);
void inc_harakiri(int);
#ifdef __BIG_ENDIAN__
@@ -2371,6 +2411,17 @@ char *uwsgi_string_get_list(struct uwsgi_string_list **, int, size_t *);
void uwsgi_fixup_fds(int, int);
void uwsgi_set_processname(char *);
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 *);
int uwsgi_farm_has_mule(struct uwsgi_farm *, int);
#ifdef UWSGI_CAP
void uwsgi_build_cap(char *);
#endif
+24
View File
@@ -84,6 +84,30 @@ class rpc(object):
uwsgi.register_rpc(self.name, f)
return f
class farm_loop(object):
def __init__(self, f, farm):
self.f = f
self.farm = farm
def __call__(self):
if uwsgi.mule_id() == 0:
return
if not uwsgi.in_farm(self.farm):
return
while True:
message = uwsgi.farm_get_msg()
if message:
self.f(message)
class farm(object):
def __init__(self, name=None, **kwargs):
self.name = name
def __call__(self, f):
postfork_chain.append(farm_loop(f, self.name))
class signal(object):
def __init__(self, num, **kwargs):
+10
View File
@@ -2,6 +2,7 @@ import uwsgi
import os
import gc
import sys
from uwsgidecorators import *
gc.set_debug(gc.DEBUG_SAVEALL)
print os.environ
@@ -42,10 +43,19 @@ routes['/logo'] = serve_logo
routes['/config'] = serve_config
routes['/options'] = serve_options
@postfork
def setprocname():
if uwsgi.worker_id() > 0:
uwsgi.setprocname("i am the worker %d" % uwsgi.worker_id())
def application(env, start_response):
uwsgi.mule_msg(env['REQUEST_URI'], 1)
req = uwsgi.workers()[uwsgi.worker_id()-1]['requests']
uwsgi.setprocname("worker %d managed %d requests" % (uwsgi.worker_id(), req))
gc.collect(2)
if DEBUG:
print env['wsgi.input'].fileno()