mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-10-04 07:58:33 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2d731532be | ||
|
|
d545cd5be0 | ||
|
|
261d05d67e | ||
|
|
639995785a | ||
|
|
817e076cae | ||
|
|
713370f343 | ||
|
|
5594503a56 | ||
|
|
d57e5d46fb | ||
|
|
ee2db5fabd | ||
|
|
0c4cd6f378 | ||
|
|
8c32b9665c | ||
|
|
7405a01654 | ||
|
|
5d27691b4b | ||
|
|
19bbcfbf4c | ||
|
|
fa1064db87 | ||
|
|
850ac69ca1 | ||
|
|
943edeb30c | ||
|
|
01815a354b | ||
|
|
2f4412aca1 | ||
|
|
2782bb7594 | ||
|
|
074c1298d5 | ||
|
|
55a7911611 | ||
|
|
013ad3044d | ||
|
|
dcf4217b9d | ||
|
|
876dfc31a9 | ||
|
|
fd03bdbe00 | ||
|
|
83612649d5 | ||
|
|
3c4db7e97e | ||
|
|
088c70efdc | ||
|
|
2b19029dd0 | ||
|
|
4fb5a84378 | ||
|
|
b3107d9a25 | ||
|
|
bf32988f3e | ||
|
|
1bf3dd7bf4 | ||
|
|
07fd84f89f | ||
|
|
3a8db4a97a | ||
|
|
aea30a26f5 | ||
|
|
8f5d1cdae7 |
@@ -35,3 +35,4 @@ fb168b0b86169219aa9b8e40f0caa6297cf34dbc 0.9.9-rc1
|
||||
7048ae11cfc8453c6bb599af3c0dfacd7ebae43c 1.0-rc1
|
||||
7a4e021ea7ba2953e474326612c3bb82d6817f0e 1.0-rc2
|
||||
4e4e781ef89ac33e71a7e92741d54549240ed156 1.0-rc3
|
||||
24f7e260ea345c2cd73e2fd8191dddbd1b039062 1.0-rc4
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import sys
|
||||
import uwsgi
|
||||
|
||||
print("i am the bootstrap for uwsgi.SymbolsImporter")
|
||||
sys.meta_path.insert(0, uwsgi.SymbolsImporter())
|
||||
|
||||
@@ -27,7 +27,7 @@ plugins =
|
||||
bin_name = uwsgi
|
||||
append_version =
|
||||
plugin_dir = .
|
||||
embedded_plugins = python, ping, cache, nagios, rpc, fastrouter, http, ugreen
|
||||
embedded_plugins = python, ping, cache, nagios, rrdtool, rpc, fastrouter, http, ugreen
|
||||
as_shared_library = false
|
||||
|
||||
locking = auto
|
||||
|
||||
@@ -94,157 +94,6 @@ void uwsgi_init_cache() {
|
||||
uwsgi_log("*** Cache subsystem initialized: %dMB preallocated ***\n", ((sizeof(uint64_t) * UMAX16) + (sizeof(uint64_t) * uwsgi.cache_max_items) + (uwsgi.cache_blocksize * uwsgi.cache_max_items) + (sizeof(struct uwsgi_cache_item) * uwsgi.cache_max_items)) / (1024 * 1024));
|
||||
}
|
||||
|
||||
struct uwsgi_subscriber_name *uwsgi_get_subscriber(struct uwsgi_dict *udict, char *key, uint16_t keylen) {
|
||||
|
||||
uint64_t ovl;
|
||||
struct uwsgi_subscriber *usub;
|
||||
struct uwsgi_subscriber_name *ret = NULL;
|
||||
|
||||
usub = (struct uwsgi_subscriber *) uwsgi_dict_get(udict, key, keylen, &ovl);
|
||||
|
||||
if (usub == NULL || !ovl) return NULL;
|
||||
|
||||
if (!usub->nodes) return NULL;
|
||||
|
||||
if (usub && ovl) {
|
||||
ret = &usub->names[usub->current];
|
||||
// dead node
|
||||
if (ret->len == 0) {
|
||||
if (usub->current == usub->nodes-1) {
|
||||
usub->nodes--;
|
||||
}
|
||||
// retry with another node (if available)
|
||||
if (usub->nodes > 0) {
|
||||
usub->current++;
|
||||
if (usub->current >= usub->nodes) usub->current = 0;
|
||||
return uwsgi_get_subscriber(udict, key, keylen);
|
||||
}
|
||||
}
|
||||
|
||||
if (usub->nodes > 1) {
|
||||
usub->current++;
|
||||
if (usub->current >= usub->nodes) usub->current = 0;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
void uwsgi_add_subscriber(struct uwsgi_dict *udict, struct uwsgi_subscribe_req *usr) {
|
||||
|
||||
char *ptr;
|
||||
uint64_t vallen = 0;
|
||||
struct uwsgi_subscriber *usub, nusub;
|
||||
int found = 0;
|
||||
int i;
|
||||
|
||||
ptr = uwsgi_dict_get(udict, usr->key, usr->keylen, &vallen);
|
||||
if (ptr && vallen) {
|
||||
usub = (struct uwsgi_subscriber *) ptr;
|
||||
for(i=0;i<(int)usub->nodes;i++) {
|
||||
if (!uwsgi_strncmp(usub->names[i].name, usub->names[i].len, usr->address, usr->address_len)) {
|
||||
found = 1;
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (!found) {
|
||||
found = usub->nodes;
|
||||
// check for unallocated slot
|
||||
for(i=0;i<(int)usub->nodes;i++) {
|
||||
if (usub->names[i].len == 0) {
|
||||
found = i;
|
||||
break;
|
||||
}
|
||||
}
|
||||
usub->names[found].len = usr->address_len;
|
||||
usub->names[found].modifier1 = usr->modifier1;
|
||||
usub->names[found].modifier2 = usr->modifier2;
|
||||
memcpy(usub->names[found].name, usr->address, usr->address_len);
|
||||
if (found == (int) usub->nodes) {
|
||||
usub->nodes++;
|
||||
}
|
||||
uwsgi_log("[uwsgi-subscription] %.*s => new node: %.*s\n", usr->keylen, usr->key, usr->address_len, usr->address);
|
||||
udict->count++;
|
||||
}
|
||||
return;
|
||||
}
|
||||
else {
|
||||
nusub.nodes = 1;
|
||||
nusub.current = 0;
|
||||
memcpy(nusub.names[0].name, usr->address, usr->address_len);
|
||||
nusub.names[0].len = usr->address_len;
|
||||
nusub.names[0].modifier1 = usr->modifier1;
|
||||
nusub.names[0].modifier2 = usr->modifier2;
|
||||
uwsgi_dict_set(udict, usr->key, usr->keylen, (char *) &nusub, sizeof(struct uwsgi_subscriber));
|
||||
uwsgi_log("[uwsgi-subscription] new pool: %.*s\n", usr->keylen, usr->key);
|
||||
uwsgi_log("[uwsgi-subscription] %.*s => new node: %.*s\n", usr->keylen, usr->key, usr->address_len, usr->address);
|
||||
udict->count++;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
struct uwsgi_dict *uwsgi_dict_create(uint64_t items, uint64_t blocksize) {
|
||||
|
||||
int i;
|
||||
|
||||
struct uwsgi_dict *udict = (struct uwsgi_dict *) mmap(NULL, sizeof(uint64_t) * UMAX16, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANON, -1, 0);
|
||||
if (!udict) {
|
||||
uwsgi_error("mmap()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
if (!blocksize) blocksize = 4096;
|
||||
|
||||
if (blocksize % uwsgi.page_size != 0) {
|
||||
uwsgi_log("invalid shared dictionary blocksize %llu: must be a multiple of memory page size (%d bytes)\n", (unsigned long long) udict->blocksize, uwsgi.page_size);
|
||||
exit(1);
|
||||
}
|
||||
|
||||
udict->blocksize = blocksize;
|
||||
udict->max_items = items;
|
||||
|
||||
udict->hashtable = (uint64_t *) mmap(NULL, sizeof(uint64_t) * UMAX16, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANON, -1, 0);
|
||||
if (!udict->hashtable) {
|
||||
uwsgi_error("mmap()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
memset(udict->hashtable, 0, sizeof(uint64_t) * UMAX16);
|
||||
|
||||
udict->unused_stack = (uint64_t *) mmap(NULL, sizeof(uint64_t) * udict->max_items, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANON, -1, 0);
|
||||
if (!udict->unused_stack) {
|
||||
uwsgi_error("mmap()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
memset(udict->unused_stack, 0, sizeof(uint64_t) * udict->max_items);
|
||||
|
||||
udict->items = (struct uwsgi_dict_item *) mmap(NULL, sizeof(struct uwsgi_dict_item) * udict->max_items, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANON, -1, 0);
|
||||
if (!udict->items) {
|
||||
uwsgi_error("mmap()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
udict->data = mmap(NULL, udict->blocksize * udict->max_items, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_ANON, -1, 0);
|
||||
if (!udict->data) {
|
||||
uwsgi_error("mmap()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
for(i=0;i< (int) udict->max_items;i++) {
|
||||
memset(&udict->items[i], 0, sizeof(struct uwsgi_dict_item));
|
||||
}
|
||||
|
||||
udict->first_available_item = 1;
|
||||
udict->unused_stack_ptr = 0;
|
||||
|
||||
udict->lock = uwsgi_mmap_shared_lock();
|
||||
uwsgi_lock_init(udict->lock);
|
||||
|
||||
return udict;
|
||||
}
|
||||
|
||||
uint32_t djb33x_hash(char *key, int keylen) {
|
||||
|
||||
register uint32_t hash = 5381;
|
||||
@@ -258,89 +107,6 @@ uint32_t djb33x_hash(char *key, int keylen) {
|
||||
}
|
||||
|
||||
|
||||
inline uint64_t uwsgi_dict_get_index(struct uwsgi_dict *udict, char *key, uint16_t keylen) {
|
||||
|
||||
uint32_t hash = djb33x_hash(key, keylen);
|
||||
|
||||
int hash_key = hash % 0xffff;
|
||||
|
||||
uint64_t slot = udict->hashtable[hash_key];
|
||||
|
||||
struct uwsgi_dict_item *udi;
|
||||
|
||||
udi = &udict->items[slot];
|
||||
|
||||
// first round
|
||||
if (udi->djbhash != hash) goto cycle;
|
||||
if (udi->keysize != keylen) goto cycle;
|
||||
if (memcmp(udi->key, key, keylen)) goto cycle;
|
||||
|
||||
return slot;
|
||||
|
||||
cycle:
|
||||
while(udi->next) {
|
||||
slot = udi->next;
|
||||
udi = &udict->items[slot];
|
||||
if (udi->djbhash != hash) continue;
|
||||
if (udi->keysize != keylen) continue;
|
||||
if (!memcmp(udi->key, key, keylen)) return slot;
|
||||
}
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
char *uwsgi_dict_get(struct uwsgi_dict *udict, char *key, uint16_t keylen, uint64_t *valsize) {
|
||||
|
||||
uint64_t index = uwsgi_dict_get_index(udict, key, keylen);
|
||||
|
||||
if (index) {
|
||||
*valsize = udict->items[index].valsize;
|
||||
udict->items[index].hits++;
|
||||
return udict->data+(index*udict->blocksize);
|
||||
}
|
||||
|
||||
return NULL;
|
||||
}
|
||||
|
||||
int uwsgi_dict_del(struct uwsgi_dict *udict, char *key, uint16_t keylen) {
|
||||
|
||||
uint64_t index = 0;
|
||||
struct uwsgi_dict_item *udi;
|
||||
int ret = -1;
|
||||
|
||||
index = uwsgi_dict_get_index(udict, key, keylen);
|
||||
if (index) {
|
||||
udi = &udict->items[index] ;
|
||||
udi->keysize = 0;
|
||||
udi->valsize = 0;
|
||||
udict->unused_stack_ptr++;
|
||||
udict->unused_stack[udict->unused_stack_ptr] = index;
|
||||
// try to return to initial condition...
|
||||
if (index == udict->first_available_item-1) {
|
||||
udict->first_available_item--;
|
||||
}
|
||||
ret = 0;
|
||||
// relink collisioned entry
|
||||
if (udi->prev) {
|
||||
udict->items[udi->prev].next = udi->next;
|
||||
}
|
||||
if (udi->next) {
|
||||
udict->items[udi->next].prev = udi->prev;
|
||||
}
|
||||
if (!udi->prev && !udi->next) {
|
||||
// reset hashtable entry
|
||||
udict->hashtable[udi->djbhash % 0xffff] = 0;
|
||||
}
|
||||
udi->djbhash = 0;
|
||||
udi->prev = 0;
|
||||
udi->next = 0;
|
||||
}
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
|
||||
|
||||
inline uint64_t uwsgi_cache_get_index(char *key, uint16_t keylen) {
|
||||
|
||||
uint32_t hash = djb33x_hash(key, keylen);
|
||||
@@ -534,71 +300,6 @@ end:
|
||||
|
||||
}
|
||||
|
||||
int uwsgi_dict_set(struct uwsgi_dict *udict, char *key, uint16_t keylen, char *val, uint64_t vallen) {
|
||||
|
||||
uint64_t index = 0, last_index = 0 ;
|
||||
|
||||
struct uwsgi_dict_item *udi, *udii;
|
||||
|
||||
int ret = -1;
|
||||
int slot;
|
||||
|
||||
if (!keylen || !vallen) return -1;
|
||||
|
||||
if (keylen > UWSGI_CACHE_MAX_KEY_SIZE) return -1;
|
||||
|
||||
if (udict->first_available_item >= udict->max_items && !udict->unused_stack_ptr) {
|
||||
uwsgi_log("*** DANGER dictionary %p is FULL !!! ***\n", udict);
|
||||
goto end;
|
||||
}
|
||||
|
||||
index = uwsgi_dict_get_index(udict, key, keylen);
|
||||
if (!index) {
|
||||
if (udict->unused_stack_ptr) {
|
||||
index = udict->unused_stack[udict->unused_stack_ptr];
|
||||
udict->unused_stack_ptr--;
|
||||
}
|
||||
else {
|
||||
index = udict->first_available_item;
|
||||
if (udict->first_available_item < udict->max_items) {
|
||||
udict->first_available_item++;
|
||||
}
|
||||
}
|
||||
udi = &udict->items[index] ;
|
||||
udi->djbhash = djb33x_hash(key, keylen);
|
||||
udi->hits = 0;
|
||||
memcpy(udi->key, key, keylen);
|
||||
memcpy(udict->data+(index*udict->blocksize), val, vallen);
|
||||
|
||||
// set this as late as possibile (to reduce races risk)
|
||||
|
||||
udi->valsize = vallen;
|
||||
udi->keysize = keylen;
|
||||
ret = 0;
|
||||
// now put the value in the 16bit hashtable
|
||||
slot = udi->djbhash % 0xffff;
|
||||
|
||||
if (udict->hashtable[slot] == 0) {
|
||||
udict->hashtable[slot] = index;
|
||||
}
|
||||
else {
|
||||
// append to first available next
|
||||
last_index = udict->hashtable[slot];
|
||||
udii = &udict->items[ last_index ];
|
||||
while(udii->next) {
|
||||
last_index = udii->next;
|
||||
udii = &udict->items[ last_index ];
|
||||
}
|
||||
udii->next = index;
|
||||
udi->prev = last_index;
|
||||
}
|
||||
}
|
||||
|
||||
end:
|
||||
return ret;
|
||||
|
||||
}
|
||||
|
||||
/* THIS PART IS HEAVILY OPTIMIZED: PERFORMANCE NOT ELEGANCE !!! */
|
||||
|
||||
void *cache_thread_loop(void *fd_ptr) {
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
#!/bin/bash
|
||||
|
||||
# uwsgi - Use uwsgi to run python and wsgi web apps.
|
||||
#
|
||||
# chkconfig: - 85 15
|
||||
# description: Use uwsgi to run python and wsgi web apps.
|
||||
# processname: uwsgi
|
||||
|
||||
# author: Roman Vasilyev
|
||||
|
||||
# Source function library.
|
||||
. /etc/rc.d/init.d/functions
|
||||
|
||||
PATH=/opt/uwsgi:/sbin:/bin:/usr/sbin:/usr/bin
|
||||
prog=/usr/sbin/uwsgi
|
||||
|
||||
OWNER=nginx
|
||||
|
||||
NAME=uwsgi
|
||||
DESC=uwsgi
|
||||
|
||||
#DAEMON_OPTS="-s 127.0.0.1:9001 -M 4 -t 30 -A 4 -p 4 -d /var/log/uwsgi.log --pidfile /var/run/$NAME.pid --pythonpath $PYTHONPATH --module $MODULE"
|
||||
DAEMON_OPTS="-s 127.0.0.1:9001 -M 4 -t 30 -A 4 -p 16 -b 32768 -d /var/log/$NAME.log --pidfile /var/run/$NAME.pid --uid $OWNER"
|
||||
|
||||
[ -f /etc/sysconfig/uwsgi ] && . /etc/sysconfig/uwsgi
|
||||
|
||||
lockfile=/var/lock/subsys/uwsgi
|
||||
|
||||
start () {
|
||||
echo -n "Starting $DESC: "
|
||||
daemon $prog $DAEMON_OPTS
|
||||
retval=$?
|
||||
echo
|
||||
[ $retval -eq 0 ] && touch $lockfile
|
||||
return $retval
|
||||
}
|
||||
|
||||
stop () {
|
||||
echo -n "Stopping $DESC: "
|
||||
killproc $prog
|
||||
retval=$?
|
||||
echo
|
||||
[ $retval -eq 0 ] && rm -f $lockfile
|
||||
return $retval
|
||||
}
|
||||
|
||||
reload () {
|
||||
echo "Reloading $NAME"
|
||||
killproc $prog -HUP
|
||||
RETVAL=$?
|
||||
echo
|
||||
}
|
||||
|
||||
force-reload () {
|
||||
echo "Reloading $NAME"
|
||||
killproc $prog -TERM
|
||||
RETVAL=$?
|
||||
echo
|
||||
}
|
||||
|
||||
restart () {
|
||||
stop
|
||||
start
|
||||
}
|
||||
|
||||
rh_status () {
|
||||
status $prog
|
||||
}
|
||||
|
||||
rh_status_q() {
|
||||
rh_status >/dev/null 2>&1
|
||||
}
|
||||
|
||||
case "$1" in
|
||||
start)
|
||||
rh_status_q && exit 0
|
||||
$1
|
||||
;;
|
||||
stop)
|
||||
rh_status_q || exit 0
|
||||
$1
|
||||
;;
|
||||
restart|force-reload)
|
||||
$1
|
||||
;;
|
||||
reload)
|
||||
rh_status_q || exit 7
|
||||
$1
|
||||
;;
|
||||
status)
|
||||
rh_status
|
||||
;;
|
||||
*)
|
||||
echo "Usage: $0 {start|stop|restart|reload|force-reload|status}" >&2
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
exit 0
|
||||
|
||||
@@ -360,6 +360,8 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
|
||||
vassal_argv[1] = "--yaml";
|
||||
if (!strcmp(name + (strlen(name) - 3), ".js"))
|
||||
vassal_argv[1] = "--json";
|
||||
if (!strcmp(name + (strlen(name) - 5), ".json"))
|
||||
vassal_argv[1] = "--json";
|
||||
|
||||
if (colon) {
|
||||
colon[0] = ':';
|
||||
@@ -660,7 +662,8 @@ reconnect:
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 4), ".ini") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 4), ".yml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 5), ".yaml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 3), ".js")
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 3), ".js") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 5), ".json")
|
||||
) {
|
||||
|
||||
|
||||
@@ -707,6 +710,7 @@ reconnect:
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".ini") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".yml") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 3), ".js") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 5), ".json") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 5), ".yaml")
|
||||
) {
|
||||
|
||||
|
||||
@@ -693,6 +693,8 @@ static int timerfd_create (clockid_t __clock_id, int __flags) {
|
||||
return syscall(283, __clock_id, __flags);
|
||||
#elif defined(__i386__)
|
||||
return syscall(322, __clock_id, __flags);
|
||||
#else
|
||||
return -1;
|
||||
#endif
|
||||
}
|
||||
|
||||
|
||||
@@ -507,9 +507,10 @@ int master_loop(char **argv, char **environ) {
|
||||
}
|
||||
|
||||
// first subscription
|
||||
for(i=0;i<uwsgi.subscriptions_cnt;i++) {
|
||||
uwsgi_log("requested subscription for %s\n", uwsgi.subscriptions[i]);
|
||||
uwsgi_subscribe(uwsgi.subscriptions[i]);
|
||||
struct uwsgi_string_list *subscriptions = uwsgi.subscriptions;
|
||||
while(subscriptions) {
|
||||
uwsgi_subscribe(subscriptions->value);
|
||||
subscriptions = subscriptions->next;
|
||||
}
|
||||
|
||||
// sync the cache store if needed
|
||||
@@ -542,6 +543,17 @@ int master_loop(char **argv, char **environ) {
|
||||
for (;;) {
|
||||
//uwsgi_log("ready_to_reload %d %d\n", ready_to_reload, uwsgi.numproc);
|
||||
|
||||
for (i = 0; i < uwsgi.gp_cnt; i++) {
|
||||
if (uwsgi.gp[i]->master_cycle) {
|
||||
uwsgi.gp[i]->master_cycle();
|
||||
}
|
||||
}
|
||||
for (i = 0; i < 0xFF; i++) {
|
||||
if (uwsgi.p[i]->master_cycle) {
|
||||
uwsgi.p[i]->master_cycle();
|
||||
}
|
||||
}
|
||||
|
||||
if (uwsgi.to_outworld) {
|
||||
//uwsgi_log("%d/%d\n", uwsgi.lazy_respawned, uwsgi.numproc);
|
||||
if (uwsgi.lazy_respawned >= uwsgi.numproc) {
|
||||
@@ -702,6 +714,7 @@ healthy:
|
||||
uwsgi_log( "closing all non-uwsgi socket fds > 2 (_SC_OPEN_MAX = %ld)...\n", sysconf(_SC_OPEN_MAX));
|
||||
for (i = 3; i < sysconf(_SC_OPEN_MAX); i++) {
|
||||
int found = 0;
|
||||
|
||||
struct uwsgi_socket *uwsgi_sock = uwsgi.sockets;
|
||||
while(uwsgi_sock) {
|
||||
if (i == uwsgi_sock->fd) {
|
||||
@@ -719,6 +732,13 @@ healthy:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (uwsgi.original_log_fd > -1) {
|
||||
if (i == uwsgi.original_log_fd) {
|
||||
dup2(uwsgi.original_log_fd, 1);
|
||||
dup2(1, 2);
|
||||
}
|
||||
}
|
||||
if (!found) {
|
||||
#ifdef __APPLE__
|
||||
fcntl(i, F_SETFD, FD_CLOEXEC);
|
||||
@@ -1396,9 +1416,11 @@ healthy:
|
||||
}
|
||||
|
||||
// resubscribe every 10 cycles
|
||||
if (uwsgi.subscriptions_cnt > 0 && ((uwsgi.master_cycles % 10) == 0 || uwsgi.master_cycles == 1)) {
|
||||
for(i=0;i<uwsgi.subscriptions_cnt;i++) {
|
||||
uwsgi_subscribe(uwsgi.subscriptions[i]);
|
||||
if (uwsgi.subscriptions && ((uwsgi.master_cycles % 10) == 0 || uwsgi.master_cycles == 1)) {
|
||||
struct uwsgi_string_list *subscriptions = uwsgi.subscriptions;
|
||||
while(subscriptions) {
|
||||
uwsgi_subscribe(subscriptions->value);
|
||||
subscriptions = subscriptions->next;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+5
-2
@@ -59,8 +59,10 @@ void uwsgi_fixup_fds(int wid, int muleid) {
|
||||
#endif
|
||||
|
||||
if (uwsgi.shared->mule_signal_pipe[0] != -1) close(uwsgi.shared->mule_signal_pipe[0]);
|
||||
|
||||
if (muleid == 0) {
|
||||
if (uwsgi.shared->mule_signal_pipe[1] != -1) close(uwsgi.shared->mule_signal_pipe[1]);
|
||||
if (uwsgi.shared->mule_queue_pipe[1] != -1) close(uwsgi.shared->mule_queue_pipe[1]);
|
||||
}
|
||||
|
||||
for(i=0;i<uwsgi.mules_cnt;i++) {
|
||||
@@ -99,8 +101,9 @@ int uwsgi_respawn_worker(int wid) {
|
||||
uwsgi.workers[uwsgi.mywid].pid = uwsgi.mypid;
|
||||
uwsgi.workers[uwsgi.mywid].id = uwsgi.mywid;
|
||||
uwsgi.workers[uwsgi.mywid].harakiri = 0;
|
||||
uwsgi.workers[uwsgi.mywid].requests = 0;
|
||||
uwsgi.workers[uwsgi.mywid].failed_requests = 0;
|
||||
// do not reset worker counters on reload !!!
|
||||
//uwsgi.workers[uwsgi.mywid].requests = 0;
|
||||
//uwsgi.workers[uwsgi.mywid].failed_requests = 0;
|
||||
uwsgi.workers[uwsgi.mywid].respawn_count++;
|
||||
uwsgi.workers[uwsgi.mywid].last_spawn = uwsgi.current_time;
|
||||
uwsgi.workers[uwsgi.mywid].manage_next_request = 1;
|
||||
|
||||
@@ -47,7 +47,7 @@ void uwsgi_mule(int id) {
|
||||
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);
|
||||
// never here
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -136,6 +136,7 @@ void uwsgi_mule_handler() {
|
||||
event_queue_add_fd_read(mule_queue, uwsgi.signal_socket);
|
||||
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]);
|
||||
event_queue_add_fd_read(mule_queue, uwsgi.shared->mule_queue_pipe[1]);
|
||||
|
||||
uwsgi_mule_add_farm_to_queue(mule_queue);
|
||||
|
||||
@@ -155,19 +156,26 @@ void uwsgi_mule_handler() {
|
||||
uwsgi_log_verbose("master sent signal %d to mule %d\n", uwsgi_signal, uwsgi.muleid);
|
||||
#endif
|
||||
if (uwsgi_signal_handler(uwsgi_signal)) {
|
||||
uwsgi_log_verbose("error managing signal %d on mule %d\n", uwsgi_signal, uwsgi.mywid);
|
||||
uwsgi_log_verbose("error managing signal %d on mule %d\n", uwsgi_signal, uwsgi.muleid);
|
||||
}
|
||||
}
|
||||
else if (interesting_fd == uwsgi.mules[uwsgi.muleid-1].queue_pipe[1] || farm_has_msg(interesting_fd)) {
|
||||
else if (interesting_fd == uwsgi.mules[uwsgi.muleid-1].queue_pipe[1] || interesting_fd == uwsgi.shared->mule_queue_pipe[1] || farm_has_msg(interesting_fd)) {
|
||||
len = read(interesting_fd, message, 65536);
|
||||
if (len < 0) {
|
||||
uwsgi_error("read()");
|
||||
}
|
||||
else if (len == 0) {
|
||||
exit(1);
|
||||
}
|
||||
else {
|
||||
uwsgi_log("*** mule %d received a %d bytes message ***\n", uwsgi.muleid, len);
|
||||
int i,found = 0;
|
||||
for(i=0;i<0xff;i++) {
|
||||
if (uwsgi.p[i]->mule_msg) {
|
||||
if (uwsgi.p[i]->mule_msg(message, len)) {
|
||||
found = 1;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!found)
|
||||
uwsgi_log("*** mule %d received a %d bytes message ***\n", uwsgi.muleid, len);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -187,6 +195,20 @@ struct uwsgi_mule *get_mule_by_id(int id) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
struct uwsgi_farm *get_farm_by_name(char *name) {
|
||||
|
||||
int i;
|
||||
|
||||
for(i=0;i<uwsgi.farms_cnt;i++) {
|
||||
if (!strcmp(uwsgi.farms[i].name, name)) {
|
||||
return &uwsgi.farms[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;
|
||||
|
||||
+25
-25
@@ -95,7 +95,7 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
char buffer[0xffff];
|
||||
|
||||
char *argv[0xff] ;
|
||||
int argc;
|
||||
int argc = 0;
|
||||
uint16_t ret;
|
||||
ei_x_buff xr;
|
||||
|
||||
@@ -103,19 +103,20 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("%d %c %c %c\n", etype, etype, ERL_SMALL_TUPLE_EXT, ERL_LARGE_TUPLE_EXT);
|
||||
#endif
|
||||
if (etype != ERL_SMALL_TUPLE_EXT && etype != ERL_LARGE_TUPLE_EXT) return;
|
||||
|
||||
uwsgi_log("decode tuple\n");
|
||||
ei_decode_tuple_header(x->buff, &x->index, &arity);
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("rpc arity %d\n", arity);
|
||||
#endif
|
||||
if (arity != 3) return ;
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
|
||||
uwsgi_log("%d %c\n", etype, etype);
|
||||
|
||||
if (etype != ERL_ATOM_EXT && etype != ERL_STRING_EXT) return ;
|
||||
|
||||
gen_call = uwsgi_malloc(esize);
|
||||
@@ -127,7 +128,9 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
ei_decode_string(x->buff, &x->index, gen_call);
|
||||
}
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("gen call = %s\n", gen_call);
|
||||
#endif
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
|
||||
@@ -138,13 +141,10 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
ei_skip_term(x->buff, &x->index);
|
||||
uwsgi_log("skip0 %d %c\n", etype, etype);
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
uwsgi_log("skip1 %d %c\n", etype, etype);
|
||||
ei_decode_ref(x->buff, &x->index, &eref);
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
uwsgi_log("%d %c\n", etype, etype);
|
||||
|
||||
module = uwsgi_malloc(esize);
|
||||
|
||||
@@ -157,17 +157,16 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
|
||||
uwsgi_log("%d %c\n", etype, etype);
|
||||
|
||||
if (etype != ERL_SMALL_TUPLE_EXT) return ;
|
||||
|
||||
ei_decode_tuple_header(x->buff, &x->index, &arity);
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("arity: %d\n", arity);
|
||||
#endif
|
||||
if (arity != 5) return ;
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
uwsgi_log("%d %c\n", etype, etype);
|
||||
|
||||
char *method = uwsgi_malloc(esize);
|
||||
|
||||
@@ -181,7 +180,6 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
if (strcmp(method, "call")) return;
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
uwsgi_log("%d %c\n", etype, etype);
|
||||
|
||||
if (etype != ERL_ATOM_EXT && etype != ERL_STRING_EXT) return ;
|
||||
|
||||
@@ -207,7 +205,9 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
ei_decode_string(x->buff, &x->index, call);
|
||||
}
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("RPC %s %s\n", module, call);
|
||||
#endif
|
||||
|
||||
ei_get_type(x->buff, &x->index, &etype, &esize);
|
||||
|
||||
@@ -224,7 +224,9 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) {
|
||||
|
||||
ret = uwsgi_rpc(call, argc, argv, buffer);
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("buffer: %.*s\n", ret, buffer);
|
||||
#endif
|
||||
|
||||
ei_x_new_with_version(&xr);
|
||||
|
||||
@@ -247,7 +249,6 @@ void erlang_loop() {
|
||||
int fd;
|
||||
|
||||
int eversion;
|
||||
int i;
|
||||
|
||||
ei_x_buff x, xr;
|
||||
|
||||
@@ -286,7 +287,7 @@ void erlang_loop() {
|
||||
if (em.msgtype == ERL_TICK)
|
||||
continue;
|
||||
|
||||
uwsgi_log("From: %s To: %s RegName: %s\n", em.from.node, em.to.node, em.toname);
|
||||
uwsgi_log("[erlang] message From: %s To (process): %s\n", em.from.node, em.toname);
|
||||
|
||||
|
||||
|
||||
@@ -300,20 +301,18 @@ void erlang_loop() {
|
||||
uwsgi_erlang_rpc(fd, &em.from, &x);
|
||||
}
|
||||
else {
|
||||
int uep = -1;
|
||||
for(i=0;i<uerl.uep_cnt;i++) {
|
||||
if (!strcmp(uerl.uep[i].name, em.toname)) {
|
||||
uep = i;
|
||||
struct uwsgi_erlang_process *uep = uerl.uep;
|
||||
while(uep) {
|
||||
if (!strcmp(uep->name, em.toname)) {
|
||||
if (uep->plugin) {
|
||||
uep->plugin(uep->func, &x);
|
||||
}
|
||||
break;
|
||||
}
|
||||
uep = uep->next;
|
||||
}
|
||||
|
||||
if (uep > -1) {
|
||||
if (uerl.uep[uep].plugin) {
|
||||
uerl.uep[uep].plugin( uerl.uep[uep].func, &x );
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (!uep) {
|
||||
uwsgi_log("!!! unregistered erlang process requested, dumping it !!!\n");
|
||||
dump_eterm(&x);
|
||||
}
|
||||
@@ -366,7 +365,6 @@ int erlang_init() {
|
||||
|
||||
if (uerl.name) {
|
||||
|
||||
uwsgi.master_process = 1;
|
||||
|
||||
host = strchr(uerl.name, '@');
|
||||
|
||||
@@ -426,6 +424,7 @@ int erlang_init() {
|
||||
uwsgi_log("unable to register the erlang gateway\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
return 0;
|
||||
@@ -435,7 +434,8 @@ int erlang_opt(int i, char *optarg) {
|
||||
|
||||
switch(i) {
|
||||
case LONG_ARGS_ERLANG:
|
||||
uerl.name = optarg;
|
||||
uwsgi.master_process = 1;
|
||||
uerl.name = optarg;
|
||||
return 1;
|
||||
case LONG_ARGS_ERLANG_COOKIE:
|
||||
uerl.cookie = optarg;
|
||||
|
||||
@@ -1,7 +1,5 @@
|
||||
#include <ei.h>
|
||||
|
||||
#define MAX_UWSGI_ERLANG_PROCESSES 64
|
||||
|
||||
#define LONG_ARGS_ERLANG 17012
|
||||
#define LONG_ARGS_ERLANG_COOKIE 17013
|
||||
|
||||
@@ -12,6 +10,8 @@ struct uwsgi_erlang_process {
|
||||
char name[0xff];
|
||||
void (*plugin)(void *, ei_x_buff *);
|
||||
void *func;
|
||||
|
||||
struct uwsgi_erlang_process *next;
|
||||
};
|
||||
|
||||
struct uwsgi_erlang {
|
||||
@@ -24,7 +24,6 @@ struct uwsgi_erlang {
|
||||
|
||||
void *lock;
|
||||
|
||||
struct uwsgi_erlang_process uep[MAX_UWSGI_ERLANG_PROCESSES];
|
||||
int uep_cnt;
|
||||
struct uwsgi_erlang_process *uep;
|
||||
};
|
||||
|
||||
|
||||
@@ -28,6 +28,8 @@
|
||||
#define FASTROUTER_STATUS_RESPONSE 4
|
||||
|
||||
#define add_timeout(x) uwsgi_add_rb_timer(ufr.timeouts, time(NULL)+ufr.socket_timeout, x)
|
||||
#define add_check_timeout(x) uwsgi_add_rb_timer(ufr.timeouts, time(NULL)+x, NULL)
|
||||
#define del_check_timeout(x) rb_erase(&x->rbt, ufr.timeouts);
|
||||
#define del_timeout(x) rb_erase(&x->timeout->rbt, ufr.timeouts); free(x->timeout);
|
||||
|
||||
struct uwsgi_fastrouter_socket {
|
||||
@@ -54,7 +56,8 @@ struct uwsgi_fastrouter {
|
||||
int base_len;
|
||||
|
||||
char *subscription_server;
|
||||
struct uwsgi_dict *subscription_dict;
|
||||
struct uwsgi_subscribe_slot *subscriptions;
|
||||
int subscription_regexp;
|
||||
|
||||
int socket_timeout;
|
||||
|
||||
@@ -64,6 +67,8 @@ struct uwsgi_fastrouter {
|
||||
|
||||
struct rb_root *timeouts;
|
||||
|
||||
struct uwsgi_rb_timer *subscriptions_check;
|
||||
|
||||
int cheap;
|
||||
int i_am_cheap;
|
||||
} ufr;
|
||||
@@ -105,6 +110,7 @@ struct option fastrouter_options[] = {
|
||||
{"fastrouter-cheap", no_argument, &ufr.cheap, 1},
|
||||
{"fastrouter-subscription-server", required_argument, 0, LONG_ARGS_FASTROUTER_SUBSCRIPTION_SERVER},
|
||||
{"fastrouter-subscription-slot", required_argument, 0, LONG_ARGS_FASTROUTER_SUBSCRIPTION_SLOT},
|
||||
{"fastrouter-subscription-use-regexp", no_argument, &ufr.subscription_regexp, 1},
|
||||
{"fastrouter-timeout", required_argument, 0, LONG_ARGS_FASTROUTER_TIMEOUT},
|
||||
{0, 0, 0, 0},
|
||||
};
|
||||
@@ -136,7 +142,6 @@ struct fastrouter_session {
|
||||
int status;
|
||||
struct uwsgi_header uh;
|
||||
uint8_t h_pos;
|
||||
char buffer[0xffff];
|
||||
uint16_t pos;
|
||||
|
||||
char *hostname;
|
||||
@@ -147,27 +152,47 @@ struct fastrouter_session {
|
||||
char *instance_address;
|
||||
uint64_t instance_address_len;
|
||||
|
||||
struct uwsgi_subscriber_name *un;
|
||||
struct uwsgi_subscribe_node *un;
|
||||
int pass_fd;
|
||||
|
||||
struct uwsgi_rb_timer *timeout;
|
||||
int instance_failed;
|
||||
|
||||
uint8_t modifier1;
|
||||
|
||||
char buffer[0xffff];
|
||||
};
|
||||
|
||||
static void close_session(struct fastrouter_session **fr_table, struct fastrouter_session *fr_session) {
|
||||
|
||||
// check timeout expired
|
||||
if (fr_session == NULL) {
|
||||
if (ufr.subscription_server) {
|
||||
//time_t current_time = time(NULL);
|
||||
uwsgi_log("checking for node health\n");
|
||||
struct uwsgi_subscribe_slot *slot = ufr.subscriptions;
|
||||
while(slot) {
|
||||
struct uwsgi_subscribe_node *node = slot->nodes;
|
||||
while(node) {
|
||||
uwsgi_log("%.*s (hits: %llu) %.*s\n", slot->keylen, slot->key, slot->hits, node->len, node->name);
|
||||
node = node->next;
|
||||
}
|
||||
slot = slot->next;
|
||||
}
|
||||
del_check_timeout(ufr.subscriptions_check);
|
||||
ufr.subscriptions_check = add_check_timeout(10);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
close(fr_session->fd);
|
||||
fr_table[fr_session->fd] = NULL;
|
||||
if (fr_session->instance_fd != -1) {
|
||||
if (ufr.subscription_server && (fr_session->instance_failed || fr_session->status == FASTROUTER_STATUS_CONNECTING)) {
|
||||
uwsgi_log("[uwsgi-fastrouter] %.*s => marking %.*s as failed\n", (int) fr_session->hostname_len, fr_session->hostname, (int) fr_session->instance_address_len,fr_session->instance_address);
|
||||
if (fr_session->un->len > 0) {
|
||||
if (ufr.subscription_dict->count > 0) {
|
||||
ufr.subscription_dict->count--;
|
||||
}
|
||||
if (ufr.subscription_dict->count == 0 && ufr.cheap && !ufr.i_am_cheap) {
|
||||
if (ufr.subscriptions && (fr_session->instance_failed || fr_session->status == FASTROUTER_STATUS_CONNECTING)) {
|
||||
if (fr_session->un && fr_session->un->len > 0) {
|
||||
uwsgi_log("[uwsgi-fastrouter] %.*s => marking %.*s as failed\n", (int) fr_session->hostname_len, fr_session->hostname, (int) fr_session->instance_address_len,fr_session->instance_address);
|
||||
uwsgi_remove_subscribe_node(&ufr.subscriptions, fr_session->un);
|
||||
if (ufr.subscriptions == NULL && ufr.cheap && !ufr.i_am_cheap) {
|
||||
uwsgi_log("[uwsgi-fastrouter] no more nodes available. Going cheap...\n");
|
||||
struct uwsgi_fastrouter_socket *ufr_sock = ufr.sockets;
|
||||
while(ufr_sock) {
|
||||
@@ -202,7 +227,6 @@ static void expire_timeouts(struct fastrouter_session **fr_table) {
|
||||
|
||||
if (urbt->key <= current) {
|
||||
close_session(fr_table, (struct fastrouter_session *)urbt->data);
|
||||
uwsgi_log("timeout !!!\n");
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -327,20 +351,22 @@ void fastrouter_loop() {
|
||||
|
||||
events = event_queue_alloc(ufr.nevents);
|
||||
|
||||
ufr.timeouts = uwsgi_init_rb_timer();
|
||||
if (!ufr.socket_timeout) ufr.socket_timeout = 30;
|
||||
|
||||
|
||||
if (ufr.subscription_server) {
|
||||
ufr_subserver = bind_to_udp(ufr.subscription_server, 0, 0);
|
||||
event_queue_add_fd_read(ufr.queue, ufr_subserver);
|
||||
if (!ufr.subscription_slot) ufr.subscription_slot = 30;
|
||||
ufr.subscription_dict = uwsgi_dict_create(ufr.subscription_slot, 0);
|
||||
// check for node status every 10 seconds
|
||||
//ufr.subscriptions_check = add_check_timeout(10);
|
||||
}
|
||||
|
||||
if (ufr.pattern) {
|
||||
init_magic_table(magic_table);
|
||||
}
|
||||
|
||||
ufr.timeouts = uwsgi_init_rb_timer();
|
||||
if (!ufr.socket_timeout) ufr.socket_timeout = 30;
|
||||
|
||||
for (;;) {
|
||||
|
||||
@@ -386,6 +412,8 @@ void fastrouter_loop() {
|
||||
fr_table[new_connection]->un = NULL;
|
||||
fr_table[new_connection]->instance_failed = 0;
|
||||
fr_table[new_connection]->instance_address_len = 0;
|
||||
fr_table[new_connection]->hostname_len = 0;
|
||||
fr_table[new_connection]->hostname = NULL;
|
||||
|
||||
fr_table[new_connection]->timeout = add_timeout(fr_table[new_connection]);
|
||||
|
||||
@@ -409,8 +437,7 @@ void fastrouter_loop() {
|
||||
if (len > 0) {
|
||||
memset(&usr, 0, sizeof(struct uwsgi_subscribe_req));
|
||||
uwsgi_hooked_parse(bbuf+4, len-4, fastrouter_manage_subscription, &usr);
|
||||
uwsgi_add_subscriber(ufr.subscription_dict, &usr);
|
||||
if (ufr.i_am_cheap) {
|
||||
if (uwsgi_add_subscribe_node(&ufr.subscriptions, &usr, ufr.subscription_regexp) && ufr.i_am_cheap) {
|
||||
struct uwsgi_fastrouter_socket *ufr_sock = ufr.sockets;
|
||||
while(ufr_sock) {
|
||||
event_queue_add_fd_read(ufr.queue, ufr_sock->fd);
|
||||
@@ -473,7 +500,7 @@ void fastrouter_loop() {
|
||||
}
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("requested domain %.*s\n", fr_session->hostname_len, fr_session->hostname);
|
||||
//uwsgi_log("requested domain %.*s\n", fr_session->hostname_len, fr_session->hostname);
|
||||
#endif
|
||||
if (ufr.use_cache) {
|
||||
fr_session->instance_address = uwsgi_cache_get(fr_session->hostname, fr_session->hostname_len, &fr_session->instance_address_len);
|
||||
@@ -491,7 +518,7 @@ void fastrouter_loop() {
|
||||
fr_session->instance_address = tmp_socket_name;
|
||||
}
|
||||
else if (ufr.subscription_server) {
|
||||
fr_session->un = uwsgi_get_subscriber(ufr.subscription_dict, fr_session->hostname, fr_session->hostname_len);
|
||||
fr_session->un = uwsgi_get_subscribe_node(&ufr.subscriptions, fr_session->hostname, fr_session->hostname_len, ufr.subscription_regexp);
|
||||
if (fr_session->un && fr_session->un->len) {
|
||||
fr_session->instance_address = fr_session->un->name;
|
||||
fr_session->instance_address_len = fr_session->un->len;
|
||||
@@ -531,14 +558,12 @@ void fastrouter_loop() {
|
||||
if (tmp_socket_name) free(tmp_socket_name);
|
||||
|
||||
if (fr_session->instance_fd < 0) {
|
||||
/*
|
||||
if (ufr.subscription_server) {
|
||||
uwsgi_log("[uwsgi-fastrouter] %.*s => marking %.*s as failed\n", (int) fr_session->hostname_len, fr_session->hostname, (int) fr_session->instance_address_len,fr_session->instance_address);
|
||||
if (fr_session->un->len > 0) {
|
||||
fr_session->un->len = 0;
|
||||
if (ufr.subscription_dict->count > 0) {
|
||||
ufr.subscription_dict->count--;
|
||||
}
|
||||
if (ufr.subscription_dict->count == 0 && ufr.cheap && !ufr.i_am_cheap) {
|
||||
if (fr_session->un && fr_session->un->len > 0) {
|
||||
uwsgi_log("[uwsgi-fastrouter] %.*s => marking %.*s as failed\n", (int) fr_session->hostname_len, fr_session->hostname, (int) fr_session->instance_address_len,fr_session->instance_address);
|
||||
uwsgi_remove_subscribe_node(&ufr.subscriptions, fr_session->un);
|
||||
if (ufr.subscriptions == NULL && ufr.cheap && !ufr.i_am_cheap) {
|
||||
uwsgi_log("[uwsgi-fastrouter] no more nodes available. Going cheap...\n");
|
||||
struct uwsgi_fastrouter_socket *ufr_sock = ufr.sockets;
|
||||
while(ufr_sock) {
|
||||
@@ -549,6 +574,8 @@ void fastrouter_loop() {
|
||||
}
|
||||
}
|
||||
}
|
||||
*/
|
||||
fr_session->instance_failed = 1;
|
||||
close_session(fr_table, fr_session);
|
||||
break;
|
||||
}
|
||||
@@ -678,7 +705,7 @@ void fastrouter_loop() {
|
||||
|
||||
// fallback to destroy !!!
|
||||
default:
|
||||
uwsgi_log("default action\n");
|
||||
uwsgi_log("unknown event: closing session\n");
|
||||
close_session(fr_table, fr_session);
|
||||
break;
|
||||
|
||||
|
||||
+9
-11
@@ -39,6 +39,7 @@ struct uwsgi_http {
|
||||
int server;
|
||||
|
||||
char *subscription_server;
|
||||
int subscription_regexp;
|
||||
|
||||
char *pattern;
|
||||
int pattern_len;
|
||||
@@ -60,7 +61,7 @@ struct uwsgi_http {
|
||||
|
||||
int socket_timeout;
|
||||
|
||||
struct uwsgi_dict *subscription_dict;
|
||||
struct uwsgi_subscribe_slot *subscriptions;
|
||||
|
||||
struct rb_root *timeouts;
|
||||
} uhttp;
|
||||
@@ -76,6 +77,7 @@ struct option http_options[] = {
|
||||
{"http-use-cluster", no_argument, &uhttp.use_cluster, 1},
|
||||
{"http-events", required_argument, 0, LONG_ARGS_HTTP_EVENTS},
|
||||
{"http-subscription-server", required_argument, 0, LONG_ARGS_HTTP_SUBSCRIPTION_SERVER},
|
||||
{"http-subscription-use-regexp", no_argument, &uhttp.subscription_regexp, 1},
|
||||
{"http-timeout", required_argument, 0, LONG_ARGS_HTTP_TIMEOUT},
|
||||
{0, 0, 0, 0},
|
||||
};
|
||||
@@ -139,7 +141,7 @@ struct http_session {
|
||||
char path_info[UMAX16];
|
||||
uint16_t path_info_len;
|
||||
|
||||
struct uwsgi_subscriber_name *un;
|
||||
struct uwsgi_subscribe_node *un;
|
||||
|
||||
in_addr_t ip_addr;
|
||||
char ip[INET_ADDRSTRLEN];
|
||||
@@ -158,9 +160,9 @@ static void close_session(struct http_session **uhttp_table, struct http_session
|
||||
close(uhttp_session->fd);
|
||||
uhttp_table[uhttp_session->fd] = NULL;
|
||||
if (uhttp_session->instance_fd != -1) {
|
||||
if (uhttp.subscription_server && (uhttp_session->instance_failed || uhttp_session->status == HTTP_STATUS_CONNECTING)) {
|
||||
if (uhttp.subscriptions && (uhttp_session->instance_failed || uhttp_session->status == HTTP_STATUS_CONNECTING)) {
|
||||
uwsgi_log("marking %.*s as failed\n", (int) uhttp_session->instance_address_len,uhttp_session->instance_address);
|
||||
uhttp_session->un->len = 0;
|
||||
uwsgi_remove_subscribe_node(&uhttp.subscriptions, uhttp_session->un);
|
||||
}
|
||||
close(uhttp_session->instance_fd);
|
||||
uhttp_table[uhttp_session->instance_fd] = NULL;
|
||||
@@ -482,7 +484,6 @@ void http_loop() {
|
||||
if (uhttp.subscription_server) {
|
||||
uhttp_subserver = bind_to_udp(uhttp.subscription_server, 0, 0);
|
||||
event_queue_add_fd_read(uhttp_queue, uhttp_subserver);
|
||||
uhttp.subscription_dict = uwsgi_dict_create(100, 0);
|
||||
}
|
||||
|
||||
if (uhttp.pattern) {
|
||||
@@ -558,7 +559,7 @@ void http_loop() {
|
||||
if (len > 0) {
|
||||
memset(&usr, 0, sizeof(struct uwsgi_subscribe_req));
|
||||
uwsgi_hooked_parse(bbuf+4, len-4, http_manage_subscription, &usr);
|
||||
uwsgi_add_subscriber(uhttp.subscription_dict, &usr);
|
||||
uwsgi_add_subscribe_node(&uhttp.subscriptions, &usr, uhttp.subscription_regexp);
|
||||
}
|
||||
}
|
||||
else {
|
||||
@@ -638,7 +639,7 @@ void http_loop() {
|
||||
uhttp_session->instance_address_len = uhttp.to_len;
|
||||
}
|
||||
else if (uhttp.subscription_server) {
|
||||
uhttp_session->un = uwsgi_get_subscriber(uhttp.subscription_dict, uhttp_session->hostname, uhttp_session->hostname_len);
|
||||
uhttp_session->un = uwsgi_get_subscribe_node(&uhttp.subscriptions, uhttp_session->hostname, uhttp_session->hostname_len, uhttp.subscription_regexp);
|
||||
if (uhttp_session->un && uhttp_session->un->len) {
|
||||
uhttp_session->instance_address = uhttp_session->un->name;
|
||||
uhttp_session->instance_address_len = uhttp_session->un->len;
|
||||
@@ -669,10 +670,7 @@ void http_loop() {
|
||||
}
|
||||
|
||||
if (uhttp_session->instance_fd < 0) {
|
||||
if (uhttp.subscription_server) {
|
||||
uwsgi_log("marking %.*s as failed\n", (int) uhttp_session->instance_address_len,uhttp_session->instance_address);
|
||||
uhttp_session->un->len = 0;
|
||||
}
|
||||
uhttp_session->instance_failed = 1;
|
||||
close_session(uhttp_table, uhttp_session);
|
||||
break;
|
||||
}
|
||||
|
||||
+21
-8
@@ -304,7 +304,7 @@ void pyerl_call_registered(void *func, ei_x_buff *x) {
|
||||
|
||||
PyTuple_SetItem(pyargs, 0, erl_to_py(x));
|
||||
|
||||
python_call((PyObject *) func, pyargs, 0);
|
||||
python_call((PyObject *) func, pyargs, 0, NULL);
|
||||
}
|
||||
|
||||
PyObject *pyerl_register_process(PyObject * self, PyObject * args) {
|
||||
@@ -316,17 +316,30 @@ PyObject *pyerl_register_process(PyObject * self, PyObject * args) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
if (uerl.uep_cnt >= MAX_UWSGI_ERLANG_PROCESSES)
|
||||
return PyErr_Format(PyExc_ValueError, "You can define max %d erlang registered processes", MAX_UWSGI_ERLANG_PROCESSES);
|
||||
|
||||
if (strlen(name) > 0xff-1)
|
||||
return PyErr_Format(PyExc_ValueError, "Invalid erlang process name");
|
||||
|
||||
strcpy(uerl.uep[uerl.uep_cnt].name, name);
|
||||
uerl.uep[uerl.uep_cnt].plugin = pyerl_call_registered;
|
||||
uerl.uep[uerl.uep_cnt].func = callable;
|
||||
struct uwsgi_erlang_process *uep = uerl.uep, *old_uep;
|
||||
|
||||
if (!uep) {
|
||||
uerl.uep = uwsgi_malloc(sizeof(struct uwsgi_erlang_process));
|
||||
uep = uerl.uep;
|
||||
}
|
||||
else {
|
||||
while(uep) {
|
||||
old_uep = uep;
|
||||
uep = uep->next;
|
||||
}
|
||||
|
||||
uep = uwsgi_malloc(sizeof(struct uwsgi_erlang_process));
|
||||
old_uep->next = uep;
|
||||
}
|
||||
|
||||
strcpy(uep->name, name);
|
||||
uep->plugin = pyerl_call_registered;
|
||||
uep->func = callable;
|
||||
uep->next = NULL;
|
||||
|
||||
uerl.uep_cnt++;
|
||||
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
|
||||
@@ -915,9 +915,11 @@ void uwsgi_python_init_apps() {
|
||||
|
||||
|
||||
#ifdef __linux__
|
||||
#if !defined(PYTHREE) && !defined(UWSGI_PYPY)
|
||||
#ifndef UWSGI_PYPY
|
||||
#ifdef UWSGI_EMBEDDED
|
||||
uwsgi_init_symbol_import();
|
||||
#endif
|
||||
#endif
|
||||
#endif
|
||||
|
||||
if (up.test_module != NULL) {
|
||||
@@ -1241,6 +1243,8 @@ clear:
|
||||
|
||||
uint16_t uwsgi_python_rpc(void *func, uint8_t argc, char **argv, char *buffer) {
|
||||
|
||||
UWSGI_GET_GIL;
|
||||
|
||||
uint8_t i;
|
||||
PyObject *pyargs = PyTuple_New(argc);
|
||||
PyObject *ret;
|
||||
@@ -1263,6 +1267,7 @@ uint16_t uwsgi_python_rpc(void *func, uint8_t argc, char **argv, char *buffer) {
|
||||
if (rl <= 0xffff) {
|
||||
memcpy(buffer, rv, rl);
|
||||
Py_DECREF(ret);
|
||||
UWSGI_RELEASE_GIL;
|
||||
return rl;
|
||||
}
|
||||
}
|
||||
@@ -1271,6 +1276,7 @@ uint16_t uwsgi_python_rpc(void *func, uint8_t argc, char **argv, char *buffer) {
|
||||
if (PyErr_Occurred())
|
||||
PyErr_Print();
|
||||
|
||||
UWSGI_RELEASE_GIL;
|
||||
|
||||
return 0;
|
||||
|
||||
@@ -1398,7 +1404,9 @@ int uwsgi_python_mule(char *opt) {
|
||||
|
||||
if (uwsgi_endswith(opt, ".py")) {
|
||||
UWSGI_GET_GIL;
|
||||
uwsgi_pyimport_by_filename("__main__", opt);
|
||||
if (uwsgi_pyimport_by_filename("__main__", opt) == NULL) {
|
||||
return 0;
|
||||
}
|
||||
UWSGI_RELEASE_GIL;
|
||||
return 1;
|
||||
}
|
||||
@@ -1407,6 +1415,33 @@ int uwsgi_python_mule(char *opt) {
|
||||
|
||||
}
|
||||
|
||||
int uwsgi_python_mule_msg(char *message, size_t len) {
|
||||
|
||||
UWSGI_GET_GIL;
|
||||
|
||||
PyObject *mule_msg_hook = PyDict_GetItemString(up.embedded_dict, "mule_msg_hook");
|
||||
if (!mule_msg_hook) {
|
||||
// ignore
|
||||
UWSGI_RELEASE_GIL;
|
||||
return 0;
|
||||
}
|
||||
|
||||
PyObject *pyargs = PyTuple_New(1);
|
||||
PyTuple_SetItem(pyargs, 0, PyString_FromStringAndSize(message, len));
|
||||
|
||||
PyObject *ret = python_call(mule_msg_hook, pyargs, 0, NULL);
|
||||
Py_DECREF(pyargs);
|
||||
if (ret) {
|
||||
Py_DECREF(ret);
|
||||
}
|
||||
|
||||
if (PyErr_Occurred())
|
||||
PyErr_Print();
|
||||
|
||||
UWSGI_RELEASE_GIL;
|
||||
return 1;
|
||||
}
|
||||
|
||||
struct uwsgi_plugin python_plugin = {
|
||||
|
||||
.name = "python",
|
||||
@@ -1444,6 +1479,7 @@ struct uwsgi_plugin python_plugin = {
|
||||
.rpc = uwsgi_python_rpc,
|
||||
|
||||
.mule = uwsgi_python_mule,
|
||||
.mule_msg = uwsgi_python_mule_msg,
|
||||
|
||||
.spooler = uwsgi_python_spooler,
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
#include "uwsgi_python.h"
|
||||
|
||||
#if !defined(PYTHREE) && !defined(UWSGI_PYPY)
|
||||
#ifndef UWSGI_PYPY
|
||||
|
||||
extern struct uwsgi_server uwsgi;
|
||||
extern struct uwsgi_python up;
|
||||
@@ -201,9 +201,14 @@ static PyObject* symzipimporter_load_module(PyObject *self, PyObject *args) {
|
||||
PyObject *source = PyObject_CallMethod(this->zip, "read", "(s)", filename);
|
||||
free(filename);
|
||||
PyObject *code = Py_CompileString(PyString_AsString(source), modname, Py_file_input);
|
||||
if (!code) {
|
||||
PyErr_Print();
|
||||
goto shit;
|
||||
}
|
||||
mod = PyImport_ExecCodeModuleEx(fullname, code, modname);
|
||||
|
||||
Py_DECREF(code);
|
||||
shit:
|
||||
Py_DECREF(source);
|
||||
free(modname);
|
||||
return mod;
|
||||
@@ -231,9 +236,14 @@ static PyObject* symzipimporter_load_module(PyObject *self, PyObject *args) {
|
||||
PyObject *source = PyObject_CallMethod(this->zip, "read", "(s)", filename);
|
||||
free(filename);
|
||||
PyObject *code = Py_CompileString(PyString_AsString(source), modname, Py_file_input);
|
||||
if (!code) {
|
||||
PyErr_Print();
|
||||
goto shit2;
|
||||
}
|
||||
mod = PyImport_ExecCodeModuleEx(fullname, code, modname);
|
||||
|
||||
Py_DECREF(code);
|
||||
shit2:
|
||||
Py_DECREF(source);
|
||||
free(modname);
|
||||
return mod;
|
||||
@@ -306,9 +316,14 @@ static PyObject* symimporter_load_module(PyObject *self, PyObject *args) {
|
||||
modname = uwsgi_concat3("sym://", fullname2, "_py");
|
||||
|
||||
code = Py_CompileString(source, modname, Py_file_input);
|
||||
if (!code) {
|
||||
PyErr_Print();
|
||||
goto shit;
|
||||
}
|
||||
mod = PyImport_ExecCodeModuleEx(fullname, code, modname);
|
||||
|
||||
Py_DECREF(code);
|
||||
shit:
|
||||
free(source);
|
||||
free(modname);
|
||||
free(fullname2);
|
||||
@@ -334,9 +349,14 @@ static PyObject* symimporter_load_module(PyObject *self, PyObject *args) {
|
||||
PyDict_SetItemString(dict, "__loader__", self);
|
||||
|
||||
code = Py_CompileString(source, modname, Py_file_input);
|
||||
if (!code) {
|
||||
PyErr_Print();
|
||||
goto shit2;
|
||||
}
|
||||
mod = PyImport_ExecCodeModuleEx(fullname, code, modname);
|
||||
|
||||
Py_DECREF(code);
|
||||
shit2:
|
||||
free(source);
|
||||
free(modname);
|
||||
free(fullname2);
|
||||
@@ -450,7 +470,10 @@ zipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds)
|
||||
}
|
||||
|
||||
|
||||
PyObject *stringio_dict = PyModule_GetDict(stringio);
|
||||
#ifdef PYTHREE
|
||||
PyObject *source_code = PyObject_CallMethodObjArgs(stringio, PyString_FromString("StringIO"), PyString_FromStringAndSize(body, len));
|
||||
#else
|
||||
PyObject *stringio_dict = PyModule_GetDict(stringio);
|
||||
if (!stringio_dict) {
|
||||
return -1;
|
||||
}
|
||||
@@ -466,7 +489,9 @@ zipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds)
|
||||
PyTuple_SetItem(stringio_args, 0, PyString_FromStringAndSize(body, len));
|
||||
|
||||
|
||||
|
||||
PyObject *source_code = PyInstance_New(stringio_stringio, stringio_args, NULL);
|
||||
#endif
|
||||
if (!source_code) {
|
||||
return -1;
|
||||
}
|
||||
@@ -480,7 +505,10 @@ zipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds)
|
||||
}
|
||||
|
||||
|
||||
PyObject *zipfile_dict = PyModule_GetDict(zipfile);
|
||||
#ifdef PYTHREE
|
||||
self->zip = PyObject_CallMethodObjArgs(zipfile, PyString_FromString("ZipFile"), source_code);
|
||||
#else
|
||||
PyObject *zipfile_dict = PyModule_GetDict(zipfile);
|
||||
if (!zipfile_dict) {
|
||||
return -1;
|
||||
}
|
||||
@@ -494,7 +522,9 @@ zipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds)
|
||||
PyObject *zipfile_args = PyTuple_New(1);
|
||||
PyTuple_SetItem(zipfile_args, 0, source_code);
|
||||
|
||||
|
||||
self->zip = PyInstance_New(zipfile_zipfile, zipfile_args, NULL);
|
||||
#endif
|
||||
if (!self->zip) {
|
||||
return -1;
|
||||
}
|
||||
@@ -556,6 +586,9 @@ symzipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds
|
||||
return -1;
|
||||
}
|
||||
|
||||
#ifdef PYTHREE
|
||||
PyObject *source_code = PyObject_CallMethodObjArgs(stringio, PyString_FromString("StringIO"), PyString_FromStringAndSize(code_start, code_end-code_start));
|
||||
#else
|
||||
PyObject *stringio_dict = PyModule_GetDict(stringio);
|
||||
if (!stringio_dict) {
|
||||
return -1;
|
||||
@@ -569,7 +602,9 @@ symzipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds
|
||||
PyObject *stringio_args = PyTuple_New(1);
|
||||
PyTuple_SetItem(stringio_args, 0, PyString_FromStringAndSize(code_start, code_end-code_start));
|
||||
|
||||
|
||||
PyObject *source_code = PyInstance_New(stringio_stringio, stringio_args, NULL);
|
||||
#endif
|
||||
if (!source_code) {
|
||||
return -1;
|
||||
}
|
||||
@@ -579,6 +614,9 @@ symzipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds
|
||||
return -1;
|
||||
}
|
||||
|
||||
#ifdef PYTHREE
|
||||
self->zip = PyObject_CallMethodObjArgs(zipfile, PyString_FromString("ZipFile"), source_code);
|
||||
#else
|
||||
PyObject *zipfile_dict = PyModule_GetDict(zipfile);
|
||||
if (!zipfile_dict) {
|
||||
return -1;
|
||||
@@ -592,7 +630,9 @@ symzipimporter_init(struct _symzipimporter *self, PyObject *args, PyObject *kwds
|
||||
PyObject *zipfile_args = PyTuple_New(1);
|
||||
PyTuple_SetItem(zipfile_args, 0, source_code);
|
||||
|
||||
|
||||
self->zip = PyInstance_New(zipfile_zipfile, zipfile_args, NULL);
|
||||
#endif
|
||||
if (!self->zip) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
+159
-20
@@ -1037,7 +1037,7 @@ PyObject *py_uwsgi_unlock(PyObject * self, PyObject * args) {
|
||||
}
|
||||
|
||||
|
||||
uwsgi_unlock(uwsgi.user_lock);
|
||||
uwsgi_unlock(uwsgi.user_lock[lock_num]);
|
||||
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
@@ -1153,46 +1153,185 @@ PyObject *py_uwsgi_mule_msg(PyObject * self, PyObject * args) {
|
||||
|
||||
char *message = NULL;
|
||||
Py_ssize_t message_len = 0;
|
||||
int mule_id = 0;
|
||||
PyObject *mule_obj = NULL;
|
||||
ssize_t len;
|
||||
int fd = -1;
|
||||
int mule_id = -1;
|
||||
|
||||
if (!PyArg_ParseTuple(args, "s#|i:mule_msg", &message, &message_len, &mule_id)) {
|
||||
if (!PyArg_ParseTuple(args, "s#|O:mule_msg", &message, &message_len, &mule_obj)) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
if (mule_id == 0) {
|
||||
}
|
||||
else if (mule_id > 0 && mule_id <= uwsgi.mules_cnt) {
|
||||
len = write(uwsgi.mules[mule_id-1].queue_pipe[0], message, message_len);
|
||||
if (uwsgi.mules_cnt < 1)
|
||||
return PyErr_Format(PyExc_ValueError, "no mule configured");
|
||||
|
||||
if (mule_obj == NULL) {
|
||||
len = write(uwsgi.shared->mule_queue_pipe[0], message, message_len);
|
||||
if (len <= 0) {
|
||||
uwsgi_error("write()");
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (PyString_Check(mule_obj)) {
|
||||
struct uwsgi_farm *uf = get_farm_by_name(PyString_AsString(mule_obj));
|
||||
if (uf == NULL) {
|
||||
return PyErr_Format(PyExc_ValueError, "unknown farm");
|
||||
}
|
||||
fd = uf->queue_pipe[0];
|
||||
}
|
||||
else if (PyInt_Check(mule_obj)) {
|
||||
mule_id = PyInt_AsLong(mule_obj);
|
||||
if (mule_id < 0 && mule_id > uwsgi.mules_cnt) {
|
||||
return PyErr_Format(PyExc_ValueError, "invalid mule number");
|
||||
}
|
||||
if (mule_id == 0) {
|
||||
fd = uwsgi.shared->mule_queue_pipe[0];
|
||||
}
|
||||
else {
|
||||
fd = uwsgi.mules[mule_id-1].queue_pipe[0];
|
||||
}
|
||||
}
|
||||
else {
|
||||
return PyErr_Format(PyExc_ValueError, "invalid mule");
|
||||
}
|
||||
|
||||
if (fd > -1) {
|
||||
len = write(fd, message, message_len);
|
||||
if (len < 0) {
|
||||
uwsgi_error("write()");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_mule_get_msg(PyObject * self, PyObject * args) {
|
||||
PyObject *py_uwsgi_mule_get_msg(PyObject * self, PyObject * args, PyObject *kwargs) {
|
||||
|
||||
ssize_t len;
|
||||
ssize_t len = 0;
|
||||
// this buffer will be configurable
|
||||
char message[65536];
|
||||
char *message;
|
||||
struct pollfd *mulepoll;
|
||||
int count = 4;
|
||||
int farms_count = 0;
|
||||
uint8_t uwsgi_signal;
|
||||
PyObject *manage_signals = NULL;
|
||||
PyObject *manage_farms = NULL;
|
||||
int buffer_size = 65536;
|
||||
int timeout = -1;
|
||||
int i;
|
||||
|
||||
static char *kwlist[] = {"signals", "buffer_size", "timeout", "farms", NULL};
|
||||
|
||||
if (uwsgi.muleid == 0) {
|
||||
return PyErr_Format(PyExc_ValueError, "you can receive mule messages only in a mule !!!");
|
||||
}
|
||||
UWSGI_RELEASE_GIL;
|
||||
len = read(uwsgi.mules[uwsgi.muleid-1].queue_pipe[1], message, 65536);
|
||||
UWSGI_GET_GIL;
|
||||
if (len <= 0) {
|
||||
uwsgi_error("read()");
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
|
||||
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|OOii:mule_get_msg", kwlist, &manage_signals, &manage_farms, &buffer_size, &timeout)) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
return PyString_FromStringAndSize(message, len);
|
||||
if (manage_signals == Py_None || manage_signals == Py_False) {
|
||||
count = 2;
|
||||
}
|
||||
|
||||
if (manage_farms == Py_None || manage_farms == Py_False) {
|
||||
goto next;
|
||||
}
|
||||
|
||||
for(i=0;i<uwsgi.farms_cnt;i++) {
|
||||
if (uwsgi_farm_has_mule(&uwsgi.farms[i], uwsgi.muleid)) farms_count++;
|
||||
}
|
||||
|
||||
|
||||
next:
|
||||
|
||||
UWSGI_RELEASE_GIL;
|
||||
if (timeout > -1) timeout = timeout*1000;
|
||||
|
||||
message = uwsgi_malloc(buffer_size);
|
||||
mulepoll = uwsgi_malloc(sizeof(struct pollfd) * (count+farms_count));
|
||||
|
||||
mulepoll[0].fd = uwsgi.mules[uwsgi.muleid-1].queue_pipe[1];
|
||||
mulepoll[0].events = POLLIN;
|
||||
mulepoll[1].fd = uwsgi.shared->mule_queue_pipe[1];
|
||||
mulepoll[1].events = POLLIN;
|
||||
if (count > 2) {
|
||||
mulepoll[2].fd = uwsgi.signal_socket;
|
||||
mulepoll[2].events = POLLIN;
|
||||
mulepoll[3].fd = uwsgi.my_signal_socket;
|
||||
mulepoll[3].events = POLLIN;
|
||||
}
|
||||
|
||||
for(i=0;i<farms_count;i++) {
|
||||
mulepoll[count+i].fd = uwsgi.farms[i].queue_pipe[1];
|
||||
mulepoll[count+i].events = POLLIN;
|
||||
}
|
||||
|
||||
int ret = poll(mulepoll, count+farms_count, timeout);
|
||||
if (ret <= 0) {
|
||||
uwsgi_error("poll");
|
||||
}
|
||||
else {
|
||||
if (mulepoll[0].revents & POLLIN) {
|
||||
len = read(uwsgi.mules[uwsgi.muleid-1].queue_pipe[1], message, buffer_size);
|
||||
}
|
||||
else if (mulepoll[1].revents & POLLIN) {
|
||||
len = read(uwsgi.shared->mule_queue_pipe[1], message, buffer_size);
|
||||
}
|
||||
else {
|
||||
if (count > 2) {
|
||||
int interesting_fd = -1;
|
||||
if (mulepoll[2].revents & POLLIN) {
|
||||
interesting_fd = mulepoll[2].fd;
|
||||
}
|
||||
else if (mulepoll[3].revents & POLLIN) {
|
||||
interesting_fd = mulepoll[3].fd;
|
||||
}
|
||||
|
||||
if (interesting_fd > -1) {
|
||||
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);
|
||||
end_me(0);
|
||||
}
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log_verbose("master sent signal %d to mule %d\n", uwsgi_signal, uwsgi.muleid);
|
||||
#endif
|
||||
if (uwsgi_signal_handler(uwsgi_signal)) {
|
||||
uwsgi_log_verbose("error managing signal %d on mule %d\n", uwsgi_signal, uwsgi.mywid);
|
||||
}
|
||||
goto clear;
|
||||
}
|
||||
}
|
||||
|
||||
for(i=0;i<farms_count;i++) {
|
||||
if (mulepoll[count+i].revents & POLLIN) {
|
||||
len = read(mulepoll[count+i].fd, message, buffer_size);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
UWSGI_GET_GIL;
|
||||
if (len < 0) {
|
||||
uwsgi_error("read()");
|
||||
goto clear2;
|
||||
}
|
||||
|
||||
PyObject *msg = PyString_FromStringAndSize(message, len);
|
||||
free(message);
|
||||
free(mulepoll);
|
||||
return msg;
|
||||
clear:
|
||||
UWSGI_GET_GIL;
|
||||
clear2:
|
||||
free(message);
|
||||
free(mulepoll);
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_farm_get_msg(PyObject * self, PyObject * args) {
|
||||
@@ -1204,7 +1343,7 @@ PyObject *py_uwsgi_farm_get_msg(PyObject * self, PyObject * args) {
|
||||
struct pollfd *farmpoll;
|
||||
|
||||
if (uwsgi.muleid == 0) {
|
||||
return PyErr_Format(PyExc_ValueError, "you can receive mule messages only in a mule !!!");
|
||||
return PyErr_Format(PyExc_ValueError, "you can receive farm messages only in a mule !!!");
|
||||
}
|
||||
UWSGI_RELEASE_GIL;
|
||||
for(i=0;i<uwsgi.farms_cnt;i++) {
|
||||
@@ -3001,7 +3140,7 @@ static PyMethodDef uwsgi_advanced_methods[] = {
|
||||
|
||||
{"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, ""},
|
||||
{"mule_get_msg", (PyCFunction) py_uwsgi_mule_get_msg, METH_VARARGS|METH_KEYWORDS, ""},
|
||||
{"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, ""},
|
||||
|
||||
@@ -275,7 +275,5 @@ char *uwsgi_pythonize(char *);
|
||||
#define uwsgi_pyexit {PyErr_Print();exit(1);}
|
||||
|
||||
#ifdef __linux__
|
||||
#ifndef PYTHREE
|
||||
int uwsgi_init_symbol_import(void);
|
||||
#endif
|
||||
#endif
|
||||
|
||||
@@ -0,0 +1,188 @@
|
||||
#include "../../uwsgi.h"
|
||||
|
||||
extern struct uwsgi_server uwsgi;
|
||||
|
||||
#define RRDTOOL_OPT_BASE 177000
|
||||
#define RRDTOOL_OPT_RRDTOOL RRDTOOL_OPT_BASE+1
|
||||
#define RRDTOOL_OPT_RRDTOOL_MAX_DS RRDTOOL_OPT_BASE+2
|
||||
|
||||
struct uwsgi_rrdtool {
|
||||
void *lib;
|
||||
int (*create)(int, char **);
|
||||
int (*update)(int, char **);
|
||||
struct uwsgi_string_list *rrd;
|
||||
int max_ds;
|
||||
|
||||
char *update_area;
|
||||
} u_rrd;
|
||||
|
||||
struct option rrdtool_options[] = {
|
||||
{"rrdtool", required_argument, 0, RRDTOOL_OPT_RRDTOOL},
|
||||
{"rrdtool-max-ds", required_argument, 0, RRDTOOL_OPT_RRDTOOL_MAX_DS},
|
||||
{0, 0, 0, 0},
|
||||
|
||||
};
|
||||
|
||||
|
||||
int rrdtool_init() {
|
||||
|
||||
u_rrd.lib = dlopen("librrd.so", RTLD_LAZY);
|
||||
if (!u_rrd.lib) return -1;
|
||||
|
||||
u_rrd.create = dlsym(u_rrd.lib, "rrd_create");
|
||||
if (!u_rrd.create) {
|
||||
dlclose(u_rrd.lib);
|
||||
return -1;
|
||||
}
|
||||
|
||||
u_rrd.update = dlsym(u_rrd.lib, "rrd_update");
|
||||
if (!u_rrd.update) {
|
||||
dlclose(u_rrd.lib);
|
||||
return -1;
|
||||
}
|
||||
|
||||
if (!u_rrd.max_ds) u_rrd.max_ds = 30;
|
||||
|
||||
uwsgi_log("*** RRDtool library available at %p ***\n", u_rrd.lib);
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
int rrdtool_opt(int i, char *optarg) {
|
||||
|
||||
switch(i) {
|
||||
case RRDTOOL_OPT_RRDTOOL:
|
||||
uwsgi.master_process = 1;
|
||||
uwsgi_string_new_list(&u_rrd.rrd, optarg);
|
||||
return 1;
|
||||
case RRDTOOL_OPT_RRDTOOL_MAX_DS:
|
||||
u_rrd.max_ds = atoi(optarg);
|
||||
return 1;
|
||||
}
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
void rrdtool_post_init() {
|
||||
|
||||
struct uwsgi_string_list *usl = u_rrd.rrd;
|
||||
char **argv;
|
||||
int i;
|
||||
|
||||
if (!u_rrd.lib || !u_rrd.create) return;
|
||||
|
||||
// do not waste time if no --rrdtool option is defiend
|
||||
if (!u_rrd.rrd) return;
|
||||
|
||||
if (uwsgi.numproc > u_rrd.max_ds) {
|
||||
uwsgi_log("!!! NOT ENOUGH SLOTS IN RRDTOOL DS TO HOST WORKERS DATA (increase them with --rrdtool-max-ds) !!!\n");
|
||||
dlclose(u_rrd.lib);
|
||||
return;
|
||||
}
|
||||
|
||||
// alloc space for DS_REQ + DS WORKER + RRA + create + filename
|
||||
argv = uwsgi_malloc( sizeof(char *) * (1 + u_rrd.max_ds + 4 + 1 +1));
|
||||
|
||||
argv[0] = "create";
|
||||
|
||||
argv[2] = "DS:requests:DERIVE:600:0:U";
|
||||
|
||||
// create DS for workers
|
||||
for(i=0;i<u_rrd.max_ds;i++) {
|
||||
int max_size = sizeof("DS:worker65536:DERIVE:600:0:U")+1;
|
||||
argv[3+i] = uwsgi_malloc( max_size );
|
||||
if (snprintf(argv[3+i], max_size, "DS:worker%d:DERIVE:600:0:U", i+1) < 25) {
|
||||
uwsgi_log("unable to create args for rrd_create()\n");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
// create RRA
|
||||
argv[3+u_rrd.max_ds] = "RRA:AVERAGE:0.5:1:288" ;
|
||||
argv[3+u_rrd.max_ds+1] = "RRA:AVERAGE:0.5:12:168" ;
|
||||
argv[3+u_rrd.max_ds+2] = "RRA:AVERAGE:0.5:288:31" ;
|
||||
argv[3+u_rrd.max_ds+3] = "RRA:AVERAGE:0.5:2016:52";
|
||||
|
||||
while(usl) {
|
||||
if (!uwsgi_file_exists(usl->value)) {
|
||||
argv[1] = usl->value;
|
||||
if (u_rrd.create((1 + u_rrd.max_ds + 4 + 1 +1), argv)) {
|
||||
uwsgi_error("rrd_create()");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
usl->value = realpath(usl->value, NULL);
|
||||
if (!usl->value) {
|
||||
uwsgi_error("realpath()");
|
||||
exit(1);
|
||||
}
|
||||
usl = usl->next;
|
||||
}
|
||||
|
||||
// free DS
|
||||
for(i=0;i<u_rrd.max_ds;i++) {
|
||||
free(argv[3+i]);
|
||||
}
|
||||
|
||||
free(argv);
|
||||
|
||||
//now allocate memory for updates
|
||||
u_rrd.update_area = uwsgi_malloc( 1+((1+sizeof(UMAX64_STR)) * (u_rrd.max_ds+1))+1 );
|
||||
memset(u_rrd.update_area, 0, 1+((1+sizeof(UMAX64_STR)) * (u_rrd.max_ds+1))+1 );
|
||||
|
||||
u_rrd.update_area[0] = 'N';
|
||||
|
||||
}
|
||||
|
||||
void rrdtool_master_cycle() {
|
||||
|
||||
static time_t last_update = 0;
|
||||
char *ptr;
|
||||
int rlen, i;
|
||||
char *argv[3];
|
||||
struct uwsgi_string_list *usl = u_rrd.rrd;
|
||||
|
||||
if (!u_rrd.lib || !u_rrd.create || !u_rrd.rrd) return ;
|
||||
|
||||
if (last_update == 0) last_update = time(NULL);
|
||||
|
||||
// update every 5 minutes
|
||||
if (uwsgi.current_time - last_update >= 300) {
|
||||
ptr = u_rrd.update_area+1;
|
||||
rlen = snprintf(ptr, 1+sizeof(UMAX64_STR), ":%llu", (unsigned long long )uwsgi.workers[0].requests);
|
||||
if (rlen < 2) return;
|
||||
ptr+=rlen;
|
||||
for(i=0;i<u_rrd.max_ds;i++) {
|
||||
if (i+1 <= uwsgi.numproc) {
|
||||
rlen = snprintf(ptr, 1+sizeof(UMAX64_STR), ":%llu", (unsigned long long )uwsgi.workers[1+i].requests);
|
||||
if (rlen < 2) return;
|
||||
}
|
||||
else {
|
||||
memcpy(ptr, ":U", 2);
|
||||
rlen = 2;
|
||||
}
|
||||
ptr+=rlen;
|
||||
}
|
||||
last_update = uwsgi.current_time;
|
||||
argv[0] = "update";
|
||||
argv[2] = u_rrd.update_area;
|
||||
while(usl) {
|
||||
argv[1] = usl->value;
|
||||
if (u_rrd.update(3, argv)) {
|
||||
uwsgi_log_verbose("ERROR: rrd_update(\"%s\", \"%s\")\n", argv[1], argv[2]);
|
||||
}
|
||||
usl = usl->next;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
struct uwsgi_plugin rrdtool_plugin = {
|
||||
|
||||
.options = rrdtool_options,
|
||||
.manage_opt = rrdtool_opt,
|
||||
|
||||
.master_cycle = rrdtool_master_cycle,
|
||||
|
||||
.post_init = rrdtool_post_init,
|
||||
.init = rrdtool_init,
|
||||
};
|
||||
@@ -0,0 +1,7 @@
|
||||
|
||||
NAME='rrdtool'
|
||||
CFLAGS = []
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
|
||||
GCC_LIST = ['rrdtool']
|
||||
@@ -1,7 +1,7 @@
|
||||
#ifdef UWSGI_PCRE
|
||||
#include "uwsgi.h"
|
||||
|
||||
void uwsgi_regexp_build(char *re, pcre **pattern, pcre_extra **pattern_extra) {
|
||||
int uwsgi_regexp_build(char *re, pcre **pattern, pcre_extra **pattern_extra) {
|
||||
|
||||
const char *errstr;
|
||||
int erroff;
|
||||
@@ -9,15 +9,17 @@ void uwsgi_regexp_build(char *re, pcre **pattern, pcre_extra **pattern_extra) {
|
||||
*pattern = pcre_compile( (const char *)re, 0, &errstr, &erroff, NULL);
|
||||
if (!*pattern) {
|
||||
uwsgi_log("pcre error: %s at offset %d\n", errstr, erroff);
|
||||
exit(1);
|
||||
return -1;
|
||||
}
|
||||
|
||||
*pattern_extra = (pcre_extra *) pcre_study((const pcre*)*pattern, 0, &errstr);
|
||||
if (!*pattern_extra) {
|
||||
pcre_free(*pattern);
|
||||
uwsgi_log("pcre (study) error: %s\n", errstr);
|
||||
exit(1);
|
||||
return -1;
|
||||
}
|
||||
|
||||
return 0;
|
||||
|
||||
}
|
||||
|
||||
@@ -26,32 +28,4 @@ int uwsgi_regexp_match(pcre *pattern, pcre_extra *pattern_extra, char *subject,
|
||||
return pcre_exec((const pcre*)pattern, (const pcre_extra *)pattern_extra, subject, length, 0, 0, NULL, 0 );
|
||||
}
|
||||
|
||||
/*
|
||||
|
||||
void uwsgi_regexp_match(regexp, what) {
|
||||
|
||||
int ret,i;
|
||||
|
||||
for(i=0;i<uwsgi->nroutes;i++) {
|
||||
|
||||
ret = pcre_exec(ur->pattern, ur->pattern_extra, wsgi_req->path_info, wsgi_req->path_info_len, 0, 0, wsgi_req->ovector, (ur->args+1)*3 );
|
||||
|
||||
if (ret >= 0) {
|
||||
if (ur->action) {
|
||||
ur->action(uwsgi, wsgi_req, ur);
|
||||
}
|
||||
else {
|
||||
uwsgi_route_action_wsgi(uwsgi, wsgi_req, ur);
|
||||
}
|
||||
}
|
||||
|
||||
// TODO check for errors if < 0 && != NO_MATCH
|
||||
}
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
*/
|
||||
|
||||
|
||||
#endif
|
||||
|
||||
@@ -9,6 +9,12 @@ from setuptools.command.install import install
|
||||
from setuptools.command.install_lib import install_lib
|
||||
from setuptools.command.build_ext import build_ext
|
||||
|
||||
"""
|
||||
This is a hack allowing you installing uWSGI and uwsgidecorators via pip and easy_install
|
||||
"""
|
||||
|
||||
uwsgi_compiled = False
|
||||
|
||||
def get_profile():
|
||||
profile = os.environ.get('UWSGI_PROFILE','buildconf/default.ini')
|
||||
if not profile.endswith('.ini'):
|
||||
@@ -35,26 +41,35 @@ def patch_bin_path(cmd, conf):
|
||||
class uWSGIBuilder(build_ext):
|
||||
|
||||
def run(self):
|
||||
conf = uc.uConf(get_profile())
|
||||
patch_bin_path(self, conf)
|
||||
uc.build_uwsgi( conf )
|
||||
global uwsgi_compiled
|
||||
if not uwsgi_compiled:
|
||||
conf = uc.uConf(get_profile())
|
||||
patch_bin_path(self, conf)
|
||||
uc.build_uwsgi( conf )
|
||||
uwsgi_compiled = True
|
||||
|
||||
|
||||
class uWSGIInstall(install):
|
||||
|
||||
def run(self):
|
||||
|
||||
conf = uc.uConf(get_profile())
|
||||
patch_bin_path(self, conf)
|
||||
uc.build_uwsgi( conf )
|
||||
global uwsgi_compiled
|
||||
if not uwsgi_compiled:
|
||||
conf = uc.uConf(get_profile())
|
||||
patch_bin_path(self, conf)
|
||||
uc.build_uwsgi( conf )
|
||||
uwsgi_compiled = True
|
||||
install.run(self)
|
||||
|
||||
class uWSGIInstallLib(install_lib):
|
||||
|
||||
def run(self):
|
||||
conf = uc.uConf(get_profile())
|
||||
patch_bin_path(self, conf)
|
||||
uc.build_uwsgi( conf )
|
||||
global uwsgi_compiled
|
||||
if not uwsgi_compiled:
|
||||
conf = uc.uConf(get_profile())
|
||||
patch_bin_path(self, conf)
|
||||
uc.build_uwsgi( conf )
|
||||
uwsgi_compiled = True
|
||||
install_lib.run(self)
|
||||
|
||||
class uWSGIDistribution(Distribution):
|
||||
|
||||
@@ -71,6 +86,7 @@ setup(name='uWSGI',
|
||||
author_email='info@unbit.it',
|
||||
url='http://projects.unbit.it/uwsgi/',
|
||||
license='GPL2',
|
||||
py_modules = ['uwsgidecorators'],
|
||||
distclass = uWSGIDistribution,
|
||||
)
|
||||
|
||||
|
||||
@@ -271,7 +271,18 @@ void uwsgi_route_signal(uint8_t sig) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
else if (!strncmp(use->receiver, "farm_", 5)) {
|
||||
char *name = use->receiver+5;
|
||||
struct uwsgi_farm *uf = get_farm_by_name(name);
|
||||
if (!uf) {
|
||||
uwsgi_log("unknown farm: %s\n", name);
|
||||
return;
|
||||
}
|
||||
if (write(uf->signal_pipe[0], &sig, 1) != 1) {
|
||||
uwsgi_error("write()");
|
||||
uwsgi_log("could not deliver signal %d to farm %d (%s)\n", sig, uf->id, uf->name);
|
||||
}
|
||||
}
|
||||
else if (!strncmp(use->receiver, "farm", 4)) {
|
||||
i = atoi(use->receiver+4);
|
||||
if (i > uwsgi.farms_cnt || i <= 0) {
|
||||
|
||||
+274
@@ -0,0 +1,274 @@
|
||||
#include "uwsgi.h"
|
||||
|
||||
/*
|
||||
|
||||
subscription subsystem
|
||||
|
||||
each subscription slot is as an auto-optmizing linked list. Originally it was a uwsgi_dict (now removed from uWSGI) but this
|
||||
would have not be able to support regexp as keys.
|
||||
|
||||
each slot has another circular linked list containing the nodes names
|
||||
|
||||
the structure and system is very similar to uwsgi_dyn_dict already used by the mime type parser
|
||||
|
||||
This system is not mean to run on shared memory. If you have multiple processes for the same app, you have to create
|
||||
a new subscriptions slot list.
|
||||
|
||||
*/
|
||||
|
||||
struct uwsgi_subscribe_slot *uwsgi_get_subscribe_slot(struct uwsgi_subscribe_slot **slot, char *key, uint16_t keylen, int regexp) {
|
||||
|
||||
struct uwsgi_subscribe_slot *current_slot = *slot;
|
||||
|
||||
if (keylen > 0xff) return NULL;
|
||||
|
||||
while(current_slot) {
|
||||
#ifdef UWSGI_PCRE
|
||||
if (regexp) {
|
||||
if (uwsgi_regexp_match(current_slot->pattern, current_slot->pattern_extra, key, keylen) >= 0) {
|
||||
return current_slot;
|
||||
}
|
||||
}
|
||||
else {
|
||||
#endif
|
||||
if (!uwsgi_strncmp(key, keylen, current_slot->key, current_slot->keylen)) {
|
||||
// auto optimization
|
||||
if (current_slot->prev) {
|
||||
if (current_slot->hits > current_slot->prev->hits) {
|
||||
struct uwsgi_subscribe_slot *slot_parent = current_slot->prev->prev, *slot_prev = current_slot->prev;
|
||||
if (slot_parent) {
|
||||
slot_parent->next = current_slot;
|
||||
}
|
||||
else {
|
||||
*slot = current_slot;
|
||||
}
|
||||
|
||||
slot_prev->prev = current_slot;
|
||||
slot_prev->next = current_slot->next;
|
||||
|
||||
current_slot->next = slot_prev;
|
||||
current_slot->prev = slot_parent;
|
||||
|
||||
}
|
||||
}
|
||||
return current_slot;
|
||||
}
|
||||
#ifdef UWSGI_PCRE
|
||||
}
|
||||
#endif
|
||||
current_slot = current_slot->next;
|
||||
}
|
||||
|
||||
return NULL;
|
||||
}
|
||||
|
||||
struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slot **slot, char *key, uint16_t keylen, int regexp) {
|
||||
|
||||
if (keylen > 0xff) return NULL;
|
||||
|
||||
struct uwsgi_subscribe_slot *current_slot = uwsgi_get_subscribe_slot(slot, key, keylen, regexp);
|
||||
uint64_t rr_pos = 0;
|
||||
|
||||
if (current_slot) {
|
||||
// node found, move up in the list increasing hits
|
||||
current_slot->hits++;
|
||||
struct uwsgi_subscribe_node *node = current_slot->nodes;
|
||||
while(node) {
|
||||
if (rr_pos == current_slot->rr) {
|
||||
current_slot->rr++;
|
||||
return node;
|
||||
}
|
||||
node = node->next;
|
||||
rr_pos++;
|
||||
}
|
||||
current_slot->rr = 0;
|
||||
return current_slot->nodes;
|
||||
}
|
||||
|
||||
return NULL;
|
||||
}
|
||||
|
||||
void uwsgi_remove_subscribe_node(struct uwsgi_subscribe_slot **slot, struct uwsgi_subscribe_node *node) {
|
||||
|
||||
struct uwsgi_subscribe_node *a_node;
|
||||
struct uwsgi_subscribe_slot *node_slot = node->slot;
|
||||
struct uwsgi_subscribe_slot *prev_slot = node_slot->prev;
|
||||
struct uwsgi_subscribe_slot *next_slot = node_slot->next;
|
||||
|
||||
// over-engineering to avoid race conditions
|
||||
node->len = 0;
|
||||
|
||||
if (node == node_slot->nodes) {
|
||||
node_slot->nodes = node->next;
|
||||
}
|
||||
else {
|
||||
a_node = node_slot->nodes;
|
||||
while(a_node) {
|
||||
if (a_node->next == node) {
|
||||
a_node->next = node->next;
|
||||
break;
|
||||
}
|
||||
a_node = a_node->next;
|
||||
}
|
||||
}
|
||||
|
||||
free(node);
|
||||
// no more nodes, remove the slot too
|
||||
if (node_slot->nodes == NULL) {
|
||||
|
||||
if (prev_slot) {
|
||||
prev_slot->next = next_slot;
|
||||
}
|
||||
if (next_slot) {
|
||||
next_slot->prev = prev_slot;
|
||||
}
|
||||
|
||||
#ifdef UWSGI_PCRE
|
||||
if (node_slot->pattern) {
|
||||
pcre_free(node_slot->pattern);
|
||||
}
|
||||
if (node_slot->pattern_extra) {
|
||||
pcre_free(node_slot->pattern_extra);
|
||||
}
|
||||
#endif
|
||||
|
||||
free(node_slot);
|
||||
// am i the only slot ?
|
||||
if (!prev_slot && !next_slot) {
|
||||
*slot = NULL;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slot **slot, struct uwsgi_subscribe_req *usr, int regexp) {
|
||||
|
||||
struct uwsgi_subscribe_slot *current_slot = uwsgi_get_subscribe_slot(slot, usr->key, usr->keylen, 0), *old_slot = NULL, *a_slot;
|
||||
struct uwsgi_subscribe_node *node, *old_node = NULL;
|
||||
|
||||
if (usr->address_len > 0xff) return NULL;
|
||||
|
||||
if (current_slot) {
|
||||
node = current_slot->nodes;
|
||||
while(node) {
|
||||
if (!uwsgi_strncmp(node->name, node->len, usr->address, usr->address_len)) {
|
||||
node->last_check = time(NULL);
|
||||
return node;
|
||||
}
|
||||
old_node = node;
|
||||
node = node->next;
|
||||
}
|
||||
|
||||
node = uwsgi_malloc(sizeof(struct uwsgi_subscribe_node));
|
||||
node->len = usr->address_len;
|
||||
node->modifier1 = usr->modifier1;
|
||||
node->modifier2 = usr->modifier2;
|
||||
node->last_check = time(NULL);
|
||||
node->slot = current_slot;
|
||||
memcpy(node->name, usr->address, usr->address_len);
|
||||
if (old_node) {
|
||||
old_node->next = node;
|
||||
}
|
||||
node->next = NULL;
|
||||
uwsgi_log("[uwsgi-subscription] %.*s => new node: %.*s\n", usr->keylen, usr->key, usr->address_len, usr->address);
|
||||
return node;
|
||||
}
|
||||
else {
|
||||
|
||||
current_slot = uwsgi_malloc(sizeof(struct uwsgi_subscribe_slot));
|
||||
current_slot->keylen = usr->keylen;
|
||||
memcpy(current_slot->key, usr->key, usr->keylen);
|
||||
current_slot->key[usr->keylen] = 0;
|
||||
current_slot->hits = 0;
|
||||
current_slot->rr = 0;
|
||||
|
||||
#ifdef UWSGI_PCRE
|
||||
current_slot->pattern = NULL;
|
||||
current_slot->pattern_extra = NULL;
|
||||
if (regexp) {
|
||||
if (uwsgi_regexp_build(current_slot->key, ¤t_slot->pattern, ¤t_slot->pattern_extra)) {
|
||||
free(current_slot);
|
||||
return NULL;
|
||||
}
|
||||
}
|
||||
#endif
|
||||
|
||||
current_slot->nodes = uwsgi_malloc(sizeof(struct uwsgi_subscribe_node));
|
||||
current_slot->nodes->slot = current_slot;
|
||||
current_slot->nodes->len = usr->address_len;
|
||||
current_slot->nodes->modifier1 = usr->modifier1;
|
||||
current_slot->nodes->modifier2 = usr->modifier2;
|
||||
memcpy(current_slot->nodes->name, usr->address, usr->address_len);
|
||||
current_slot->nodes->last_check = time(NULL);
|
||||
|
||||
current_slot->nodes->next = NULL;
|
||||
|
||||
#ifdef UWSGI_PCRE
|
||||
// if key is a regexp, order it by keylen
|
||||
if (regexp) {
|
||||
old_slot = NULL;
|
||||
a_slot = *slot;
|
||||
while(a_slot) {
|
||||
if (a_slot->keylen > current_slot->keylen) {
|
||||
old_slot = a_slot;
|
||||
break;
|
||||
}
|
||||
a_slot = a_slot->next;
|
||||
}
|
||||
|
||||
if (old_slot) {
|
||||
current_slot->prev = old_slot->prev;
|
||||
old_slot->prev = current_slot;
|
||||
if (current_slot->prev) {
|
||||
old_slot->prev->next = current_slot;
|
||||
}
|
||||
|
||||
current_slot->next = old_slot;
|
||||
}
|
||||
else {
|
||||
a_slot = *slot;
|
||||
while(a_slot) {
|
||||
old_slot = a_slot;
|
||||
a_slot = a_slot->next;
|
||||
}
|
||||
|
||||
|
||||
if (old_slot) {
|
||||
old_slot->next = current_slot;
|
||||
}
|
||||
|
||||
current_slot->prev = old_slot;
|
||||
current_slot->next = NULL;
|
||||
}
|
||||
}
|
||||
else {
|
||||
#endif
|
||||
a_slot = *slot;
|
||||
while(a_slot) {
|
||||
old_slot = a_slot;
|
||||
a_slot = a_slot->next;
|
||||
}
|
||||
|
||||
|
||||
if (old_slot) {
|
||||
old_slot->next = current_slot;
|
||||
}
|
||||
|
||||
current_slot->prev = old_slot;
|
||||
current_slot->next = NULL;
|
||||
|
||||
#ifdef UWSGI_PCRE
|
||||
}
|
||||
#endif
|
||||
|
||||
if (!*slot || current_slot->prev == NULL) {
|
||||
*slot = current_slot;
|
||||
}
|
||||
|
||||
uwsgi_log("[uwsgi-subscription] new pool: %.*s\n", usr->keylen, usr->key);
|
||||
uwsgi_log("[uwsgi-subscription] %.*s => new node: %.*s\n", usr->keylen, usr->key, usr->address_len, usr->address);
|
||||
return current_slot->nodes;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -246,8 +246,8 @@ void logto(char *logfile) {
|
||||
uwsgi_error_open(logfile);
|
||||
exit(1);
|
||||
}
|
||||
#ifdef UWSGI_UDP
|
||||
uwsgi.logfile = logfile;
|
||||
#ifdef UWSGI_UDP
|
||||
}
|
||||
#endif
|
||||
|
||||
@@ -271,61 +271,6 @@ void logto(char *logfile) {
|
||||
#ifdef UWSGI_ZEROMQ
|
||||
void log_zeromq(char *node) {
|
||||
|
||||
if (socketpair(AF_UNIX, SOCK_DGRAM, 0, uwsgi.shared->worker_log_pipe)) {
|
||||
uwsgi_error("socketpair()\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
uwsgi_socket_nb(uwsgi.shared->worker_log_pipe[0]);
|
||||
uwsgi_socket_nb(uwsgi.shared->worker_log_pipe[1]);
|
||||
#ifdef UWSGI_DEBUG
|
||||
int so_bufsize;
|
||||
socklen_t so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[0], SOL_SOCKET, SO_RCVBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("master logger SO_RCVBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[0], SOL_SOCKET, SO_SNDBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("master logger SO_SNDBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[1], SOL_SOCKET, SO_RCVBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("worker logger SO_RCVBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[1], SOL_SOCKET, SO_SNDBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("worker logger SO_SNDBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
#endif
|
||||
|
||||
if (uwsgi.shared->worker_log_pipe[1] != 1) {
|
||||
if (dup2(uwsgi.shared->worker_log_pipe[1], 1) < 0) {
|
||||
uwsgi_error("dup2()");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
if (dup2(1, 2) < 0) {
|
||||
uwsgi_error("dup2()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
void *ctx = zmq_init(1);
|
||||
if (ctx == NULL) {
|
||||
uwsgi_error("zmq_init()");
|
||||
@@ -365,10 +310,6 @@ void log_socket(char *socket_name) {
|
||||
uwsgi_nuclear_blast();
|
||||
}
|
||||
|
||||
// create log connection with the master
|
||||
|
||||
create_logpipe();
|
||||
|
||||
}
|
||||
|
||||
void create_logpipe(void) {
|
||||
@@ -403,71 +344,12 @@ void log_syslog(char *syslog_opts) {
|
||||
syslog_opts = "uwsgi";
|
||||
}
|
||||
|
||||
if (socketpair(AF_UNIX, SOCK_DGRAM, 0, uwsgi.shared->worker_log_pipe)) {
|
||||
uwsgi_error("socketpair()\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
uwsgi_socket_nb(uwsgi.shared->worker_log_pipe[0]);
|
||||
uwsgi_socket_nb(uwsgi.shared->worker_log_pipe[1]);
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
int so_bufsize;
|
||||
uwsgi_log("log pipe %d %d\n", uwsgi.shared->worker_log_pipe[0], uwsgi.shared->worker_log_pipe[1]);
|
||||
socklen_t so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[0], SOL_SOCKET, SO_RCVBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("master logger SO_RCVBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[0], SOL_SOCKET, SO_SNDBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("master logger SO_SNDBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[1], SOL_SOCKET, SO_RCVBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("worker logger SO_RCVBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
so_bufsize_len = sizeof(int);
|
||||
if (getsockopt(uwsgi.shared->worker_log_pipe[1], SOL_SOCKET, SO_SNDBUF, &so_bufsize, &so_bufsize_len)) {
|
||||
uwsgi_error("getsockopt()");
|
||||
}
|
||||
else {
|
||||
uwsgi_debug("worker logger SO_SNDBUF size: %d\n", so_bufsize);
|
||||
}
|
||||
|
||||
#endif
|
||||
|
||||
if (uwsgi.shared->worker_log_pipe[1] != 1) {
|
||||
if (dup2(uwsgi.shared->worker_log_pipe[1], 1) < 0) {
|
||||
uwsgi_error("dup2()");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("opening syslog\n");
|
||||
#endif
|
||||
|
||||
if (dup2(1, 2) < 0) {
|
||||
uwsgi_error("dup2()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
openlog(syslog_opts, 0, LOG_DAEMON);
|
||||
|
||||
|
||||
}
|
||||
|
||||
char *uwsgi_get_cwd() {
|
||||
@@ -2684,6 +2566,22 @@ struct uwsgi_dyn_dict *uwsgi_dyn_dict_new(struct uwsgi_dyn_dict **dd, char *key,
|
||||
return uwsgi_dd;
|
||||
}
|
||||
|
||||
void uwsgi_dyn_dict_del(struct uwsgi_dyn_dict *item) {
|
||||
|
||||
struct uwsgi_dyn_dict *prev = item->prev;
|
||||
struct uwsgi_dyn_dict *next = item->next;
|
||||
|
||||
if (prev) {
|
||||
prev->next = next;
|
||||
}
|
||||
|
||||
if (next) {
|
||||
next->prev = prev;
|
||||
}
|
||||
|
||||
free(item);
|
||||
}
|
||||
|
||||
|
||||
struct uwsgi_string_list *uwsgi_string_new_list(struct uwsgi_string_list **list, char *value) {
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
|
||||
_socket_nb(uwsgi.shared->worker_log_pipe[1]
|
||||
*** uWSGI ***
|
||||
|
||||
Copyright (C) 2009-2011 Unbit S.a.s. <info@unbit.it>
|
||||
@@ -70,6 +70,7 @@ static struct option long_base_options[] = {
|
||||
{"abstract-socket", no_argument, 0, 'a'},
|
||||
{"chmod-socket", optional_argument, 0, 'C'},
|
||||
{"chown-socket", required_argument, 0, LONG_ARGS_CHOWN_SOCKET},
|
||||
{"umask", required_argument, 0, LONG_ARGS_UMASK},
|
||||
#ifdef __linux__
|
||||
{"freebind", no_argument, &uwsgi.freebind, 1},
|
||||
#endif
|
||||
@@ -80,6 +81,7 @@ static struct option long_base_options[] = {
|
||||
#endif
|
||||
{"auto-procname", no_argument, &uwsgi.auto_procname, 1},
|
||||
{"procname-prefix", required_argument, 0, LONG_ARGS_PROCNAME_PREFIX},
|
||||
{"procname-prefix-spaced", required_argument, 0, LONG_ARGS_PROCNAME_PREFIX_SP},
|
||||
{"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},
|
||||
@@ -126,6 +128,7 @@ static struct option long_base_options[] = {
|
||||
{"spooler-chdir", required_argument, 0, LONG_ARGS_SPOOLER_CHDIR},
|
||||
#endif
|
||||
{"mule", optional_argument, 0, LONG_ARGS_MULE},
|
||||
{"mules", required_argument, 0, LONG_ARGS_MULES},
|
||||
{"farm", required_argument, 0, LONG_ARGS_FARM},
|
||||
{"disable-logging", no_argument, 0, 'L'},
|
||||
|
||||
@@ -992,6 +995,8 @@ int main(int argc, char *argv[], char *envp[]) {
|
||||
uwsgi.cache_server_fd = -1;
|
||||
uwsgi.stats_fd = -1;
|
||||
|
||||
uwsgi.original_log_fd = -1;
|
||||
|
||||
uwsgi.emperor_fd_config = -1;
|
||||
uwsgi.emperor_throttle = 1000;
|
||||
uwsgi.emperor_pid = -1;
|
||||
@@ -1020,6 +1025,12 @@ int main(int argc, char *argv[], char *envp[]) {
|
||||
uwsgi.shared->mule_signal_pipe[0] = -1;
|
||||
uwsgi.shared->mule_signal_pipe[1] = -1;
|
||||
|
||||
uwsgi.shared->mule_queue_pipe[0] = -1;
|
||||
uwsgi.shared->mule_queue_pipe[1] = -1;
|
||||
|
||||
uwsgi.shared->worker_log_pipe[0] = -1;
|
||||
uwsgi.shared->worker_log_pipe[1] = -1;
|
||||
|
||||
uwsgi.mime_file = "/etc/mime.types";
|
||||
|
||||
|
||||
@@ -1180,6 +1191,9 @@ int main(int argc, char *argv[], char *envp[]) {
|
||||
else if (!strcmp(lazy + strlen(lazy) - 3, ".js")) {
|
||||
uwsgi.json = lazy;
|
||||
}
|
||||
else if (!strcmp(lazy + strlen(lazy) - 5, ".json")) {
|
||||
uwsgi.json = lazy;
|
||||
}
|
||||
#endif
|
||||
#ifdef UWSGI_SQLITE3
|
||||
else if (!strcmp(lazy + strlen(lazy) - 3, ".db")) {
|
||||
@@ -1291,6 +1305,9 @@ int main(int argc, char *argv[], char *envp[]) {
|
||||
if (!strcmp(uct->filename + strlen(uct->filename) - 3, ".js")) {
|
||||
uwsgi_json_config(uct->filename, uwsgi.magic_table);
|
||||
}
|
||||
if (!strcmp(uct->filename + strlen(uct->filename) - 5, ".json")) {
|
||||
uwsgi_json_config(uct->filename, uwsgi.magic_table);
|
||||
}
|
||||
#endif
|
||||
#ifdef UWSGI_SQLITE3
|
||||
if (!strcmp(uct->filename + strlen(uct->filename) - 3, ".db")) {
|
||||
@@ -1319,6 +1336,11 @@ int main(int argc, char *argv[], char *envp[]) {
|
||||
|
||||
uwsgi_configure();
|
||||
|
||||
if (uwsgi.log_master) {
|
||||
uwsgi.original_log_fd = dup(1);
|
||||
create_logpipe();
|
||||
}
|
||||
|
||||
/* uWSGI IS CONFIGURED !!! */
|
||||
|
||||
if (uwsgi.build_mime_dict) {
|
||||
@@ -1764,7 +1786,7 @@ int uwsgi_start(void *v_argv) {
|
||||
|
||||
// application generic lock
|
||||
uwsgi.user_lock = uwsgi_malloc(sizeof(void *) * (uwsgi.locks+1));
|
||||
for(i=0;i<uwsgi.locks;i++) {
|
||||
for(i=0;i<uwsgi.locks+1;i++) {
|
||||
uwsgi.user_lock[i] = uwsgi_mmap_shared_lock();
|
||||
uwsgi_lock_init(uwsgi.user_lock[i]);
|
||||
}
|
||||
@@ -1920,6 +1942,21 @@ int uwsgi_start(void *v_argv) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (uwsgi.shared->worker_log_pipe[0] > -1) {
|
||||
if (j == uwsgi.shared->worker_log_pipe[0])
|
||||
continue;
|
||||
}
|
||||
|
||||
if (uwsgi.shared->worker_log_pipe[1] > -1) {
|
||||
if (j == uwsgi.shared->worker_log_pipe[1])
|
||||
continue;
|
||||
}
|
||||
|
||||
if (uwsgi.original_log_fd > -1) {
|
||||
if (j == uwsgi.original_log_fd)
|
||||
continue;
|
||||
}
|
||||
|
||||
if (uwsgi.cache_server && uwsgi.cache_server_fd != -1) {
|
||||
if (j == uwsgi.cache_server_fd)
|
||||
continue;
|
||||
@@ -2243,6 +2280,11 @@ skipzero:
|
||||
exit(1);
|
||||
}
|
||||
|
||||
if (socketpair(AF_UNIX, SOCK_DGRAM, 0, uwsgi.shared->mule_queue_pipe)) {
|
||||
uwsgi_error("socketpair()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
for(i=0;i<uwsgi.mules_cnt;i++) {
|
||||
// create the socket pipe
|
||||
if (socketpair(AF_UNIX, SOCK_STREAM, 0, uwsgi.mules[i].signal_pipe)) {
|
||||
@@ -2491,12 +2533,12 @@ skipzero:
|
||||
|
||||
if (getpid() == masterpid && uwsgi.master_process == 1) {
|
||||
#ifdef UWSGI_AS_SHARED_LIBRARY
|
||||
int ml_ret = master_loop(uwsgi.orig_argv, uwsgi.environ);
|
||||
int ml_ret = master_loop(uwsgi.argv, uwsgi.environ);
|
||||
if (ml_ret == -1) {
|
||||
return 0;
|
||||
}
|
||||
#else
|
||||
(void) master_loop(uwsgi.orig_argv, uwsgi.environ);
|
||||
(void) master_loop(uwsgi.argv, uwsgi.environ);
|
||||
#endif
|
||||
//from now on the process is a real worker
|
||||
}
|
||||
@@ -2912,6 +2954,7 @@ static int manage_base_opt(int i, char *optarg) {
|
||||
struct uwsgi_cron *uc, *old_uc;
|
||||
struct uwsgi_socket *uwsgi_sock = NULL;
|
||||
int zerg_fd;
|
||||
mode_t umask_mode;
|
||||
|
||||
switch (i) {
|
||||
|
||||
@@ -2921,6 +2964,10 @@ static int manage_base_opt(int i, char *optarg) {
|
||||
uwsgi.auto_procname = 1;
|
||||
uwsgi.procname_prefix = optarg;
|
||||
return 1;
|
||||
case LONG_ARGS_PROCNAME_PREFIX_SP:
|
||||
uwsgi.auto_procname = 1;
|
||||
uwsgi.procname_prefix = uwsgi_concat2(optarg, " ");
|
||||
return 1;
|
||||
case LONG_ARGS_PROCNAME_APPEND:
|
||||
uwsgi.auto_procname = 1;
|
||||
uwsgi.procname_append = optarg;
|
||||
@@ -3064,10 +3111,6 @@ static int manage_base_opt(int i, char *optarg) {
|
||||
uwsgi.lazy = 1;
|
||||
return 1;
|
||||
case LONG_ARGS_LOG_MAXSIZE:
|
||||
if (!uwsgi.log_master) {
|
||||
uwsgi.original_log_fd = dup(1);
|
||||
create_logpipe();
|
||||
}
|
||||
uwsgi.log_master = 1;
|
||||
uwsgi.log_maxsize = atoi(optarg);
|
||||
return 1;
|
||||
@@ -3075,10 +3118,6 @@ static int manage_base_opt(int i, char *optarg) {
|
||||
uwsgi.log_backupname = optarg;
|
||||
return 1;
|
||||
case LONG_ARGS_LOG_MASTER:
|
||||
if (!uwsgi.log_master) {
|
||||
uwsgi.original_log_fd = dup(1);
|
||||
create_logpipe();
|
||||
}
|
||||
uwsgi.log_master = 1;
|
||||
return 1;
|
||||
case LONG_ARGS_LOG_SOCKET:
|
||||
@@ -3346,13 +3385,8 @@ static int manage_base_opt(int i, char *optarg) {
|
||||
}
|
||||
return 1;
|
||||
case LONG_ARGS_SUBSCRIBE_TO:
|
||||
if (uwsgi.subscriptions_cnt < MAX_SUBSCRIPTIONS) {
|
||||
uwsgi.subscriptions[uwsgi.subscriptions_cnt] = optarg;
|
||||
uwsgi.subscriptions_cnt++;
|
||||
}
|
||||
else {
|
||||
uwsgi_log("you can specify at most %d --attach-daemons options\n", MAX_SUBSCRIPTIONS);
|
||||
}
|
||||
uwsgi.master_process = 1;
|
||||
uwsgi_string_new_list(&uwsgi.subscriptions, optarg);
|
||||
return 1;
|
||||
#ifdef __linux__
|
||||
case LONG_ARGS_CGROUP:
|
||||
@@ -3435,6 +3469,13 @@ 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_MULES:
|
||||
uwsgi.master_process = 1;
|
||||
for(i=0;i<atoi(optarg);i++) {
|
||||
uwsgi.mules_cnt++;
|
||||
uwsgi_string_new_list(&uwsgi.mules_patches, optarg);
|
||||
}
|
||||
return 1;
|
||||
case LONG_ARGS_FARM:
|
||||
uwsgi.master_process = 1;
|
||||
uwsgi.farms_cnt++;
|
||||
@@ -3690,6 +3731,23 @@ static int manage_base_opt(int i, char *optarg) {
|
||||
case LONG_ARGS_CHOWN_SOCKET:
|
||||
uwsgi.chown_socket = optarg;
|
||||
return 1;
|
||||
case LONG_ARGS_UMASK:
|
||||
if (strlen(optarg) < 3) {
|
||||
uwsgi_log("invalid umask: %s\n", optarg);
|
||||
}
|
||||
umask_mode = 0;
|
||||
if (strlen(optarg) == 3) {
|
||||
umask_mode = (umask_mode << 3) + (optarg[0] - '0');
|
||||
umask_mode = (umask_mode << 3) + (optarg[1] - '0');
|
||||
umask_mode = (umask_mode << 3) + (optarg[2] - '0');
|
||||
}
|
||||
else {
|
||||
umask_mode = (umask_mode << 3) + (optarg[1] - '0');
|
||||
umask_mode = (umask_mode << 3) + (optarg[2] - '0');
|
||||
umask_mode = (umask_mode << 3) + (optarg[3] - '0');
|
||||
}
|
||||
umask(umask_mode);
|
||||
return 1;
|
||||
case 'C':
|
||||
uwsgi.chmod_socket = 1;
|
||||
if (optarg) {
|
||||
|
||||
@@ -8,6 +8,8 @@ extern "C" {
|
||||
|
||||
#define UMAX16 65536
|
||||
|
||||
#define UMAX64_STR "18446744073709551616"
|
||||
|
||||
#define uwsgi_error(x) uwsgi_log("%s: %s [%s line %d]\n", x, strerror(errno), __FILE__, __LINE__);
|
||||
#define uwsgi_fatal_error(x) uwsgi_error(x); exit(1);
|
||||
#define uwsgi_error_open(x) uwsgi_log("open(\"%s\"): %s [%s line %d]\n", x, strerror(errno), __FILE__, __LINE__);
|
||||
@@ -37,7 +39,6 @@ extern "C" {
|
||||
#define MAX_RPC 64
|
||||
#define MAX_GATEWAYS 64
|
||||
#define MAX_DAEMONS 8
|
||||
#define MAX_SUBSCRIPTIONS 8
|
||||
#define MAX_CRONS 64
|
||||
|
||||
#ifndef UWSGI_LOAD_EMBEDDED_PLUGINS
|
||||
@@ -533,6 +534,9 @@ struct uwsgi_opt {
|
||||
#define LONG_ARGS_PROCNAME 17155
|
||||
#define LONG_ARGS_PROCNAME_MASTER 17156
|
||||
#define LONG_ARGS_FARM 17157
|
||||
#define LONG_ARGS_MULES 17158
|
||||
#define LONG_ARGS_PROCNAME_PREFIX_SP 17159
|
||||
#define LONG_ARGS_UMASK 17160
|
||||
|
||||
|
||||
#define UWSGI_OK 0
|
||||
@@ -656,6 +660,7 @@ struct uwsgi_plugin {
|
||||
void (*init_apps) (void);
|
||||
void (*fixup) (void);
|
||||
void (*master_fixup) (int);
|
||||
void (*master_cycle) (void);
|
||||
int (*mount_app) (char *, char *, int);
|
||||
int (*manage_udp) (char *, int, char *, int);
|
||||
int (*manage_xml) (char *, char *);
|
||||
@@ -679,13 +684,14 @@ struct uwsgi_plugin {
|
||||
void (*jail) (int (*)(void *), char **);
|
||||
|
||||
int (*mule)(char *);
|
||||
int (*mule_msg)(char *, size_t);
|
||||
struct uwsgi_help_item *help;
|
||||
|
||||
};
|
||||
|
||||
#ifdef UWSGI_PCRE
|
||||
#include <pcre.h>
|
||||
void uwsgi_regexp_build(char *re, pcre **pattern, pcre_extra **pattern_extra);
|
||||
int uwsgi_regexp_build(char *re, pcre **pattern, pcre_extra **pattern_extra);
|
||||
int uwsgi_regexp_match(pcre *pattern, pcre_extra *pattern_extra, char *subject, int length);
|
||||
#endif
|
||||
|
||||
@@ -1513,8 +1519,7 @@ struct uwsgi_server {
|
||||
int startup_daemons_cnt;
|
||||
|
||||
// subscription client
|
||||
char *subscriptions[MAX_SUBSCRIPTIONS];
|
||||
int subscriptions_cnt;
|
||||
struct uwsgi_string_list *subscriptions;
|
||||
|
||||
};
|
||||
|
||||
@@ -1636,6 +1641,7 @@ struct uwsgi_shared {
|
||||
int spooler_signal_pipe[2];
|
||||
#endif
|
||||
int mule_signal_pipe[2];
|
||||
int mule_queue_pipe[2];
|
||||
|
||||
struct uwsgi_signal_entry signal_table[256];
|
||||
|
||||
@@ -2153,30 +2159,6 @@ struct uwsgi_dict {
|
||||
struct uwsgi_dict_item *items;
|
||||
};
|
||||
|
||||
#define SUBSCRIBER_PAGESIZE 4096
|
||||
#define SUBSCRIBER_NODES (SUBSCRIBER_PAGESIZE/128)-4
|
||||
|
||||
struct uwsgi_subscriber_name {
|
||||
uint16_t len;
|
||||
char name[128];
|
||||
|
||||
uint8_t modifier1;
|
||||
uint8_t modifier2;
|
||||
|
||||
// total requests managed by this node
|
||||
uint64_t requests;
|
||||
// bytes transferred by this node
|
||||
uint64_t transferred;
|
||||
|
||||
};
|
||||
|
||||
struct uwsgi_subscriber {
|
||||
uint64_t nodes;
|
||||
uint64_t current;
|
||||
// support upto md5
|
||||
char auth[32];
|
||||
struct uwsgi_subscriber_name names[SUBSCRIBER_NODES];
|
||||
};
|
||||
|
||||
struct uwsgi_subscribe_req {
|
||||
char *key;
|
||||
@@ -2192,13 +2174,6 @@ struct uwsgi_subscribe_req {
|
||||
uint8_t modifier2;
|
||||
};
|
||||
|
||||
struct uwsgi_dict *uwsgi_dict_create(uint64_t, uint64_t);
|
||||
void uwsgi_add_subscriber(struct uwsgi_dict *, struct uwsgi_subscribe_req *);
|
||||
char *uwsgi_dict_get(struct uwsgi_dict *, char *, uint16_t, uint64_t *);
|
||||
int uwsgi_dict_set(struct uwsgi_dict *, char *, uint16_t, char *, uint64_t);
|
||||
|
||||
struct uwsgi_subscriber_name *uwsgi_get_subscriber(struct uwsgi_dict *, char *, uint16_t);
|
||||
|
||||
#ifndef _NO_UWSGI_RB
|
||||
#include "lib/rbtree.h"
|
||||
|
||||
@@ -2400,6 +2375,7 @@ int uwsgi_simple_parse_vars(struct wsgi_request *, char *, char *);
|
||||
|
||||
void uwsgi_build_mime_dict(char *);
|
||||
struct uwsgi_dyn_dict *uwsgi_dyn_dict_new(struct uwsgi_dyn_dict **, char *, int, char *, int);
|
||||
void uwsgi_dyn_dict_del(struct uwsgi_dyn_dict *);
|
||||
|
||||
void uwsgi_send_stats(int);
|
||||
|
||||
@@ -2421,6 +2397,53 @@ 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);
|
||||
struct uwsgi_farm *get_farm_by_name(char *);
|
||||
|
||||
struct uwsgi_subscribe_slot;
|
||||
|
||||
struct uwsgi_subscribe_node {
|
||||
|
||||
char name[0xff];
|
||||
uint16_t len;
|
||||
uint8_t modifier1;
|
||||
uint8_t modifier2;
|
||||
|
||||
time_t last_check;
|
||||
|
||||
uint64_t requests;
|
||||
uint64_t transferred;
|
||||
|
||||
struct uwsgi_subscribe_slot *slot;
|
||||
|
||||
struct uwsgi_subscribe_node *next;
|
||||
};
|
||||
|
||||
struct uwsgi_subscribe_slot {
|
||||
|
||||
char key[0xff];
|
||||
uint16_t keylen;
|
||||
|
||||
#ifdef UWSGI_PCRE
|
||||
pcre *pattern;
|
||||
pcre_extra *pattern_extra;
|
||||
#endif
|
||||
|
||||
uint64_t hits;
|
||||
|
||||
// used for round robin
|
||||
uint64_t rr;
|
||||
|
||||
struct uwsgi_subscribe_node *nodes;
|
||||
|
||||
struct uwsgi_subscribe_slot *prev;
|
||||
struct uwsgi_subscribe_slot *next;
|
||||
};
|
||||
|
||||
|
||||
struct uwsgi_subscribe_slot *uwsgi_get_subscribe_slot(struct uwsgi_subscribe_slot **, char *, uint16_t, int);
|
||||
struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slot **, char *, uint16_t, int);
|
||||
void uwsgi_remove_subscribe_node(struct uwsgi_subscribe_slot **, struct uwsgi_subscribe_node *);
|
||||
struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slot **, struct uwsgi_subscribe_req *, int);
|
||||
|
||||
#ifdef UWSGI_CAP
|
||||
void uwsgi_build_cap(char *);
|
||||
|
||||
+1
-1
@@ -231,7 +231,7 @@ class uConf(object):
|
||||
ulp.close()
|
||||
|
||||
self.config.read(filename)
|
||||
self.gcc_list = ['utils', 'protocol', 'socket', 'logging', 'master', 'master_utils', 'emperor', 'notify', 'mule',
|
||||
self.gcc_list = ['utils', 'protocol', 'socket', 'logging', 'master', 'master_utils', 'emperor', 'notify', 'mule', 'subscription',
|
||||
'plugins', 'lock', 'cache', 'queue', 'event', 'signal', 'rpc', 'gateway', 'loop', 'lib/rbtree', 'lib/amqp', 'rb_timers', 'uwsgi']
|
||||
# add protocols
|
||||
self.gcc_list.append('proto/base')
|
||||
|
||||
@@ -1,6 +1,12 @@
|
||||
import uwsgi
|
||||
from threading import Thread
|
||||
|
||||
try:
|
||||
import cPickle as pickle
|
||||
except:
|
||||
import pickle
|
||||
|
||||
|
||||
if uwsgi.masterpid() == 0:
|
||||
raise Exception("you have to enable the uWSGI master process to use this module")
|
||||
|
||||
@@ -8,6 +14,7 @@ if uwsgi.opt.get('lazy'):
|
||||
raise Exception("uWSGI lazy mode is not supported by this module")
|
||||
|
||||
spooler_functions = {}
|
||||
mule_functions = {}
|
||||
postfork_chain = []
|
||||
|
||||
def get_free_signal():
|
||||
@@ -75,6 +82,46 @@ class spoolraw(spool):
|
||||
return uwsgi.spool(arguments)
|
||||
|
||||
|
||||
class mulefunc(object):
|
||||
|
||||
def __init__(self, f):
|
||||
if callable(f):
|
||||
self.fname = f.__name__
|
||||
self.mule = 0
|
||||
mule_functions[f.__name__] = f
|
||||
else:
|
||||
self.mule = f
|
||||
self.fname = None
|
||||
|
||||
def real_call(self, *args, **kwargs):
|
||||
uwsgi.mule_msg(pickle.dumps(
|
||||
{
|
||||
'service': 'uwsgi_mulefunc',
|
||||
'func':self.fname,
|
||||
'args': args,
|
||||
'kwargs': kwargs
|
||||
}
|
||||
), self.mule)
|
||||
|
||||
|
||||
def __call__(self, *args, **kwargs):
|
||||
|
||||
if not self.fname:
|
||||
self.fname = args[0].__name__
|
||||
mule_functions[self.fname] = args[0]
|
||||
return self.real_call
|
||||
|
||||
return self.real_call(*args, **kwargs)
|
||||
|
||||
|
||||
def mule_msg_dispatcher(message):
|
||||
msg = pickle.loads(message)
|
||||
if msg['service'] == 'uwsgi_mulefunc':
|
||||
return mule_functions[msg['func']](*msg['args'],**msg['kwargs'])
|
||||
|
||||
uwsgi.mule_msg_hook = mule_msg_dispatcher
|
||||
|
||||
|
||||
class rpc(object):
|
||||
|
||||
def __init__(self, name):
|
||||
@@ -108,6 +155,44 @@ class farm(object):
|
||||
def __call__(self, f):
|
||||
postfork_chain.append(farm_loop(f, self.name))
|
||||
|
||||
class mule_loop(object):
|
||||
|
||||
def __init__(self, f, num):
|
||||
self.f = f
|
||||
self.num = num
|
||||
|
||||
def __call__(self):
|
||||
if uwsgi.mule_id() == self.num:
|
||||
self.f()
|
||||
|
||||
class mule(object):
|
||||
def __init__(self, num):
|
||||
self.num = num
|
||||
|
||||
def __call__(self, f):
|
||||
postfork_chain.append(mule_loop(f, self.num))
|
||||
|
||||
class mulemsg_loop(object):
|
||||
|
||||
def __init__(self, f, num):
|
||||
self.f = f
|
||||
self.num = num
|
||||
|
||||
def __call__(self):
|
||||
if uwsgi.mule_id() == self.num:
|
||||
print " i am the mule"
|
||||
while True:
|
||||
message = uwsgi.mule_get_msg()
|
||||
if message:
|
||||
self.f(message)
|
||||
|
||||
class mulemsg(object):
|
||||
def __init__(self, num):
|
||||
self.num = num
|
||||
|
||||
def __call__(self, f):
|
||||
postfork_chain.append(mulemsg_loop(f, self.num))
|
||||
|
||||
class signal(object):
|
||||
|
||||
def __init__(self, num, **kwargs):
|
||||
@@ -172,6 +257,15 @@ class filemon(object):
|
||||
uwsgi.add_file_monitor(self.num, self.fsobj)
|
||||
return f
|
||||
|
||||
class erlang(object):
|
||||
|
||||
def __init__(self, name):
|
||||
self.name = name
|
||||
|
||||
def __call__(self, f):
|
||||
uwsgi.erlang_register_process(self.name, f)
|
||||
return f
|
||||
|
||||
class lock(object):
|
||||
def __init__(self, f):
|
||||
self.f = f
|
||||
|
||||
+10
-7
@@ -5,9 +5,9 @@ import sys
|
||||
from uwsgidecorators import *
|
||||
gc.set_debug(gc.DEBUG_SAVEALL)
|
||||
|
||||
print os.environ
|
||||
print sys.modules
|
||||
print sys.argv
|
||||
print(os.environ)
|
||||
print(sys.modules)
|
||||
print(sys.argv)
|
||||
|
||||
try:
|
||||
if sys.argv[1] == 'debug':
|
||||
@@ -50,7 +50,10 @@ def setprocname():
|
||||
|
||||
def application(env, start_response):
|
||||
|
||||
uwsgi.mule_msg(env['REQUEST_URI'], 1)
|
||||
try:
|
||||
uwsgi.mule_msg(env['REQUEST_URI'], 1)
|
||||
except:
|
||||
pass
|
||||
|
||||
req = uwsgi.workers()[uwsgi.worker_id()-1]['requests']
|
||||
|
||||
@@ -58,7 +61,7 @@ def application(env, start_response):
|
||||
|
||||
gc.collect(2)
|
||||
if DEBUG:
|
||||
print env['wsgi.input'].fileno()
|
||||
print(env['wsgi.input'].fileno())
|
||||
|
||||
if routes.has_key(env['PATH_INFO']):
|
||||
return routes[env['PATH_INFO']](env, start_response)
|
||||
@@ -66,12 +69,12 @@ def application(env, start_response):
|
||||
start_response('200 OK', [('Content-Type', 'text/html')])
|
||||
|
||||
if DEBUG:
|
||||
print env['wsgi.input'].fileno()
|
||||
print(env['wsgi.input'].fileno())
|
||||
|
||||
gc.collect(2)
|
||||
|
||||
if DEBUG:
|
||||
print len(gc.get_objects())
|
||||
print(len(gc.get_objects()))
|
||||
|
||||
workers = ''
|
||||
for w in uwsgi.workers():
|
||||
|
||||
Reference in New Issue
Block a user