mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-10-04 16:08:31 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7604c67018 | ||
|
|
afc962709f | ||
|
|
c11e0654a0 | ||
|
|
6e4f47cc76 | ||
|
|
20c9f0ef40 | ||
|
|
ef3c8a5bd9 | ||
|
|
76d0438fee | ||
|
|
18fd338908 | ||
|
|
2b35e95706 | ||
|
|
4761146599 | ||
|
|
3924cc639e | ||
|
|
5f09de0469 | ||
|
|
4f0e3fb79f | ||
|
|
05b6baf173 | ||
|
|
149823aac4 | ||
|
|
85d86dc10e | ||
|
|
e97d921909 | ||
|
|
c5e20bf4da | ||
|
|
a09d25ccaa | ||
|
|
bf01cea062 | ||
|
|
b6459be3e4 | ||
|
|
d42788c17b | ||
|
|
872730c116 | ||
|
|
6fce0aee84 | ||
|
|
e3d8a02614 | ||
|
|
b58696399e | ||
|
|
519c9f58b3 | ||
|
|
2c77c5c639 | ||
|
|
fb9e5a1891 | ||
|
|
5aa7d1a79f | ||
|
|
bc3442ba29 | ||
|
|
3061e36089 | ||
|
|
4f8c41cf61 | ||
|
|
b0065c379e | ||
|
|
57f71c974d | ||
|
|
b820fefbbd | ||
|
|
5f759f0d26 | ||
|
|
f7eb9ed1cd | ||
|
|
3a82732b19 | ||
|
|
d74e826f71 | ||
|
|
c846a9b388 | ||
|
|
001e4b6e39 | ||
|
|
a51a396f47 | ||
|
|
17965ea870 | ||
|
|
8a032b41ec | ||
|
|
a489cfc2c1 | ||
|
|
3b32672ff9 | ||
|
|
858131ba57 | ||
|
|
8a8d2f8675 | ||
|
|
9d37725323 | ||
|
|
e973492267 | ||
|
|
105a68bce9 | ||
|
|
0587a0ca53 | ||
|
|
d4b259d5e0 | ||
|
|
9e591c587f | ||
|
|
e7cce7b785 | ||
|
|
3285b224c2 | ||
|
|
2bba4cf15a |
@@ -27,3 +27,4 @@ Roberto Leandrini
|
||||
Ryan Petrello
|
||||
Danila Shtan <danila@shtan.ru>
|
||||
Ævar Arnfjörð Bjarmason
|
||||
Yu Zhao (getcwd)
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
[uwsgi]
|
||||
main_plugin = python,asyncio,greenlet
|
||||
inherit = base
|
||||
+140
-33
@@ -472,6 +472,38 @@ uint32_t uwsgi_cache_exists2(struct uwsgi_cache *uc, char *key, uint16_t keylen)
|
||||
return uwsgi_cache_get_index(uc, key, keylen);
|
||||
}
|
||||
|
||||
static void lru_remove_item(struct uwsgi_cache *uc, uint64_t index)
|
||||
{
|
||||
struct uwsgi_cache_item *prev, *next, *curr = cache_item(index);
|
||||
|
||||
if (curr->next) {
|
||||
next = cache_item(curr->next);
|
||||
next->prev = curr->prev;
|
||||
} else
|
||||
uc->lru_tail = curr->prev;
|
||||
|
||||
if (curr->prev) {
|
||||
prev = cache_item(curr->prev);
|
||||
prev->next = curr->next;
|
||||
} else
|
||||
uc->lru_head = curr->next;
|
||||
}
|
||||
|
||||
static void lru_add_item(struct uwsgi_cache *uc, uint64_t index)
|
||||
{
|
||||
struct uwsgi_cache_item *prev, *curr = cache_item(index);
|
||||
|
||||
if (uc->lru_tail) {
|
||||
prev = cache_item(uc->lru_tail);
|
||||
prev->next = index;
|
||||
} else
|
||||
uc->lru_head = index;
|
||||
|
||||
curr->next = 0;
|
||||
curr->prev = uc->lru_tail;
|
||||
uc->lru_tail = index;
|
||||
}
|
||||
|
||||
char *uwsgi_cache_get2(struct uwsgi_cache *uc, char *key, uint16_t keylen, uint64_t * valsize) {
|
||||
|
||||
uint64_t index = uwsgi_cache_get_index(uc, key, keylen);
|
||||
@@ -481,6 +513,10 @@ char *uwsgi_cache_get2(struct uwsgi_cache *uc, char *key, uint16_t keylen, uint6
|
||||
if (uci->flags & UWSGI_CACHE_FLAG_UNGETTABLE)
|
||||
return NULL;
|
||||
*valsize = uci->valsize;
|
||||
if (uc->purge_lru) {
|
||||
lru_remove_item(uc, index);
|
||||
lru_add_item(uc, index);
|
||||
}
|
||||
uci->hits++;
|
||||
uc->hits++;
|
||||
return uc->data + (uci->first_block * uc->blocksize);
|
||||
@@ -520,6 +556,10 @@ char *uwsgi_cache_get3(struct uwsgi_cache *uc, char *key, uint16_t keylen, uint6
|
||||
*valsize = uci->valsize;
|
||||
if (expires)
|
||||
*expires = uci->expires;
|
||||
if (uc->purge_lru) {
|
||||
lru_remove_item(uc, index);
|
||||
lru_add_item(uc, index);
|
||||
}
|
||||
uci->hits++;
|
||||
uc->hits++;
|
||||
return uc->data + (uci->first_block * uc->blocksize);
|
||||
@@ -588,6 +628,10 @@ int uwsgi_cache_del2(struct uwsgi_cache *uc, char *key, uint16_t keylen, uint64_
|
||||
// reset hashtable entry
|
||||
uc->hashtable[uci->hash % uc->hashsize] = 0;
|
||||
}
|
||||
|
||||
if (uc->purge_lru)
|
||||
lru_remove_item(uc, index);
|
||||
|
||||
uc->n_items--;
|
||||
}
|
||||
|
||||
@@ -666,9 +710,17 @@ int uwsgi_cache_set2(struct uwsgi_cache *uc, char *key, uint16_t keylen, char *v
|
||||
index = uwsgi_cache_get_index(uc, key, keylen);
|
||||
if (!index) {
|
||||
if (!uc->unused_blocks_stack_ptr) {
|
||||
uwsgi_log("*** DANGER cache \"%s\" is FULL !!! ***\n", uc->name);
|
||||
if (!uc->ignore_full) {
|
||||
if (uc->purge_lru)
|
||||
uwsgi_log("LRU item will be purged from cache \"%s\"\n", uc->name);
|
||||
else
|
||||
uwsgi_log("*** DANGER cache \"%s\" is FULL !!! ***\n", uc->name);
|
||||
}
|
||||
uc->full++;
|
||||
goto end;
|
||||
if (uc->purge_lru && uc->lru_head)
|
||||
uwsgi_cache_del2(uc, NULL, 0, uc->lru_head, UWSGI_CACHE_FLAG_LOCAL);
|
||||
if (!uc->unused_blocks_stack_ptr)
|
||||
goto end;
|
||||
}
|
||||
|
||||
index = uc->unused_blocks_stack[uc->unused_blocks_stack_ptr];
|
||||
@@ -681,7 +733,8 @@ int uwsgi_cache_set2(struct uwsgi_cache *uc, char *key, uint16_t keylen, char *v
|
||||
else {
|
||||
uci->first_block = uwsgi_cache_find_free_blocks(uc, vallen);
|
||||
if (uci->first_block == 0xffffffffffffffffLLU) {
|
||||
uwsgi_log("*** DANGER cache \"%s\" is FULL !!! ***\n", uc->name);
|
||||
if (!uc->ignore_full)
|
||||
uwsgi_log("*** DANGER cache \"%s\" is FULL !!! ***\n", uc->name);
|
||||
uc->full++;
|
||||
uc->unused_blocks_stack_ptr++;
|
||||
goto end;
|
||||
@@ -696,9 +749,13 @@ int uwsgi_cache_set2(struct uwsgi_cache *uc, char *key, uint16_t keylen, char *v
|
||||
uc->blocks_bitmap_pos = uci->first_block + needed_blocks;
|
||||
}
|
||||
}
|
||||
if (expires && !(flags & UWSGI_CACHE_FLAG_ABSEXPIRE)) {
|
||||
if (uc->purge_lru)
|
||||
lru_add_item(uc, index);
|
||||
else if (expires && !(flags & UWSGI_CACHE_FLAG_ABSEXPIRE)) {
|
||||
now = uwsgi_now();
|
||||
expires += now;
|
||||
if (!uc->next_scan || uc->next_scan > expires)
|
||||
uc->next_scan = expires;
|
||||
}
|
||||
uci->expires = expires;
|
||||
uci->hash = uc->hash->func(key, keylen);
|
||||
@@ -763,9 +820,16 @@ int uwsgi_cache_set2(struct uwsgi_cache *uc, char *key, uint16_t keylen, char *v
|
||||
}
|
||||
else if (flags & UWSGI_CACHE_FLAG_UPDATE) {
|
||||
uci = cache_item(index);
|
||||
if (expires && !(flags & UWSGI_CACHE_FLAG_ABSEXPIRE) && !(flags & UWSGI_CACHE_FLAG_FIXEXPIRE)) {
|
||||
now = uwsgi_now();
|
||||
expires += now;
|
||||
if (!(flags & UWSGI_CACHE_FLAG_FIXEXPIRE)) {
|
||||
if (uc->purge_lru) {
|
||||
lru_remove_item(uc, index);
|
||||
lru_add_item(uc, index);
|
||||
} else if (expires && !(flags & UWSGI_CACHE_FLAG_ABSEXPIRE)) {
|
||||
now = uwsgi_now();
|
||||
expires += now;
|
||||
if (!uc->next_scan || uc->next_scan > expires)
|
||||
uc->next_scan = expires;
|
||||
}
|
||||
uci->expires = expires;
|
||||
}
|
||||
if (uc->blocks_bitmap) {
|
||||
@@ -773,7 +837,8 @@ int uwsgi_cache_set2(struct uwsgi_cache *uc, char *key, uint16_t keylen, char *v
|
||||
uint64_t old_first_block = uci->first_block;
|
||||
uci->first_block = uwsgi_cache_find_free_blocks(uc, vallen);
|
||||
if (uci->first_block == 0xffffffffffffffffLLU) {
|
||||
uwsgi_log("*** DANGER cache \"%s\" is FULL !!! ***\n", uc->name);
|
||||
if (!uc->ignore_full)
|
||||
uwsgi_log("*** DANGER cache \"%s\" is FULL !!! ***\n", uc->name);
|
||||
uc->full++;
|
||||
uci->first_block = old_first_block;
|
||||
goto end;
|
||||
@@ -983,39 +1048,66 @@ void *cache_udp_server_loop(void *ucache) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
static uint64_t cache_sweeper_free_items(struct uwsgi_cache *uc) {
|
||||
uint64_t i;
|
||||
uint64_t freed_items = 0;
|
||||
|
||||
if (uc->no_expire || uc->purge_lru)
|
||||
return 0;
|
||||
|
||||
uwsgi_rlock(uc->lock);
|
||||
if (!uc->next_scan || uc->next_scan > (uint64_t)uwsgi.current_time) {
|
||||
uwsgi_rwunlock(uc->lock);
|
||||
return 0;
|
||||
}
|
||||
uwsgi_rwunlock(uc->lock);
|
||||
|
||||
// skip the first slot
|
||||
for (i = 1; i < uc->max_items; i++) {
|
||||
struct uwsgi_cache_item *uci = cache_item(i);
|
||||
|
||||
uwsgi_wlock(uc->lock);
|
||||
// we reset next scan time first, then we find the least
|
||||
// expiration time from those that are NOT expired yet.
|
||||
if (i == 1)
|
||||
uc->next_scan = 0;
|
||||
|
||||
if (uci->expires) {
|
||||
if (uci->expires <= (uint64_t)uwsgi.current_time) {
|
||||
uwsgi_cache_del2(uc, NULL, 0, i, UWSGI_CACHE_FLAG_LOCAL);
|
||||
freed_items++;
|
||||
} else if (!uc->next_scan || uc->next_scan > uci->expires) {
|
||||
uc->next_scan = uci->expires;
|
||||
}
|
||||
}
|
||||
uwsgi_rwunlock(uc->lock);
|
||||
}
|
||||
|
||||
return freed_items;
|
||||
}
|
||||
|
||||
static void *cache_sweeper_loop(void *ucache) {
|
||||
|
||||
uint64_t i;
|
||||
// block all signals
|
||||
sigset_t smask;
|
||||
sigfillset(&smask);
|
||||
pthread_sigmask(SIG_BLOCK, &smask, NULL);
|
||||
|
||||
struct uwsgi_cache *uc = (struct uwsgi_cache *) ucache;
|
||||
|
||||
if (!uwsgi.cache_expire_freq)
|
||||
uwsgi.cache_expire_freq = 3;
|
||||
|
||||
// remove expired cache items TODO use rb_tree timeouts
|
||||
for (;;) {
|
||||
struct uwsgi_cache *uc;
|
||||
|
||||
for (uc = (struct uwsgi_cache *)ucache; uc; uc = uc->next) {
|
||||
uint64_t freed_items = cache_sweeper_free_items(uc);
|
||||
if (uwsgi.cache_report_freed_items && freed_items)
|
||||
uwsgi_log("freed %llu items for cache \"%s\"\n", (unsigned long long)freed_items, uc->name);
|
||||
}
|
||||
|
||||
sleep(uwsgi.cache_expire_freq);
|
||||
uint64_t freed_items = 0;
|
||||
// skip the first slot
|
||||
for (i = 1; i < uc->max_items; i++) {
|
||||
uwsgi_wlock(uc->lock);
|
||||
struct uwsgi_cache_item *uci = cache_item(i);
|
||||
if (uci->expires) {
|
||||
if (uci->expires < (uint64_t) uwsgi.current_time) {
|
||||
uwsgi_cache_del2(uc, NULL, 0, i, UWSGI_CACHE_FLAG_LOCAL);
|
||||
freed_items++;
|
||||
}
|
||||
}
|
||||
uwsgi_rwunlock(uc->lock);
|
||||
}
|
||||
if (uwsgi.cache_report_freed_items && freed_items > 0) {
|
||||
uwsgi_log("freed %llu items for cache \"%s\"\n", (unsigned long long) freed_items, uc->name);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
return NULL;
|
||||
}
|
||||
@@ -1035,16 +1127,21 @@ void uwsgi_cache_sync_all() {
|
||||
|
||||
void uwsgi_cache_start_sweepers() {
|
||||
struct uwsgi_cache *uc = uwsgi.caches;
|
||||
|
||||
if (uwsgi.cache_no_expire)
|
||||
return;
|
||||
|
||||
while(uc) {
|
||||
pthread_t cache_sweeper;
|
||||
if (!uwsgi.cache_no_expire && !uc->no_expire) {
|
||||
if (pthread_create(&cache_sweeper, NULL, cache_sweeper_loop, (void *) uc)) {
|
||||
if (!uc->no_expire && !uc->purge_lru) {
|
||||
if (pthread_create(&cache_sweeper, NULL, cache_sweeper_loop, uwsgi.caches)) {
|
||||
uwsgi_error("pthread_create()");
|
||||
uwsgi_log("unable to run the sweeper for cache \"%s\" !!!\n", uc->name);
|
||||
uwsgi_log("unable to run the sweeper!!!\n");
|
||||
}
|
||||
else {
|
||||
uwsgi_log("sweeper thread enabled for cache \"%s\"\n", uc->name);
|
||||
uwsgi_log("sweeper thread enabled\n");
|
||||
}
|
||||
break;
|
||||
}
|
||||
uc = uc->next;
|
||||
}
|
||||
@@ -1121,6 +1218,8 @@ struct uwsgi_cache *uwsgi_cache_create(char *arg) {
|
||||
char *c_bitmap = NULL;
|
||||
char *c_use_last_modified = NULL;
|
||||
char *c_math_initial = NULL;
|
||||
char *c_ignore_full = NULL;
|
||||
char *c_purge_lru = NULL;
|
||||
|
||||
if (uwsgi_kvlist_parse(arg, strlen(arg), ',', '=',
|
||||
"name", &c_name,
|
||||
@@ -1148,6 +1247,8 @@ struct uwsgi_cache *uwsgi_cache_create(char *arg) {
|
||||
"bitmap", &c_bitmap,
|
||||
"lastmod", &c_use_last_modified,
|
||||
"math_initial", &c_math_initial,
|
||||
"ignore_full", &c_ignore_full,
|
||||
"purge_lru", &c_purge_lru,
|
||||
NULL)) {
|
||||
uwsgi_log("unable to parse cache definition\n");
|
||||
exit(1);
|
||||
@@ -1195,6 +1296,7 @@ struct uwsgi_cache *uwsgi_cache_create(char *arg) {
|
||||
uc->max_item_size = uc->blocksize * uc->blocks;
|
||||
}
|
||||
if (c_use_last_modified) uc->use_last_modified = 1;
|
||||
if (c_ignore_full) uc->ignore_full = 1;
|
||||
|
||||
if (c_math_initial) uc->math_initial = strtol(c_math_initial, NULL, 10);
|
||||
|
||||
@@ -1229,6 +1331,8 @@ struct uwsgi_cache *uwsgi_cache_create(char *arg) {
|
||||
}
|
||||
}
|
||||
|
||||
if (c_purge_lru)
|
||||
uc->purge_lru = 1;
|
||||
}
|
||||
|
||||
uwsgi_cache_init(uc);
|
||||
@@ -1484,7 +1588,10 @@ char *uwsgi_cache_magic_get(char *key, uint16_t keylen, uint64_t *vallen, uint64
|
||||
|
||||
// we have a local cache !!!
|
||||
if (uc) {
|
||||
uwsgi_rlock(uc->lock);
|
||||
if (uc->purge_lru)
|
||||
uwsgi_wlock(uc->lock);
|
||||
else
|
||||
uwsgi_rlock(uc->lock);
|
||||
char *value = uwsgi_cache_get3(uc, key, keylen, vallen, expires);
|
||||
if (!value) {
|
||||
uwsgi_rwunlock(uc->lock);
|
||||
|
||||
+58
-2
@@ -701,6 +701,23 @@ void emperor_del(struct uwsgi_instance *c_ui) {
|
||||
|
||||
}
|
||||
|
||||
void emperor_back_to_ondemand(struct uwsgi_instance *c_ui) {
|
||||
if (c_ui->status > 0) return;
|
||||
|
||||
// remove uWSGI instance
|
||||
|
||||
if (c_ui->pid != -1) {
|
||||
if (write(c_ui->pipe[0], "\0", 1) != 1) {
|
||||
uwsgi_error("emperor_stop()/write()");
|
||||
}
|
||||
}
|
||||
|
||||
c_ui->status = 2;
|
||||
c_ui->cursed_at = uwsgi_now();
|
||||
|
||||
uwsgi_log_verbose("[emperor] bringing back instance %s to on-demand mode\n", c_ui->name);
|
||||
}
|
||||
|
||||
void emperor_stop(struct uwsgi_instance *c_ui) {
|
||||
if (c_ui->status == 1) return;
|
||||
// remove uWSGI instance
|
||||
@@ -721,7 +738,8 @@ void emperor_curse(struct uwsgi_instance *c_ui) {
|
||||
if (c_ui->status == 1) return;
|
||||
// curse uWSGI instance
|
||||
|
||||
c_ui->status = 1;
|
||||
// take in account on-demand mode
|
||||
if (c_ui->status == 0) c_ui->status = 1;
|
||||
c_ui->cursed_at = uwsgi_now();
|
||||
|
||||
uwsgi_log_verbose("[emperor] curse the uwsgi instance %s (pid: %d)\n", c_ui->name, (int) c_ui->pid);
|
||||
@@ -749,10 +767,22 @@ static void emperor_push_config(struct uwsgi_instance *c_ui) {
|
||||
|
||||
void emperor_respawn(struct uwsgi_instance *c_ui, time_t mod) {
|
||||
|
||||
// if the vassal is being destroyed, do not honour respawns
|
||||
if (c_ui->status > 0) return;
|
||||
|
||||
// check if we are in on_demand mode (the respawn will be ignored)
|
||||
if (c_ui->pid == -1 && c_ui->on_demand_fd > -1) {
|
||||
c_ui->last_mod = mod;
|
||||
uwsgi_log_verbose("[emperor] updated configuration for \"on demand\" instance %s\n", c_ui->name);
|
||||
return;
|
||||
}
|
||||
|
||||
// reload the uWSGI instance
|
||||
if (write(c_ui->pipe[0], "\1", 1) != 1) {
|
||||
// the vassal could be already dead, better to curse it
|
||||
uwsgi_error("emperor_respawn/write()");
|
||||
emperor_curse(c_ui);
|
||||
return;
|
||||
}
|
||||
|
||||
// push the config to the config pipe (if needed)
|
||||
@@ -940,11 +970,14 @@ int uwsgi_emperor_vassal_start(struct uwsgi_instance *n_ui) {
|
||||
n_ui->pid = pid;
|
||||
// close the right side of the pipe
|
||||
close(n_ui->pipe[1]);
|
||||
/* THE ON-DEMAND file descriptir is left mapped to the emperor to allow fast-respawn
|
||||
// TODO add an option to force closing it
|
||||
// close the "on demand" socket
|
||||
if (n_ui->on_demand_fd > -1) {
|
||||
close(n_ui->on_demand_fd);
|
||||
n_ui->on_demand_fd = -1;
|
||||
}
|
||||
*/
|
||||
if (n_ui->use_config) {
|
||||
close(n_ui->pipe_config[1]);
|
||||
}
|
||||
@@ -1684,6 +1717,10 @@ void emperor_loop() {
|
||||
if (rlen <= 0) {
|
||||
// SAFE
|
||||
event_queue_del_fd(uwsgi.emperor_queue, interesting_fd, event_queue_read());
|
||||
if (ui_current->status > 0) {
|
||||
// temporarily set frequency to a low value , so we can eventually fast-restart the instance
|
||||
freq = ui_current->status;
|
||||
}
|
||||
emperor_curse(ui_current);
|
||||
}
|
||||
else {
|
||||
@@ -1700,7 +1737,13 @@ void emperor_loop() {
|
||||
ui_current->last_heartbeat = uwsgi_now();
|
||||
}
|
||||
else if (byte == 22) {
|
||||
emperor_stop(ui_current);
|
||||
// command 22 changes meaning when in "on_demand" mode
|
||||
if (ui_current->on_demand_fd != -1) {
|
||||
emperor_back_to_ondemand(ui_current);
|
||||
}
|
||||
else {
|
||||
emperor_stop(ui_current);
|
||||
}
|
||||
}
|
||||
else if (byte == 30 && uwsgi.emperor_broodlord > 0 && uwsgi.emperor_broodlord_count < uwsgi.emperor_broodlord) {
|
||||
uwsgi_log_verbose("[emperor] going in broodlord mode: launching zergs for %s\n", ui_current->name);
|
||||
@@ -1822,6 +1865,8 @@ recheck:
|
||||
// UNSAFE
|
||||
emperor_add(ui_current->scanner, ui_current->name, ui_current->last_mod, ui_current->config, ui_current->config_len, ui_current->uid, ui_current->gid, ui_current->socket_name);
|
||||
emperor_del(ui_current);
|
||||
// temporarily set frequency to 0, so we can eventually fast-restart the instance
|
||||
freq = 0;
|
||||
}
|
||||
break;
|
||||
}
|
||||
@@ -1831,6 +1876,17 @@ recheck:
|
||||
free(ui_current->config);
|
||||
// SAFE
|
||||
emperor_del(ui_current);
|
||||
// temporarily set frequency to 0, so we can eventually fast-restart the instance
|
||||
freq = 0;
|
||||
break;
|
||||
}
|
||||
// back to on_demand mode ...
|
||||
else if (ui_current->status == 2) {
|
||||
event_queue_add_fd_read(uwsgi.emperor_queue, ui_current->on_demand_fd);
|
||||
ui_current->pid = -1;
|
||||
ui_current->status = 0;
|
||||
ui_current->cursed_at = 0;
|
||||
uwsgi_log("[uwsgi-emperor] %s -> back to \"on demand\" mode, waiting for connections on socket \"%s\" ...\n", ui_current->name, ui_current->socket_name);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -485,6 +485,60 @@ static int uwsgi_hook_callret(char *arg) {
|
||||
return func();
|
||||
}
|
||||
|
||||
static int uwsgi_hook_rpc(char *arg) {
|
||||
|
||||
int ret = -1;
|
||||
size_t i, argc = 0;
|
||||
char **rargv = uwsgi_split_quoted(arg, strlen(arg), " \t", &argc);
|
||||
if (!argc) goto end;
|
||||
if (argc > 256) goto destroy;
|
||||
|
||||
char *argv[256];
|
||||
uint16_t argvs[256];
|
||||
|
||||
char *node = NULL;
|
||||
char *func = rargv[0];
|
||||
|
||||
char *at = strchr(func, '@');
|
||||
if (at) {
|
||||
*at = 0;
|
||||
node = at + 1;
|
||||
}
|
||||
|
||||
for(i=0;i<(argc-1);i++) {
|
||||
size_t a_len = strlen(rargv[i+1]);
|
||||
if (a_len > 0xffff) goto destroy;
|
||||
argv[i] = rargv[i+1] ;
|
||||
argvs[i] = a_len;
|
||||
}
|
||||
|
||||
uint64_t size = 0;
|
||||
// response must be always freed
|
||||
char *response = uwsgi_do_rpc(node, func, argc-1, argv, argvs, &size);
|
||||
if (response) {
|
||||
if (at) *at = '@';
|
||||
uwsgi_log("[rpc result from \"%s\"] %.*s\n", rargv[0], size, response);
|
||||
free(response);
|
||||
ret = 0;
|
||||
}
|
||||
|
||||
destroy:
|
||||
for(i=0;i<argc;i++) {
|
||||
free(rargv[i]);
|
||||
}
|
||||
end:
|
||||
free(rargv);
|
||||
return ret;
|
||||
}
|
||||
|
||||
static int uwsgi_hook_retryrpc(char *arg) {
|
||||
for(;;) {
|
||||
int ret = uwsgi_hook_rpc(arg);
|
||||
if (!ret) break;
|
||||
sleep(2);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
void uwsgi_register_base_hooks() {
|
||||
uwsgi_register_hook("cd", uwsgi_hook_chdir);
|
||||
@@ -523,6 +577,9 @@ void uwsgi_register_base_hooks() {
|
||||
|
||||
uwsgi_register_hook("alarm", uwsgi_hook_alarm);
|
||||
|
||||
uwsgi_register_hook("rpc", uwsgi_hook_rpc);
|
||||
uwsgi_register_hook("retryrpc", uwsgi_hook_retryrpc);
|
||||
|
||||
// for testing
|
||||
uwsgi_register_hook("exit", uwsgi_hook_exit);
|
||||
uwsgi_register_hook("print", uwsgi_hook_print);
|
||||
@@ -561,6 +618,22 @@ void uwsgi_hooks_run(struct uwsgi_string_list *l, char *phase, int fatal) {
|
||||
|
||||
int ret = uh->func(colon+1);
|
||||
if (fatal && ret != 0) {
|
||||
uwsgi_log_verbose("FATAL hook failed, destroying instance\n");
|
||||
if (uwsgi.master_process) {
|
||||
if (uwsgi.workers) {
|
||||
if (uwsgi.workers[0].pid == getpid()) {
|
||||
kill_them_all(0);
|
||||
return;
|
||||
}
|
||||
else {
|
||||
if (kill(uwsgi.workers[0].pid, SIGINT)) {
|
||||
uwsgi_error("uwsgi_hooks_run()/kill()");
|
||||
exit(1);
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
+9
-3
@@ -477,13 +477,19 @@ void uwsgi_check_logrotate(void) {
|
||||
|
||||
int need_rotation = 0;
|
||||
int need_reopen = 0;
|
||||
off_t logsize;
|
||||
|
||||
if (uwsgi.log_master) {
|
||||
uwsgi.shared->logsize = lseek(uwsgi.original_log_fd, 0, SEEK_CUR);
|
||||
logsize = lseek(uwsgi.original_log_fd, 0, SEEK_CUR);
|
||||
}
|
||||
else {
|
||||
uwsgi.shared->logsize = lseek(2, 0, SEEK_CUR);
|
||||
logsize = lseek(2, 0, SEEK_CUR);
|
||||
}
|
||||
if (logsize < 0) {
|
||||
uwsgi_error("uwsgi_check_logrotate()/lseek()");
|
||||
return;
|
||||
}
|
||||
uwsgi.shared->logsize = logsize;
|
||||
|
||||
if (uwsgi.log_maxsize > 0 && (uint64_t) uwsgi.shared->logsize > uwsgi.log_maxsize) {
|
||||
need_rotation = 1;
|
||||
@@ -917,7 +923,7 @@ struct uwsgi_logger *uwsgi_get_logger_from_id(char *id) {
|
||||
struct uwsgi_logger *ul = uwsgi.choosen_logger;
|
||||
|
||||
while (ul) {
|
||||
if (!strcmp(ul->id, id)) {
|
||||
if (ul->id && !strcmp(ul->id, id)) {
|
||||
return ul;
|
||||
}
|
||||
ul = ul->next;
|
||||
|
||||
+2
-2
@@ -330,10 +330,10 @@ next:
|
||||
}
|
||||
|
||||
int ret = poll(mulepoll, count + farms_count, timeout);
|
||||
if (ret <= 0) {
|
||||
if (ret < 0) {
|
||||
uwsgi_error("poll");
|
||||
}
|
||||
else {
|
||||
else if (ret > 0 ) {
|
||||
if (mulepoll[0].revents & POLLIN) {
|
||||
len = read(uwsgi.mules[uwsgi.muleid - 1].queue_pipe[1], message, buffer_size);
|
||||
}
|
||||
|
||||
@@ -237,6 +237,7 @@ static void uwsgi_offload_loop(struct uwsgi_thread *ut) {
|
||||
void *events = event_queue_alloc(uwsgi.offload_threads_events);
|
||||
|
||||
for (;;) {
|
||||
// TODO make timeout tunable
|
||||
int nevents = event_queue_wait_multi(ut->queue, -1, events, uwsgi.offload_threads_events);
|
||||
for (i = 0; i < nevents; i++) {
|
||||
int interesting_fd = event_queue_interesting_fd(events, i);
|
||||
@@ -456,6 +457,13 @@ static int u_offload_transfer_do(struct uwsgi_thread *ut, struct uwsgi_offload_r
|
||||
// write event (or just connected)
|
||||
case 1:
|
||||
if (fd == uor->fd) {
|
||||
// maybe we want only a connection...
|
||||
if (uor->ubuf->pos == 0) {
|
||||
uor->status = 2;
|
||||
if (event_queue_add_fd_read(ut->queue, uor->s)) return -1;
|
||||
if (event_queue_fd_write_to_read(ut->queue, uor->fd)) return -1;
|
||||
return 0;
|
||||
}
|
||||
rlen = write(uor->fd, uor->ubuf->buf + uor->written, uor->ubuf->pos-uor->written);
|
||||
if (rlen > 0) {
|
||||
uor->written += rlen;
|
||||
|
||||
+24
-1
@@ -1159,6 +1159,26 @@ static int uwsgi_router_setscheme(struct uwsgi_route *ur, char *arg) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
// setmodifiers
|
||||
static int uwsgi_router_setmodifier1_func(struct wsgi_request *wsgi_req, struct uwsgi_route *ur) {
|
||||
wsgi_req->uh->modifier1 = ur->custom;
|
||||
return UWSGI_ROUTE_NEXT;
|
||||
}
|
||||
static int uwsgi_router_setmodifier1(struct uwsgi_route *ur, char *arg) {
|
||||
ur->func = uwsgi_router_setmodifier1_func;
|
||||
ur->custom = atoi(arg);
|
||||
return 0;
|
||||
}
|
||||
static int uwsgi_router_setmodifier2_func(struct wsgi_request *wsgi_req, struct uwsgi_route *ur) {
|
||||
wsgi_req->uh->modifier2 = ur->custom;
|
||||
return UWSGI_ROUTE_NEXT;
|
||||
}
|
||||
static int uwsgi_router_setmodifier2(struct uwsgi_route *ur, char *arg) {
|
||||
ur->func = uwsgi_router_setmodifier2_func;
|
||||
ur->custom = atoi(arg);
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
// setuser route
|
||||
static int uwsgi_router_setuser_func(struct wsgi_request *wsgi_req, struct uwsgi_route *ur) {
|
||||
@@ -1684,7 +1704,7 @@ static char *uwsgi_route_var_mime(struct wsgi_request *wsgi_req, char *key, uint
|
||||
char *var_value = uwsgi_get_var(wsgi_req, key, keylen, &var_vallen);
|
||||
if (var_value) {
|
||||
size_t mime_type_len = 0;
|
||||
char *ret = uwsgi_get_mime_type(key, keylen, &mime_type_len);
|
||||
ret = uwsgi_get_mime_type(var_value, var_vallen, &mime_type_len);
|
||||
if (ret) *vallen = mime_type_len;
|
||||
}
|
||||
return ret;
|
||||
@@ -1769,6 +1789,9 @@ void uwsgi_register_embedded_routers() {
|
||||
uwsgi_register_router("setprocname", uwsgi_router_setprocname);
|
||||
uwsgi_register_router("alarm", uwsgi_router_alarm);
|
||||
|
||||
uwsgi_register_router("setmodifier1", uwsgi_router_setmodifier1);
|
||||
uwsgi_register_router("setmodifier2", uwsgi_router_setmodifier2);
|
||||
|
||||
uwsgi_register_router("+", uwsgi_router_simple_math_plus);
|
||||
uwsgi_register_router("-", uwsgi_router_simple_math_minus);
|
||||
uwsgi_register_router("*", uwsgi_router_simple_math_multiply);
|
||||
|
||||
+5
-1
@@ -104,8 +104,12 @@ char *uwsgi_do_rpc(char *node, char *func, uint8_t argc, char *argv[], uint16_t
|
||||
|
||||
if (node == NULL || !strcmp(node, "")) {
|
||||
// allocate the whole buffer
|
||||
if (!uwsgi.rpc_table) {
|
||||
uwsgi_log("local rpc subsystem is still not initialized !!!\n");
|
||||
return NULL;
|
||||
}
|
||||
*len = uwsgi_rpc(func, argc, argv, argvs, &buffer);
|
||||
if (*buffer)
|
||||
if (buffer)
|
||||
return buffer;
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -243,6 +243,7 @@ SSL_CTX *uwsgi_ssl_new_server_context(char *name, char *crt, char *key, char *ci
|
||||
DH *dh = PEM_read_bio_DHparams(bio, NULL, NULL, NULL);
|
||||
BIO_free(bio);
|
||||
if (dh) {
|
||||
SSL_CTX_set_options(ctx, SSL_OP_SINGLE_DH_USE);
|
||||
SSL_CTX_set_tmp_dh(ctx, dh);
|
||||
DH_free(dh);
|
||||
}
|
||||
@@ -252,6 +253,7 @@ SSL_CTX *uwsgi_ssl_new_server_context(char *name, char *crt, char *key, char *ci
|
||||
#ifdef NID_X9_62_prime256v1
|
||||
EC_KEY *ecdh = EC_KEY_new_by_curve_name(NID_X9_62_prime256v1);
|
||||
if (ecdh) {
|
||||
SSL_CTX_set_options(ctx, SSL_OP_SINGLE_ECDH_USE);
|
||||
SSL_CTX_set_tmp_ecdh(ctx, ecdh);
|
||||
EC_KEY_free(ecdh);
|
||||
}
|
||||
|
||||
@@ -427,6 +427,8 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
|
||||
uwsgi_log("[uwsgi-subscription for pid %d] invalid (sniffed ?) packet sent for slot: %.*s node: %.*s unix_check: %lu\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, (unsigned long) usr->unix_check);
|
||||
return NULL;
|
||||
}
|
||||
// eventually the packet could be upgraded to sni...
|
||||
uwsgi_subscription_sni_check(current_slot, usr);
|
||||
#endif
|
||||
// remove death mark and update cores and load
|
||||
node->death_mark = 0;
|
||||
|
||||
+9
-1
@@ -2198,7 +2198,7 @@ void uwsgi_setup(int argc, char *argv[], char *envp[]) {
|
||||
uwsgi_opt_flock(NULL, uwsgi.flock_wait2, NULL);
|
||||
|
||||
// setup master logging
|
||||
if (uwsgi.log_master && !uwsgi.daemonize2)
|
||||
if (uwsgi.log_master && !uwsgi.daemonize2 && !uwsgi.logto2)
|
||||
uwsgi_setup_log_master();
|
||||
|
||||
// setup offload engines
|
||||
@@ -2481,6 +2481,7 @@ int uwsgi_start(void *v_argv) {
|
||||
if (uwsgi.logto2) {
|
||||
if (!uwsgi.is_a_reload || uwsgi.log_reopen) {
|
||||
logto(uwsgi.logto2);
|
||||
uwsgi_setup_log_master();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3663,6 +3664,13 @@ void uwsgi_init_all_apps() {
|
||||
if (uwsgi.need_app) {
|
||||
if (!uwsgi.lazy)
|
||||
uwsgi_log("*** no app loaded. GAME OVER ***\n");
|
||||
if (uwsgi.lazy_apps) {
|
||||
if (uwsgi.master_process) {
|
||||
if (kill(uwsgi.workers[0].pid, SIGINT)) {
|
||||
uwsgi_error("kill()");
|
||||
}
|
||||
}
|
||||
}
|
||||
exit(UWSGI_FAILED_APP_CODE);
|
||||
}
|
||||
else {
|
||||
|
||||
@@ -0,0 +1,378 @@
|
||||
#include "../python/uwsgi_python.h"
|
||||
|
||||
/*
|
||||
|
||||
python >= 3.4 asyncio (PEP 3156) loop engine
|
||||
|
||||
EXPERIMENTAL !!!
|
||||
|
||||
Author: Roberto De Ioris
|
||||
|
||||
*/
|
||||
|
||||
extern struct uwsgi_server uwsgi;
|
||||
extern struct uwsgi_python up;
|
||||
|
||||
static struct uwsgi_asyncio {
|
||||
PyObject *mod;
|
||||
PyObject *loop;
|
||||
PyObject *request;
|
||||
PyObject *hook_fd;
|
||||
PyObject *hook_timeout;
|
||||
PyObject *hook_fix;
|
||||
} uasyncio;
|
||||
|
||||
#define free_req_queue uwsgi.async_queue_unused_ptr++; uwsgi.async_queue_unused[uwsgi.async_queue_unused_ptr] = wsgi_req
|
||||
|
||||
static void uwsgi_opt_setup_asyncio(char *opt, char *value, void *null) {
|
||||
|
||||
// set async mode
|
||||
uwsgi_opt_set_int(opt, value, &uwsgi.async);
|
||||
if (uwsgi.socket_timeout < 30) {
|
||||
uwsgi.socket_timeout = 30;
|
||||
}
|
||||
// set loop engine
|
||||
uwsgi.loop = "asyncio";
|
||||
|
||||
}
|
||||
|
||||
static struct uwsgi_option asyncio_options[] = {
|
||||
{"asyncio", required_argument, 0, "a shortcut enabling asyncio loop engine with the specified number of async cores and optimal parameters", uwsgi_opt_setup_asyncio, NULL, UWSGI_OPT_THREADS},
|
||||
{0, 0, 0, 0, 0, 0, 0},
|
||||
|
||||
};
|
||||
|
||||
static void gil_asyncio_get() {
|
||||
pthread_setspecific(up.upt_gil_key, (void *) PyGILState_Ensure());
|
||||
}
|
||||
|
||||
static void gil_asyncio_release() {
|
||||
PyGILState_Release((PyGILState_STATE) pthread_getspecific(up.upt_gil_key));
|
||||
}
|
||||
|
||||
static int uwsgi_asyncio_wait_read_hook(int fd, int timeout) {
|
||||
|
||||
struct wsgi_request *wsgi_req = current_wsgi_req();
|
||||
|
||||
if (PyObject_CallMethod(uasyncio.loop, "add_reader", "iOl", fd, uasyncio.hook_fd,(long) wsgi_req) == NULL) {
|
||||
goto error;
|
||||
}
|
||||
|
||||
PyObject *ob_timeout = PyObject_CallMethod(uasyncio.loop, "call_later", "iOl", timeout, uasyncio.hook_timeout, (long)wsgi_req);
|
||||
if (!ob_timeout) {
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_reader", "i", fd) == NULL) PyErr_Print();
|
||||
goto error;
|
||||
}
|
||||
// back to loop
|
||||
if (uwsgi.schedule_to_main) uwsgi.schedule_to_main(wsgi_req);
|
||||
// back from loop
|
||||
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_reader", "i", fd) == NULL) PyErr_Print();
|
||||
if (PyObject_CallMethod(ob_timeout, "cancel", NULL) == NULL) PyErr_Print();
|
||||
|
||||
Py_DECREF(ob_timeout);
|
||||
|
||||
if (wsgi_req->async_timed_out) return 0;
|
||||
|
||||
return 1;
|
||||
|
||||
error:
|
||||
PyErr_Print();
|
||||
return -1;
|
||||
}
|
||||
|
||||
static int uwsgi_asyncio_wait_write_hook(int fd, int timeout) {
|
||||
|
||||
struct wsgi_request *wsgi_req = current_wsgi_req();
|
||||
|
||||
if (PyObject_CallMethod(uasyncio.loop, "add_writer", "iOl", fd, uasyncio.hook_fd,(long) wsgi_req) == NULL) {
|
||||
goto error;
|
||||
}
|
||||
|
||||
PyObject *ob_timeout = PyObject_CallMethod(uasyncio.loop, "call_later", "iOl", timeout, uasyncio.hook_timeout, (long)wsgi_req);
|
||||
if (!ob_timeout) {
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_writer", "i", fd) == NULL) PyErr_Print();
|
||||
goto error;
|
||||
}
|
||||
// back to loop
|
||||
if (uwsgi.schedule_to_main) {
|
||||
uwsgi.schedule_to_main(wsgi_req);
|
||||
}
|
||||
// back from loop
|
||||
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_writer", "i", fd) == NULL) PyErr_Print();
|
||||
if (PyObject_CallMethod(ob_timeout, "cancel", NULL) == NULL) PyErr_Print();
|
||||
|
||||
Py_DECREF(ob_timeout);
|
||||
|
||||
if (wsgi_req->async_timed_out) return 0;
|
||||
|
||||
return 1;
|
||||
|
||||
error:
|
||||
PyErr_Print();
|
||||
return -1;
|
||||
}
|
||||
|
||||
static PyObject *py_uwsgi_asyncio_request(PyObject *self, PyObject *args) {
|
||||
long wsgi_req_ptr = 0;
|
||||
int timed_out = 0;
|
||||
if (!PyArg_ParseTuple(args, "l|i:uwsgi_asyncio_request", &wsgi_req_ptr, &timed_out)) {
|
||||
uwsgi_log_verbose("[BUG] invalid arguments for asyncio callback !!!\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
struct wsgi_request *wsgi_req = (struct wsgi_request *) wsgi_req_ptr;
|
||||
uwsgi.wsgi_req = wsgi_req;
|
||||
|
||||
PyObject *ob_timeout = (PyObject *) wsgi_req->async_timeout;
|
||||
if (PyObject_CallMethod(ob_timeout, "cancel", NULL) == NULL) PyErr_Print();
|
||||
Py_DECREF(ob_timeout);
|
||||
// avoid mess when closing the request
|
||||
wsgi_req->async_timeout = NULL;
|
||||
|
||||
if (timed_out > 0) {
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_reader", "i", wsgi_req->fd) == NULL) PyErr_Print();
|
||||
goto end;
|
||||
}
|
||||
|
||||
int status = wsgi_req->socket->proto(wsgi_req);
|
||||
if (status > 0) {
|
||||
ob_timeout = PyObject_CallMethod(uasyncio.loop, "call_later", "iOli", uwsgi.socket_timeout, uasyncio.request, wsgi_req_ptr, 1);
|
||||
if (!ob_timeout) {
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_reader", "i", wsgi_req->fd) == NULL) PyErr_Print();
|
||||
goto end;
|
||||
}
|
||||
// trick for reference counting
|
||||
wsgi_req->async_timeout = (struct uwsgi_rb_timer *) ob_timeout;
|
||||
goto again;
|
||||
}
|
||||
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_reader", "i", wsgi_req->fd) == NULL) {
|
||||
PyErr_Print();
|
||||
goto end;
|
||||
}
|
||||
|
||||
if (status == 0) {
|
||||
// we call this two time... overengineering :(
|
||||
uwsgi.async_proto_fd_table[wsgi_req->fd] = NULL;
|
||||
uwsgi.schedule_to_req();
|
||||
goto again;
|
||||
}
|
||||
|
||||
end:
|
||||
uwsgi.async_proto_fd_table[wsgi_req->fd] = NULL;
|
||||
uwsgi_close_request(uwsgi.wsgi_req);
|
||||
free_req_queue;
|
||||
again:
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
}
|
||||
|
||||
|
||||
static PyObject *py_uwsgi_asyncio_accept(PyObject *self, PyObject *args) {
|
||||
long uwsgi_sock_ptr = 0;
|
||||
if (!PyArg_ParseTuple(args, "l:uwsgi_asyncio_accept", &uwsgi_sock_ptr)) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
struct wsgi_request *wsgi_req = find_first_available_wsgi_req();
|
||||
|
||||
if (wsgi_req == NULL) {
|
||||
uwsgi_async_queue_is_full(uwsgi_now());
|
||||
goto end;
|
||||
}
|
||||
|
||||
uwsgi.wsgi_req = wsgi_req;
|
||||
struct uwsgi_socket *uwsgi_sock = (struct uwsgi_socket *) uwsgi_sock_ptr;
|
||||
|
||||
// fill wsgi_request structure
|
||||
wsgi_req_setup(wsgi_req, wsgi_req->async_id, uwsgi_sock );
|
||||
|
||||
// mark core as used
|
||||
uwsgi.workers[uwsgi.mywid].cores[wsgi_req->async_id].in_request = 1;
|
||||
|
||||
// accept the connection (since uWSGI 1.5 all of the sockets are non-blocking)
|
||||
if (wsgi_req_simple_accept(wsgi_req, uwsgi_sock->fd)) {
|
||||
// in case of errors (or thundering herd, just reset it)
|
||||
uwsgi.workers[uwsgi.mywid].cores[wsgi_req->async_id].in_request = 0;
|
||||
free_req_queue;
|
||||
goto end;
|
||||
}
|
||||
|
||||
wsgi_req->start_of_request = uwsgi_micros();
|
||||
wsgi_req->start_of_request_in_sec = wsgi_req->start_of_request/1000000;
|
||||
|
||||
// enter harakiri mode
|
||||
if (uwsgi.harakiri_options.workers > 0) {
|
||||
set_harakiri(uwsgi.harakiri_options.workers);
|
||||
}
|
||||
|
||||
uwsgi.async_proto_fd_table[wsgi_req->fd] = wsgi_req;
|
||||
|
||||
// add callback for protocol
|
||||
if (PyObject_CallMethod(uasyncio.loop, "add_reader", "iOl", wsgi_req->fd, uasyncio.request, (long) wsgi_req) == NULL) {
|
||||
free_req_queue;
|
||||
PyErr_Print();
|
||||
}
|
||||
|
||||
// add timeout
|
||||
PyObject *ob_timeout = PyObject_CallMethod(uasyncio.loop, "call_later", "iOli", uwsgi.socket_timeout, uasyncio.request, (long)wsgi_req, 1);
|
||||
if (!ob_timeout) {
|
||||
if (PyObject_CallMethod(uasyncio.loop, "remove_reader", "i", wsgi_req->fd) == NULL) PyErr_Print();
|
||||
free_req_queue;
|
||||
}
|
||||
else {
|
||||
// trick for reference counting
|
||||
wsgi_req->async_timeout = (struct uwsgi_rb_timer *) ob_timeout;
|
||||
}
|
||||
end:
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_asyncio_hook_fd(PyObject *self, PyObject *args) {
|
||||
long wsgi_req_ptr = 0;
|
||||
if (!PyArg_ParseTuple(args, "l:uwsgi_asyncio_hook_fd", &wsgi_req_ptr)) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
uwsgi.wsgi_req = (struct wsgi_request *) wsgi_req_ptr;
|
||||
uwsgi.schedule_to_req();
|
||||
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_asyncio_hook_timeout(PyObject *self, PyObject *args) {
|
||||
long wsgi_req_ptr = 0;
|
||||
if (!PyArg_ParseTuple(args, "l", &wsgi_req_ptr)) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
uwsgi.wsgi_req = (struct wsgi_request *) wsgi_req_ptr;
|
||||
uwsgi.wsgi_req->async_timed_out = 1;
|
||||
uwsgi.schedule_to_req();
|
||||
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_asyncio_hook_fix(PyObject *self, PyObject *args) {
|
||||
long wsgi_req_ptr = 0;
|
||||
if (!PyArg_ParseTuple(args, "l", &wsgi_req_ptr)) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
uwsgi.wsgi_req = (struct wsgi_request *) wsgi_req_ptr;
|
||||
uwsgi.schedule_to_req();
|
||||
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
}
|
||||
|
||||
|
||||
PyMethodDef uwsgi_asyncio_accept_def[] = { {"uwsgi_asyncio_accept", py_uwsgi_asyncio_accept, METH_VARARGS, ""} };
|
||||
PyMethodDef uwsgi_asyncio_request_def[] = { {"uwsgi_asyncio_request", py_uwsgi_asyncio_request, METH_VARARGS, ""} };
|
||||
PyMethodDef uwsgi_asyncio_hook_fd_def[] = { {"uwsgi_asyncio_hook_fd", py_uwsgi_asyncio_hook_fd, METH_VARARGS, ""} };
|
||||
PyMethodDef uwsgi_asyncio_hook_timeout_def[] = { {"uwsgi_asyncio_hook_timeout", py_uwsgi_asyncio_hook_timeout, METH_VARARGS, ""} };
|
||||
PyMethodDef uwsgi_asyncio_hook_fix_def[] = { {"uwsgi_asyncio_hook_fix", py_uwsgi_asyncio_hook_fix, METH_VARARGS, ""} };
|
||||
|
||||
static void uwsgi_asyncio_schedule_fix(struct wsgi_request *wsgi_req) {
|
||||
PyObject *cb = PyObject_CallMethod(uasyncio.loop, "call_soon", "Ol", uasyncio.hook_fix, (long) wsgi_req);
|
||||
if (!cb) goto error;
|
||||
Py_DECREF(cb);
|
||||
return;
|
||||
|
||||
error:
|
||||
PyErr_Print();
|
||||
}
|
||||
|
||||
static void asyncio_loop() {
|
||||
|
||||
if (!uwsgi.has_threads && uwsgi.mywid == 1) {
|
||||
uwsgi_log("!!! Running asyncio without threads IS NOT recommended, enable them with --enable-threads !!!\n");
|
||||
}
|
||||
|
||||
if (uwsgi.socket_timeout < 30) {
|
||||
uwsgi_log("!!! Running asyncio with a socket-timeout lower than 30 seconds is not recommended, tune it with --socket-timeout !!!\n");
|
||||
}
|
||||
|
||||
if (!uwsgi.async_waiting_fd_table)
|
||||
uwsgi.async_waiting_fd_table = uwsgi_calloc(sizeof(struct wsgi_request *) * uwsgi.max_fd);
|
||||
if (!uwsgi.async_proto_fd_table)
|
||||
uwsgi.async_proto_fd_table = uwsgi_calloc(sizeof(struct wsgi_request *) * uwsgi.max_fd);
|
||||
|
||||
// get the GIL
|
||||
UWSGI_GET_GIL
|
||||
|
||||
up.gil_get = gil_asyncio_get;
|
||||
up.gil_release = gil_asyncio_release;
|
||||
|
||||
uwsgi.wait_write_hook = uwsgi_asyncio_wait_write_hook;
|
||||
uwsgi.wait_read_hook = uwsgi_asyncio_wait_read_hook;
|
||||
|
||||
uwsgi.schedule_fix = uwsgi_asyncio_schedule_fix;
|
||||
|
||||
if (uwsgi.async < 2) {
|
||||
uwsgi_log("the asyncio loop engine requires async mode (--async <n>)\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
if (!uwsgi.schedule_to_main) {
|
||||
uwsgi_log("*** DANGER *** asyncio mode without coroutine/greenthread engine loaded !!!\n");
|
||||
}
|
||||
|
||||
PyObject *asyncio = PyImport_ImportModule("asyncio");
|
||||
if (!asyncio) uwsgi_pyexit;
|
||||
|
||||
uasyncio.mod = asyncio;
|
||||
|
||||
uasyncio.loop = PyObject_CallMethod(asyncio, "get_event_loop", NULL);
|
||||
if (!uasyncio.loop) uwsgi_pyexit;
|
||||
|
||||
// main greenlet waiting for connection (one greenlet per-socket)
|
||||
PyObject *asyncio_accept = PyCFunction_New(uwsgi_asyncio_accept_def, NULL);
|
||||
Py_INCREF(asyncio_accept);
|
||||
|
||||
uasyncio.request = PyCFunction_New(uwsgi_asyncio_request_def, NULL);
|
||||
if (!uasyncio.request) uwsgi_pyexit;
|
||||
|
||||
uasyncio.hook_fd = PyCFunction_New(uwsgi_asyncio_hook_fd_def, NULL);
|
||||
if (!uasyncio.hook_fd) uwsgi_pyexit;
|
||||
uasyncio.hook_timeout = PyCFunction_New(uwsgi_asyncio_hook_timeout_def, NULL);
|
||||
if (!uasyncio.hook_timeout) uwsgi_pyexit;
|
||||
uasyncio.hook_fix = PyCFunction_New(uwsgi_asyncio_hook_fix_def, NULL);
|
||||
if (!uasyncio.hook_fix) uwsgi_pyexit;
|
||||
|
||||
Py_INCREF(uasyncio.request);
|
||||
Py_INCREF(uasyncio.hook_fd);
|
||||
Py_INCREF(uasyncio.hook_timeout);
|
||||
Py_INCREF(uasyncio.hook_fix);
|
||||
|
||||
// call add_handler on each socket
|
||||
struct uwsgi_socket *uwsgi_sock = uwsgi.sockets;
|
||||
while(uwsgi_sock) {
|
||||
if (PyObject_CallMethod(uasyncio.loop, "add_reader", "iOi", uwsgi_sock->fd, asyncio_accept, (long) uwsgi_sock) == NULL) {
|
||||
uwsgi_pyexit;
|
||||
}
|
||||
uwsgi_sock = uwsgi_sock->next;
|
||||
}
|
||||
|
||||
if (PyObject_CallMethod(uasyncio.loop, "run_forever", NULL) == NULL) {
|
||||
uwsgi_pyexit;
|
||||
}
|
||||
|
||||
// never here ?
|
||||
}
|
||||
|
||||
static void asyncio_init() {
|
||||
uwsgi_register_loop( (char *) "asyncio", asyncio_loop);
|
||||
}
|
||||
|
||||
|
||||
struct uwsgi_plugin asyncio_plugin = {
|
||||
.name = "asyncio",
|
||||
.options = asyncio_options,
|
||||
.on_load = asyncio_init,
|
||||
};
|
||||
@@ -0,0 +1,8 @@
|
||||
from distutils import sysconfig
|
||||
|
||||
NAME='asyncio'
|
||||
CFLAGS = ['-I' + sysconfig.get_python_inc(), '-I' + sysconfig.get_python_inc(plat_specific=True)]
|
||||
LDFLAGS = []
|
||||
LIBS = []
|
||||
|
||||
GCC_LIST = ['asyncio']
|
||||
@@ -45,6 +45,7 @@ struct corerouter_peer *uwsgi_cr_peer_add(struct corerouter_session *cs) {
|
||||
if (!bufsize) bufsize = uwsgi.page_size;
|
||||
peers->in = uwsgi_buffer_new(bufsize);
|
||||
// add timeout
|
||||
peers->current_timeout = cs->corerouter->socket_timeout;
|
||||
peers->timeout = cr_add_timeout(cs->corerouter, peers);
|
||||
peers->prev = old_peers;
|
||||
|
||||
@@ -582,8 +583,8 @@ struct corerouter_session *corerouter_alloc_session(struct uwsgi_corerouter *ucr
|
||||
cs->corerouter = ucr;
|
||||
cs->ugs = ugs;
|
||||
|
||||
// set initial timeout
|
||||
peer->timeout = cr_add_timeout(ucr, ucr->cr_table[new_connection]);
|
||||
// set initial timeout (could be overridden)
|
||||
peer->current_timeout = ucr->socket_timeout;
|
||||
|
||||
ucr->active_sessions++;
|
||||
|
||||
@@ -622,6 +623,10 @@ struct corerouter_session *corerouter_alloc_session(struct uwsgi_corerouter *ucr
|
||||
corerouter_close_session(ucr, cs);
|
||||
cs = NULL;
|
||||
}
|
||||
else {
|
||||
// truly set the timeout
|
||||
peer->timeout = cr_add_timeout(ucr, ucr->cr_table[new_connection]);
|
||||
}
|
||||
|
||||
return cs;
|
||||
}
|
||||
|
||||
@@ -3,12 +3,12 @@
|
||||
#define COREROUTER_STATUS_RECV_HDR 2
|
||||
#define COREROUTER_STATUS_RESPONSE 3
|
||||
|
||||
#define cr_add_timeout(u, x) uwsgi_add_rb_timer(u->timeouts, uwsgi_now()+u->socket_timeout, x)
|
||||
#define cr_add_timeout_fast(u, x, t) uwsgi_add_rb_timer(u->timeouts, t+u->socket_timeout, x)
|
||||
#define cr_add_timeout(u, x) uwsgi_add_rb_timer(u->timeouts, uwsgi_now()+x->current_timeout, x)
|
||||
#define cr_add_timeout_fast(u, x, t) uwsgi_add_rb_timer(u->timeouts, t+x->current_timeout, x)
|
||||
#define cr_del_timeout(u, x) uwsgi_del_rb_timer(u->timeouts, x->timeout); free(x->timeout);
|
||||
|
||||
#define uwsgi_cr_error(x, y) uwsgi_log("[uwsgi-%s key: %.*s client_addr: %s client_port: %s] %s: %s [%s line %d]\n", x->session->corerouter->short_name, x->session->main_peer ? x->session->main_peer->key_len : 0, x->session->main_peer ? x->session->main_peer->key: "", x->session->client_address, x->session->client_port, y, strerror(errno), __FILE__, __LINE__)
|
||||
#define uwsgi_cr_log(x, y, ...) uwsgi_log("[uwsgi-%s key: %.*s client_addr: %s client_port: %s]" y, x->session->corerouter->short_name, x->session->main_peer ? x->session->main_peer->key_len : 0, x->session->main_peer ? x->session->main_peer->key : "", x->session->client_address, x->session->client_port, __VA_ARGS__)
|
||||
#define uwsgi_cr_error(x, y) uwsgi_log("[uwsgi-%s key: %.*s client_addr: %s client_port: %s] %s: %s [%s line %d]\n", x->session->corerouter->short_name, (x == x->session->main_peer) ? (x->session->peers ? x->session->peers->key_len: 0) : x->key_len, (x == x->session->main_peer) ? (x->session->peers ? x->session->peers->key: "") : x->key, x->session->client_address, x->session->client_port, y, strerror(errno), __FILE__, __LINE__)
|
||||
#define uwsgi_cr_log(x, y, ...) uwsgi_log("[uwsgi-%s key: %.*s client_addr: %s client_port: %s]" y, x->session->corerouter->short_name, (x == x->session->main_peer) ? (x->session->peers ? x->session->peers->key_len: 0) : x->key_len, (x == x->session->main_peer) ? (x->session->peers ? x->session->peers->key: "") : x->key, x->session->client_address, x->session->client_port, __VA_ARGS__)
|
||||
|
||||
#define cr_try_again if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) {\
|
||||
errno = EINPROGRESS;\
|
||||
@@ -188,6 +188,8 @@ struct corerouter_peer {
|
||||
|
||||
struct corerouter_peer *prev;
|
||||
struct corerouter_peer *next;
|
||||
|
||||
int current_timeout;
|
||||
};
|
||||
|
||||
struct uwsgi_corerouter {
|
||||
|
||||
@@ -10,6 +10,10 @@ struct uwsgi_emperor_mongodb_state {
|
||||
char *address;
|
||||
char *collection;
|
||||
char *json;
|
||||
char *database;
|
||||
char *username;
|
||||
char *password;
|
||||
char *predigest;
|
||||
};
|
||||
|
||||
|
||||
@@ -29,6 +33,14 @@ extern "C" void uwsgi_imperial_monitor_mongodb(struct uwsgi_emperor_scanner *ues
|
||||
// connect
|
||||
c.connect(uems->address);
|
||||
|
||||
if (uems->database && uems->username && uems->password) {
|
||||
std::string err;
|
||||
if (c.auth(uems->database, uems->username, uems->password, err, uems->predigest ? false : true) == false) {
|
||||
uwsgi_log_verbose("[emperor-mongodb] unabel to authenticate to db %s: %s\n", uems->database, err.c_str());
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
// run the query
|
||||
std::auto_ptr<mongo::DBClientCursor> cursor = c.query(uems->collection, q, 0, 0, &p);
|
||||
while(cursor.get() && cursor->more() ) {
|
||||
@@ -124,3 +136,40 @@ done:
|
||||
uwsgi_log("[emperor] enabled emperor MongoDB monitor for %s on collection %s\n", uems->address, uems->collection);
|
||||
}
|
||||
|
||||
// setup a new mongodb imperial monitor (keyval based)
|
||||
extern "C" void uwsgi_imperial_monitor_mongodb_init2(struct uwsgi_emperor_scanner *ues) {
|
||||
|
||||
// allocate a new state
|
||||
ues->data = uwsgi_calloc(sizeof(struct uwsgi_emperor_mongodb_state));
|
||||
size_t arg_len = strlen(ues->arg);
|
||||
struct uwsgi_emperor_mongodb_state *uems = (struct uwsgi_emperor_mongodb_state *) ues->data;
|
||||
|
||||
// parse args/ set defaults
|
||||
uems->address = (char *) "127.0.0.1:27017";
|
||||
uems->collection = (char *) "uwsgi.emperor.vassals";
|
||||
uems->json = (char *) "";
|
||||
char *args = NULL;
|
||||
if (arg_len <= 11) goto done;
|
||||
args = ues->arg+11;
|
||||
if (uwsgi_kvlist_parse(args, strlen(args), ',', '=',
|
||||
"addr", &uems->address,
|
||||
"address", &uems->address,
|
||||
"server", &uems->address,
|
||||
"collection", &uems->collection,
|
||||
"coll", &uems->collection,
|
||||
"json", &uems->json,
|
||||
"database", &uems->database,
|
||||
"db", &uems->database,
|
||||
"username", &uems->username,
|
||||
"password", &uems->password,
|
||||
"predigest", &uems->predigest,
|
||||
NULL)) {
|
||||
|
||||
uwsgi_log("[emperor-mongodb] invalid keyval syntax !\n");
|
||||
exit(1);
|
||||
}
|
||||
done:
|
||||
uwsgi_log("[emperor] enabled emperor MongoDB monitor for %s on collection %s\n", uems->address, uems->collection);
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -2,9 +2,11 @@
|
||||
|
||||
void uwsgi_imperial_monitor_mongodb(struct uwsgi_emperor_scanner *);
|
||||
void uwsgi_imperial_monitor_mongodb_init(struct uwsgi_emperor_scanner *);
|
||||
void uwsgi_imperial_monitor_mongodb_init2(struct uwsgi_emperor_scanner *);
|
||||
|
||||
void emperor_mongodb_init(void) {
|
||||
uwsgi_register_imperial_monitor("mongodb", uwsgi_imperial_monitor_mongodb_init, uwsgi_imperial_monitor_mongodb);
|
||||
uwsgi_register_imperial_monitor("mongodb2", uwsgi_imperial_monitor_mongodb_init2, uwsgi_imperial_monitor_mongodb);
|
||||
}
|
||||
|
||||
struct uwsgi_plugin emperor_mongodb_plugin = {
|
||||
|
||||
@@ -4,6 +4,8 @@
|
||||
extern struct uwsgi_server uwsgi;
|
||||
extern struct uwsgi_python up;
|
||||
|
||||
#define is_not_python(x) strcmp(uwsgi.p[x]->name, "python")
|
||||
|
||||
struct ugreenlet {
|
||||
int enabled;
|
||||
PyObject *callable;
|
||||
@@ -69,7 +71,7 @@ static void greenlet_schedule_to_req() {
|
||||
}
|
||||
|
||||
// call it in the main core
|
||||
if (uwsgi.p[modifier1]->suspend) {
|
||||
if (is_not_python(modifier1) && uwsgi.p[modifier1]->suspend) {
|
||||
uwsgi.p[modifier1]->suspend(NULL);
|
||||
}
|
||||
|
||||
@@ -81,7 +83,7 @@ static void greenlet_schedule_to_req() {
|
||||
}
|
||||
Py_DECREF(ret);
|
||||
|
||||
if (uwsgi.p[modifier1]->resume) {
|
||||
if (is_not_python(modifier1) && uwsgi.p[modifier1]->resume) {
|
||||
uwsgi.p[modifier1]->resume(NULL);
|
||||
}
|
||||
|
||||
@@ -93,7 +95,7 @@ static void greenlet_schedule_to_main(struct wsgi_request *wsgi_req) {
|
||||
// ensure gil
|
||||
UWSGI_GET_GIL
|
||||
|
||||
if (uwsgi.p[wsgi_req->uh->modifier1]->suspend) {
|
||||
if (is_not_python(wsgi_req->uh->modifier1) && uwsgi.p[wsgi_req->uh->modifier1]->suspend) {
|
||||
uwsgi.p[wsgi_req->uh->modifier1]->suspend(wsgi_req);
|
||||
}
|
||||
PyObject *ret = PyGreenlet_Switch(ugl.main, NULL, NULL);
|
||||
@@ -103,7 +105,7 @@ static void greenlet_schedule_to_main(struct wsgi_request *wsgi_req) {
|
||||
exit(1);
|
||||
}
|
||||
Py_DECREF(ret);
|
||||
if (uwsgi.p[wsgi_req->uh->modifier1]->resume) {
|
||||
if (is_not_python(wsgi_req->uh->modifier1) && uwsgi.p[wsgi_req->uh->modifier1]->resume) {
|
||||
uwsgi.p[wsgi_req->uh->modifier1]->resume(wsgi_req);
|
||||
}
|
||||
uwsgi.wsgi_req = wsgi_req;
|
||||
|
||||
@@ -42,6 +42,8 @@ struct uwsgi_http {
|
||||
|
||||
int server_name_as_http_host;
|
||||
|
||||
int headers_timeout;
|
||||
int connect_timeout;
|
||||
};
|
||||
|
||||
struct http_session {
|
||||
|
||||
+120
-42
@@ -60,43 +60,56 @@ struct uwsgi_option http_options[] = {
|
||||
{"http-buffer-size", required_argument, 0, "set internal buffer size (default: page size)", uwsgi_opt_set_64bit, &uhttp.cr.buffer_size, 0},
|
||||
|
||||
{"http-server-name-as-http-host", required_argument, 0, "force SERVER_NAME to HTTP_HOST", uwsgi_opt_true, &uhttp.server_name_as_http_host, 0},
|
||||
{"http-headers-timeout", required_argument, 0, "set internal http socket timeout for headers", uwsgi_opt_set_int, &uhttp.headers_timeout, 0},
|
||||
{"http-connect-timeout", required_argument, 0, "set internal http socket timeout for backend connections", uwsgi_opt_set_int, &uhttp.connect_timeout, 0},
|
||||
{0, 0, 0, 0, 0, 0, 0},
|
||||
};
|
||||
|
||||
int http_add_uwsgi_header(struct corerouter_peer *peer, char *hh, uint16_t hhlen) {
|
||||
static void http_set_timeout(struct corerouter_peer *peer, int timeout) {
|
||||
if (peer->current_timeout == timeout) return;
|
||||
peer->current_timeout = timeout;
|
||||
peer->timeout = corerouter_reset_timeout(peer->session->corerouter, peer);
|
||||
}
|
||||
|
||||
static char * http_header_to_cgi(char *hh, size_t hhlen, size_t *keylen, size_t *vallen, int *has_prefix) {
|
||||
size_t i;
|
||||
char *val = hh;
|
||||
int status = 0;
|
||||
for (i = 0; i < hhlen; i++) {
|
||||
if (!status) {
|
||||
hh[i] = toupper((int) hh[i]);
|
||||
if (hh[i] == '-')
|
||||
hh[i] = '_';
|
||||
if (hh[i] == ':') {
|
||||
status = 1;
|
||||
*keylen = i;
|
||||
}
|
||||
}
|
||||
else if (status == 1 && hh[i] != ' ') {
|
||||
status = 2;
|
||||
val += i;
|
||||
*vallen+=1;
|
||||
}
|
||||
else if (status == 2) {
|
||||
*vallen+=1;
|
||||
}
|
||||
}
|
||||
|
||||
if (!(*keylen))
|
||||
return NULL;
|
||||
|
||||
if (uwsgi_strncmp("CONTENT_LENGTH", 14, hh, *keylen) && uwsgi_strncmp("CONTENT_TYPE", 12, hh, *keylen)) {
|
||||
*has_prefix = 0x02;
|
||||
}
|
||||
|
||||
return val;
|
||||
}
|
||||
|
||||
static int http_add_uwsgi_header(struct corerouter_peer *peer, char *hh, size_t keylen, char *val, size_t vallen, int prefix) {
|
||||
|
||||
struct uwsgi_buffer *out = peer->out;
|
||||
struct http_session *hr = (struct http_session *) peer->session;
|
||||
|
||||
int i;
|
||||
int status = 0;
|
||||
char *val = hh;
|
||||
uint16_t keylen = 0, vallen = 0;
|
||||
int prefix = 0;
|
||||
|
||||
for (i = 0; i < hhlen; i++) {
|
||||
if (!status) {
|
||||
hh[i] = toupper((int) hh[i]);
|
||||
if (hh[i] == '-')
|
||||
hh[i] = '_';
|
||||
if (hh[i] == ':') {
|
||||
status = 1;
|
||||
keylen = i;
|
||||
}
|
||||
}
|
||||
else if (status == 1 && hh[i] != ' ') {
|
||||
status = 2;
|
||||
val += i;
|
||||
vallen++;
|
||||
}
|
||||
else if (status == 2) {
|
||||
vallen++;
|
||||
}
|
||||
}
|
||||
|
||||
if (!keylen)
|
||||
return -1;
|
||||
|
||||
if (hr->websockets) {
|
||||
if (!uwsgi_strncmp("UPGRADE", 7, hh, keylen)) {
|
||||
if (!uwsgi_strnicmp(val, vallen, "websocket", 9)) {
|
||||
@@ -151,11 +164,7 @@ int http_add_uwsgi_header(struct corerouter_peer *peer, char *hh, uint16_t hhlen
|
||||
#endif
|
||||
|
||||
done:
|
||||
|
||||
if (uwsgi_strncmp("CONTENT_TYPE", 12, hh, keylen) && uwsgi_strncmp("CONTENT_LENGTH", 14, hh, keylen)) {
|
||||
keylen += 5;
|
||||
prefix = 1;
|
||||
}
|
||||
if (prefix) keylen += 5;
|
||||
|
||||
if (uwsgi_buffer_u16le(out, keylen)) return -1;
|
||||
|
||||
@@ -163,7 +172,7 @@ done:
|
||||
if (uwsgi_buffer_append(out, "HTTP_", 5)) return -1;
|
||||
}
|
||||
|
||||
if (uwsgi_buffer_append(out, hh, keylen - (prefix * 5))) return -1;
|
||||
if (uwsgi_buffer_append(out, hh, keylen - (prefix ? 5 : 0))) return -1;
|
||||
|
||||
if (uwsgi_buffer_u16le(out, vallen)) return -1;
|
||||
if (uwsgi_buffer_append(out, val, vallen)) return -1;
|
||||
@@ -325,6 +334,8 @@ int http_headers_parse(struct corerouter_peer *peer) {
|
||||
//HEADERS
|
||||
base = ptr;
|
||||
|
||||
struct uwsgi_string_list *headers = NULL, *usl = NULL;
|
||||
|
||||
while (ptr < watermark) {
|
||||
if (*ptr == '\r') {
|
||||
if (ptr + 1 >= watermark)
|
||||
@@ -345,13 +356,53 @@ int http_headers_parse(struct corerouter_peer *peer) {
|
||||
hr->send_expect_100 = 1;
|
||||
}
|
||||
}
|
||||
if (http_add_uwsgi_header(peer, base, ptr - base)) return -1;
|
||||
|
||||
size_t key_len = 0, value_len = 0;
|
||||
int has_prefix = 0;
|
||||
// last line, do not waste time
|
||||
if (ptr - base == 0) break;
|
||||
char *value = http_header_to_cgi(base, ptr - base, &key_len, &value_len, &has_prefix);
|
||||
if (!value) goto clear;
|
||||
usl = uwsgi_string_list_has_item(headers, base, key_len);
|
||||
// there is already a HTTP header with the same name, let's merge them
|
||||
if (usl) {
|
||||
char *old_value = usl->custom_ptr;
|
||||
usl->custom_ptr = uwsgi_concat3n(old_value, (size_t) usl->custom, ", ", 2, value, value_len);
|
||||
usl->custom += 2 + value_len;
|
||||
if (usl->custom2 & 0x01) free(old_value);
|
||||
usl->custom2 |= 0x01;
|
||||
}
|
||||
else {
|
||||
// add an entry
|
||||
usl = uwsgi_string_new_list(&headers, NULL);
|
||||
usl->value = base;
|
||||
usl->len = key_len;
|
||||
usl->custom_ptr = value;
|
||||
usl->custom = value_len;
|
||||
usl->custom2 = has_prefix;
|
||||
}
|
||||
ptr++;
|
||||
base = ptr + 1;
|
||||
}
|
||||
ptr++;
|
||||
}
|
||||
|
||||
usl = headers;
|
||||
int broken = 0;
|
||||
while(usl) {
|
||||
if (!broken) {
|
||||
if (http_add_uwsgi_header(peer, usl->value, usl->len, usl->custom_ptr, (size_t) usl->custom, usl->custom2 & 0x02)) broken = 1;
|
||||
}
|
||||
if (usl->custom2 & 0x01) {
|
||||
free(usl->custom_ptr);
|
||||
}
|
||||
struct uwsgi_string_list *tmp_usl = usl;
|
||||
usl = usl->next;
|
||||
free(tmp_usl);
|
||||
}
|
||||
|
||||
if (broken) return -1;
|
||||
|
||||
struct uwsgi_string_list *hv = uhttp.http_vars;
|
||||
while (hv) {
|
||||
char *equal = strchr(hv->value, '=');
|
||||
@@ -363,6 +414,18 @@ int http_headers_parse(struct corerouter_peer *peer) {
|
||||
|
||||
return 0;
|
||||
|
||||
clear:
|
||||
usl = headers;
|
||||
while(usl) {
|
||||
if (usl->custom2 & 0x01) {
|
||||
free(usl->custom_ptr);
|
||||
}
|
||||
struct uwsgi_string_list *tmp_usl = usl;
|
||||
usl = usl->next;
|
||||
free(tmp_usl);
|
||||
}
|
||||
return -1;
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -446,6 +509,7 @@ ssize_t hr_write(struct corerouter_peer *main_peer) {
|
||||
return 0;
|
||||
}
|
||||
if (main_peer->session->connect_peer_after_write) {
|
||||
http_set_timeout(main_peer->session->connect_peer_after_write, uhttp.connect_timeout);
|
||||
cr_connect(main_peer->session->connect_peer_after_write, hr_instance_connected);
|
||||
main_peer->session->connect_peer_after_write = NULL;
|
||||
return len;
|
||||
@@ -459,6 +523,9 @@ ssize_t hr_write(struct corerouter_peer *main_peer) {
|
||||
ssize_t hr_instance_connected(struct corerouter_peer* peer) {
|
||||
|
||||
cr_peer_connected(peer, "hr_instance_connected()");
|
||||
|
||||
// set the default timeout
|
||||
http_set_timeout(peer, uhttp.cr.socket_timeout);
|
||||
|
||||
// we are connected, we cannot retry anymore
|
||||
peer->can_retry = 0;
|
||||
@@ -523,10 +590,7 @@ ssize_t hr_instance_read(struct corerouter_peer *peer) {
|
||||
hr->has_gzip = 0;
|
||||
#endif
|
||||
if (uhttp.keepalive > 1) {
|
||||
int orig_timeout = peer->session->corerouter->socket_timeout;
|
||||
peer->session->corerouter->socket_timeout = uhttp.keepalive;
|
||||
peer->session->main_peer->timeout = corerouter_reset_timeout(peer->session->corerouter, peer->session->main_peer);
|
||||
peer->session->corerouter->socket_timeout = orig_timeout;
|
||||
http_set_timeout(peer->session->main_peer, uhttp.keepalive);
|
||||
}
|
||||
}
|
||||
#ifdef UWSGI_ZLIB
|
||||
@@ -655,6 +719,9 @@ ssize_t http_parse(struct corerouter_peer *main_peer) {
|
||||
return 1;
|
||||
}
|
||||
|
||||
// ensure the headers timeout is honoured
|
||||
http_set_timeout(main_peer, uhttp.headers_timeout);
|
||||
|
||||
// read until \r\n\r\n is found
|
||||
size_t j;
|
||||
size_t len = main_peer->in->pos;
|
||||
@@ -744,6 +811,10 @@ ssize_t http_parse(struct corerouter_peer *main_peer) {
|
||||
hr->raw_body = 1;
|
||||
}
|
||||
new_peer->can_retry = 1;
|
||||
// reset main timeout
|
||||
http_set_timeout(main_peer, uhttp.cr.socket_timeout);
|
||||
// set peer timeout
|
||||
http_set_timeout(new_peer, uhttp.connect_timeout);
|
||||
cr_connect(new_peer, hr_instance_connected);
|
||||
break;
|
||||
}
|
||||
@@ -830,12 +901,17 @@ static int hr_retry(struct corerouter_peer *peer) {
|
||||
|
||||
retry:
|
||||
// start async connect (again)
|
||||
http_set_timeout(peer, uhttp.connect_timeout);
|
||||
cr_connect(peer, hr_instance_connected);
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
int http_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) {
|
||||
|
||||
if (!uhttp.headers_timeout) uhttp.headers_timeout = uhttp.cr.socket_timeout;
|
||||
if (!uhttp.connect_timeout) uhttp.connect_timeout = uhttp.cr.socket_timeout;
|
||||
|
||||
// set the retry hook
|
||||
cs->retry = hr_retry;
|
||||
struct http_session *hr = (struct http_session *) cs;
|
||||
@@ -845,6 +921,9 @@ int http_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket
|
||||
// default hook
|
||||
cs->main_peer->last_hook_read = hr_read;
|
||||
|
||||
// headers timeout
|
||||
cs->main_peer->current_timeout = uhttp.headers_timeout;
|
||||
|
||||
if (uhttp.raw_body) {
|
||||
hr->raw_body = 1;
|
||||
}
|
||||
@@ -906,7 +985,6 @@ int http_init() {
|
||||
uhttp.cr.socket_num = 0;
|
||||
}
|
||||
uwsgi_corerouter_init((struct uwsgi_corerouter *) &uhttp);
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
+46
-5
@@ -180,8 +180,10 @@ struct uwsgi_buffer *spdy_http_to_spdy(char *buf, size_t len, uint32_t *hh) {
|
||||
if (!key) return ub;
|
||||
|
||||
uint32_t h_len = 0;
|
||||
|
||||
// headers (key lowercase...)
|
||||
// merge header values and ensure keys are all lowercase
|
||||
struct uwsgi_string_list *hr=NULL, *usl=NULL;
|
||||
char *line_value;
|
||||
size_t line_key_len, line_value_len;
|
||||
for(i=next;i<len;i++) {
|
||||
if (key) {
|
||||
if (buf[i] == '\r' || buf[i] == '\n') {
|
||||
@@ -192,10 +194,37 @@ struct uwsgi_buffer *spdy_http_to_spdy(char *buf, size_t len, uint32_t *hh) {
|
||||
// tolower !!!
|
||||
size_t j;
|
||||
for(j=0;j<h_len;j++) {
|
||||
key[j] = tolower((int) key[j]);
|
||||
// don't lowercase values, only keys
|
||||
if (key[j] == ':') break;
|
||||
key[j] = tolower((int) key[j]);
|
||||
}
|
||||
line_key_len = colon - key;
|
||||
key[line_key_len] = 0;
|
||||
line_value_len = h_len - 2 - line_key_len;
|
||||
line_value = uwsgi_strncopy(colon+2, line_value_len);
|
||||
if (hr) {
|
||||
// check if we already store values for this key
|
||||
usl = uwsgi_string_list_has_item(hr, key, line_key_len);
|
||||
if (usl) {
|
||||
// we have this key, append new value
|
||||
char *oldval = usl->custom_ptr;
|
||||
usl->custom_ptr = uwsgi_concat3n(usl->custom_ptr, usl->custom, "\0", 1, line_value, line_value_len);
|
||||
usl->custom = usl->custom + 1 + line_value_len;
|
||||
free(oldval);
|
||||
}
|
||||
else {
|
||||
// this key is new
|
||||
usl = uwsgi_string_new_list(&hr, key);
|
||||
usl->custom_ptr = line_value;
|
||||
usl->custom = line_value_len;
|
||||
}
|
||||
}
|
||||
else {
|
||||
// this is first key
|
||||
usl = uwsgi_string_new_list(&hr, key);
|
||||
usl->custom_ptr = line_value;
|
||||
usl->custom = line_value_len;
|
||||
}
|
||||
if (uwsgi_buffer_append_keyval32(ub, key, colon-key, colon+2, h_len-((colon-key)+2))) goto end;
|
||||
*hh+=1;
|
||||
key = NULL;
|
||||
h_len = 0;
|
||||
}
|
||||
@@ -210,6 +239,18 @@ struct uwsgi_buffer *spdy_http_to_spdy(char *buf, size_t len, uint32_t *hh) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// append all merged header lines to buffer and free memory
|
||||
struct uwsgi_string_list *ohr;
|
||||
while (hr) {
|
||||
if (uwsgi_buffer_append_keyval32(ub, hr->value, hr->len, hr->custom_ptr, hr->custom)) goto end;
|
||||
*hh+=1;
|
||||
ohr = hr;
|
||||
hr = hr->next;
|
||||
free(ohr->custom_ptr);
|
||||
free(ohr);
|
||||
}
|
||||
|
||||
return ub;
|
||||
|
||||
end:
|
||||
|
||||
@@ -0,0 +1,12 @@
|
||||
#include "php.h"
|
||||
#include "SAPI.h"
|
||||
#include "php_main.h"
|
||||
#include "php_variables.h"
|
||||
|
||||
#include "ext/standard/php_smart_str.h"
|
||||
#include "ext/standard/info.h"
|
||||
|
||||
#include "ext/session/php_session.h"
|
||||
|
||||
#include <uwsgi.h>
|
||||
|
||||
@@ -1,12 +1,4 @@
|
||||
#include "php.h"
|
||||
#include "SAPI.h"
|
||||
#include "php_main.h"
|
||||
#include "php_variables.h"
|
||||
|
||||
#include "ext/standard/php_smart_str.h"
|
||||
#include "ext/standard/info.h"
|
||||
|
||||
#include "../../uwsgi.h"
|
||||
#include "common.h"
|
||||
|
||||
extern struct uwsgi_server uwsgi;
|
||||
|
||||
@@ -233,8 +225,9 @@ void uwsgi_php_set(char *opt) {
|
||||
uwsgi_sapi_module.ini_entries[uphp.ini_size] = 0;
|
||||
}
|
||||
|
||||
// future implementation...
|
||||
extern ps_module ps_mod_uwsgi;
|
||||
PHP_MINIT_FUNCTION(uwsgi_php_minit) {
|
||||
php_session_register_module(&ps_mod_uwsgi);
|
||||
return SUCCESS;
|
||||
}
|
||||
|
||||
@@ -341,7 +334,7 @@ PHP_FUNCTION(uwsgi_cache_set) {
|
||||
RETURN_NULL();
|
||||
}
|
||||
|
||||
if (uwsgi_cache_magic_set(key, keylen, value, vallen, expires, 0, cache)) {
|
||||
if (!uwsgi_cache_magic_set(key, keylen, value, vallen, expires, 0, cache)) {
|
||||
RETURN_TRUE;
|
||||
}
|
||||
RETURN_NULL();
|
||||
@@ -363,7 +356,7 @@ PHP_FUNCTION(uwsgi_cache_update) {
|
||||
RETURN_NULL();
|
||||
}
|
||||
|
||||
if (uwsgi_cache_magic_set(key, keylen, value, vallen, expires, UWSGI_CACHE_FLAG_UPDATE, cache)) {
|
||||
if (!uwsgi_cache_magic_set(key, keylen, value, vallen, expires, UWSGI_CACHE_FLAG_UPDATE, cache)) {
|
||||
RETURN_TRUE;
|
||||
}
|
||||
RETURN_NULL();
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
#include "common.h"
|
||||
|
||||
PS_OPEN_FUNC(uwsgi) {
|
||||
PS_SET_MOD_DATA((char *)save_path);
|
||||
return SUCCESS;
|
||||
}
|
||||
|
||||
PS_CLOSE_FUNC(uwsgi) {
|
||||
return SUCCESS;
|
||||
}
|
||||
|
||||
PS_READ_FUNC(uwsgi) {
|
||||
char *cache = PS_GET_MOD_DATA();
|
||||
uint64_t valsize = 0;
|
||||
char *value = uwsgi_cache_magic_get((char *)key, strlen(key), &valsize, NULL, cache);
|
||||
if (!value) return FAILURE;
|
||||
char *new_val = emalloc(valsize);
|
||||
memcpy(new_val, value, valsize);
|
||||
free(value);
|
||||
*val = new_val;
|
||||
*vallen = valsize;
|
||||
return SUCCESS;
|
||||
|
||||
}
|
||||
|
||||
PS_WRITE_FUNC(uwsgi) {
|
||||
char *cache = PS_GET_MOD_DATA();
|
||||
if (vallen == 0) return SUCCESS;
|
||||
if (!uwsgi_cache_magic_set((char *)key, strlen(key), (char *)val, vallen, 0, UWSGI_CACHE_FLAG_UPDATE, cache)) {
|
||||
return SUCCESS;
|
||||
}
|
||||
return FAILURE;
|
||||
}
|
||||
|
||||
PS_DESTROY_FUNC(uwsgi) {
|
||||
char *cache = PS_GET_MOD_DATA();
|
||||
if (!uwsgi_cache_magic_del((char *)key, strlen(key), cache)) {
|
||||
return SUCCESS;
|
||||
}
|
||||
return FAILURE;
|
||||
}
|
||||
|
||||
PS_GC_FUNC(uwsgi) {
|
||||
return SUCCESS;
|
||||
}
|
||||
|
||||
ps_module ps_mod_uwsgi = {
|
||||
PS_MOD(uwsgi)
|
||||
};
|
||||
|
||||
@@ -25,4 +25,4 @@ phplibdir = os.environ.get('UWSGICONFIG_PHPLIBDIR')
|
||||
if phplibdir:
|
||||
LIBS.append('-Wl,-rpath=%s' % phplibdir)
|
||||
|
||||
GCC_LIST = ['php_plugin']
|
||||
GCC_LIST = ['php_plugin', 'session']
|
||||
|
||||
@@ -23,6 +23,9 @@ static PyObject *py_uwsgi_add_var(PyObject * self, PyObject * args) {
|
||||
return Py_True;
|
||||
}
|
||||
|
||||
static PyObject *py_uwsgi_micros(PyObject * self, PyObject * args) {
|
||||
return PyLong_FromUnsignedLongLong(uwsgi_micros());
|
||||
}
|
||||
|
||||
static PyObject *py_uwsgi_signal_wait(PyObject * self, PyObject * args) {
|
||||
|
||||
@@ -2439,6 +2442,8 @@ static PyMethodDef uwsgi_advanced_methods[] = {
|
||||
|
||||
{"add_var", py_uwsgi_add_var, METH_VARARGS, ""},
|
||||
|
||||
{"micros", py_uwsgi_micros, METH_VARARGS, ""},
|
||||
|
||||
{NULL, NULL},
|
||||
};
|
||||
|
||||
|
||||
@@ -5,6 +5,8 @@
|
||||
extern struct uwsgi_server uwsgi;
|
||||
|
||||
static int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_route *ur) {
|
||||
|
||||
struct uwsgi_buffer *ub = NULL;
|
||||
|
||||
// mark a route request
|
||||
wsgi_req->via = UWSGI_VIA_ROUTE;
|
||||
@@ -24,9 +26,14 @@ static int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_r
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// convert the wsgi_request to an http proxy request
|
||||
struct uwsgi_buffer *ub = uwsgi_to_http(wsgi_req, ur->data2, ur->data2_len, ub_url ? ub_url->buf : NULL, ub_url ? ub_url->pos : 0);
|
||||
if (ur->custom & 0x02) {
|
||||
ub = uwsgi_buffer_new(uwsgi.page_size);
|
||||
}
|
||||
else {
|
||||
ub = uwsgi_to_http(wsgi_req, ur->data2, ur->data2_len, ub_url ? ub_url->buf : NULL, ub_url ? ub_url->pos : 0);
|
||||
}
|
||||
|
||||
if (!ub) {
|
||||
if (ub_url) uwsgi_buffer_destroy(ub_url);
|
||||
uwsgi_log("unable to generate http request for %s\n", ub_addr->buf);
|
||||
@@ -46,11 +53,12 @@ static int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_r
|
||||
uwsgi_buffer_destroy(ub_addr);
|
||||
return UWSGI_ROUTE_NEXT;
|
||||
}
|
||||
wsgi_req->post_pos += wsgi_req->proto_parser_remains;
|
||||
wsgi_req->proto_parser_remains = 0;
|
||||
}
|
||||
|
||||
// ok now if have offload threads, directly use them
|
||||
if (!wsgi_req->post_file && !ur->custom && wsgi_req->socket->can_offload) {
|
||||
if (!wsgi_req->post_file && !(ur->custom & 0x01) && wsgi_req->socket->can_offload) {
|
||||
// append buffered body
|
||||
if (uwsgi.post_buffering > 0 && wsgi_req->post_cl > 0) {
|
||||
if (uwsgi_buffer_append(ub, wsgi_req->post_buffering_buf, wsgi_req->post_cl)) {
|
||||
@@ -60,6 +68,14 @@ static int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_r
|
||||
return UWSGI_ROUTE_NEXT;
|
||||
}
|
||||
}
|
||||
|
||||
// if we have a CONNECT request, let's confirm it to the client
|
||||
if (ur->custom & 0x02) {
|
||||
if (uwsgi_response_prepare_headers(wsgi_req, "200 Connection established", 26)) goto end;
|
||||
// no need to check for return value
|
||||
uwsgi_response_write_headers_do(wsgi_req);
|
||||
}
|
||||
|
||||
if (!uwsgi_offload_request_net_do(wsgi_req, ub_addr->buf, ub)) {
|
||||
wsgi_req->via = UWSGI_VIA_OFFLOAD;
|
||||
wsgi_req->status = 202;
|
||||
@@ -72,6 +88,7 @@ static int uwsgi_routing_func_http(struct wsgi_request *wsgi_req, struct uwsgi_r
|
||||
uwsgi_log("error routing request to http server %s\n", ub_addr->buf);
|
||||
}
|
||||
|
||||
end:
|
||||
uwsgi_buffer_destroy(ub);
|
||||
uwsgi_buffer_destroy(ub_addr);
|
||||
|
||||
@@ -102,15 +119,27 @@ static int uwsgi_router_http(struct uwsgi_route *ur, char *args) {
|
||||
}
|
||||
|
||||
static int uwsgi_router_proxyhttp(struct uwsgi_route *ur, char *args) {
|
||||
ur->custom = 1;
|
||||
ur->custom = 0x1;
|
||||
return uwsgi_router_http(ur, args);
|
||||
}
|
||||
|
||||
static int uwsgi_router_proxyhttp_connect(struct uwsgi_route *ur, char *args) {
|
||||
ur->custom = 0x1|0x02;
|
||||
return uwsgi_router_http(ur, args);
|
||||
}
|
||||
|
||||
static int uwsgi_router_http_connect(struct uwsgi_route *ur, char *args) {
|
||||
ur->custom = 0x02;
|
||||
return uwsgi_router_http(ur, args);
|
||||
}
|
||||
|
||||
|
||||
static void router_http_register(void) {
|
||||
|
||||
uwsgi_register_router("http", uwsgi_router_http);
|
||||
uwsgi_register_router("proxyhttp", uwsgi_router_proxyhttp);
|
||||
uwsgi_register_router("httpconnect", uwsgi_router_http_connect);
|
||||
uwsgi_register_router("proxyhttpconnect", uwsgi_router_proxyhttp_connect);
|
||||
}
|
||||
|
||||
struct uwsgi_plugin router_http_plugin = {
|
||||
|
||||
@@ -18,6 +18,8 @@ struct uwsgi_router_file_conf {
|
||||
char *mime;
|
||||
|
||||
char *no_cl;
|
||||
|
||||
char *no_headers;
|
||||
};
|
||||
|
||||
int uwsgi_routing_func_static(struct wsgi_request *wsgi_req, struct uwsgi_route *ur) {
|
||||
@@ -38,6 +40,7 @@ int uwsgi_routing_func_file(struct wsgi_request *wsgi_req, struct uwsgi_route *u
|
||||
char buf[32768];
|
||||
struct stat st;
|
||||
int ret = UWSGI_ROUTE_BREAK;
|
||||
size_t remains = 0;
|
||||
|
||||
struct uwsgi_router_file_conf *urfc = (struct uwsgi_router_file_conf *) ur->data2;
|
||||
|
||||
@@ -61,6 +64,8 @@ int uwsgi_routing_func_file(struct wsgi_request *wsgi_req, struct uwsgi_route *u
|
||||
struct uwsgi_buffer *ub_s = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urfc->status, urfc->status_len);
|
||||
if (!ub_s) goto end2;
|
||||
|
||||
if (urfc->no_headers) goto send;
|
||||
|
||||
if (uwsgi_response_prepare_headers(wsgi_req, ub_s->buf, ub_s->pos)) {
|
||||
uwsgi_buffer_destroy(ub_s);
|
||||
goto end2;
|
||||
@@ -82,8 +87,10 @@ int uwsgi_routing_func_file(struct wsgi_request *wsgi_req, struct uwsgi_route *u
|
||||
else {
|
||||
if (uwsgi_response_add_content_type(wsgi_req, urfc->content_type, urfc->content_type_len)) goto end2;
|
||||
}
|
||||
|
||||
send:
|
||||
|
||||
size_t remains = st.st_size;
|
||||
remains = st.st_size;
|
||||
while(remains) {
|
||||
ssize_t rlen = read(fd, buf, UMIN(32768, remains));
|
||||
if (rlen <= 0) goto end2;
|
||||
@@ -125,6 +132,8 @@ int uwsgi_routing_func_sendfile(struct wsgi_request *wsgi_req, struct uwsgi_rout
|
||||
struct uwsgi_buffer *ub_s = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urfc->status, urfc->status_len);
|
||||
if (!ub_s) goto end2;
|
||||
|
||||
if (urfc->no_headers) goto send;
|
||||
|
||||
if (uwsgi_response_prepare_headers(wsgi_req, ub_s->buf, ub_s->pos)) {
|
||||
uwsgi_buffer_destroy(ub_s);
|
||||
goto end2;
|
||||
@@ -145,6 +154,8 @@ int uwsgi_routing_func_sendfile(struct wsgi_request *wsgi_req, struct uwsgi_rout
|
||||
if (uwsgi_response_add_content_type(wsgi_req, urfc->content_type, urfc->content_type_len)) goto end2;
|
||||
}
|
||||
|
||||
send:
|
||||
|
||||
if (!wsgi_req->headers_sent) {
|
||||
if (uwsgi_response_write_headers_do(wsgi_req)) goto end2;
|
||||
}
|
||||
@@ -188,6 +199,8 @@ int uwsgi_routing_func_fastfile(struct wsgi_request *wsgi_req, struct uwsgi_rout
|
||||
struct uwsgi_buffer *ub_s = uwsgi_routing_translate(wsgi_req, ur, *subject, *subject_len, urfc->status, urfc->status_len);
|
||||
if (!ub_s) goto end2;
|
||||
|
||||
if (urfc->no_headers) goto send;
|
||||
|
||||
if (uwsgi_response_prepare_headers(wsgi_req, ub_s->buf, ub_s->pos)) {
|
||||
uwsgi_buffer_destroy(ub_s);
|
||||
goto end2;
|
||||
@@ -208,6 +221,8 @@ int uwsgi_routing_func_fastfile(struct wsgi_request *wsgi_req, struct uwsgi_rout
|
||||
if (uwsgi_response_add_content_type(wsgi_req, urfc->content_type, urfc->content_type_len)) goto end2;
|
||||
}
|
||||
|
||||
send:
|
||||
|
||||
if (!wsgi_req->headers_sent) {
|
||||
if (uwsgi_response_write_headers_do(wsgi_req)) goto end2;
|
||||
}
|
||||
@@ -256,6 +271,7 @@ static int uwsgi_router_file(struct uwsgi_route *ur, char *args) {
|
||||
"no_cl", &urfc->no_cl,
|
||||
"no_content_length", &urfc->no_cl,
|
||||
"mime", &urfc->mime,
|
||||
"no_headers", &urfc->no_headers,
|
||||
NULL)) {
|
||||
uwsgi_log("invalid file route syntax: %s\n", args);
|
||||
return -1;
|
||||
|
||||
@@ -54,6 +54,7 @@ static int uwsgi_routing_func_uwsgi_remote(struct wsgi_request *wsgi_req, struct
|
||||
if (uwsgi_buffer_append(ub, wsgi_req->proto_parser_remains_buf, wsgi_req->proto_parser_remains)) {
|
||||
goto end;
|
||||
}
|
||||
wsgi_req->post_pos += wsgi_req->proto_parser_remains;
|
||||
wsgi_req->proto_parser_remains = 0;
|
||||
}
|
||||
|
||||
|
||||
@@ -326,6 +326,11 @@ static void tornado_loop() {
|
||||
uwsgi_log("!!! Running tornado with a socket-timeout lower than 30 seconds is not recommended, tune it with --socket-timeout !!!\n");
|
||||
}
|
||||
|
||||
if (!uwsgi.async_waiting_fd_table)
|
||||
uwsgi.async_waiting_fd_table = uwsgi_calloc(sizeof(struct wsgi_request *) * uwsgi.max_fd);
|
||||
if (!uwsgi.async_proto_fd_table)
|
||||
uwsgi.async_proto_fd_table = uwsgi_calloc(sizeof(struct wsgi_request *) * uwsgi.max_fd);
|
||||
|
||||
// get the GIL
|
||||
UWSGI_GET_GIL
|
||||
|
||||
|
||||
+110
-45
@@ -4,52 +4,55 @@
|
||||
|
||||
extern struct uwsgi_server uwsgi;
|
||||
|
||||
static uint16_t http_add_uwsgi_header(struct wsgi_request *wsgi_req, char *hh, int hhlen) {
|
||||
static char * http_header_to_cgi(char *hh, size_t hhlen, size_t *keylen, size_t *vallen, int *has_prefix) {
|
||||
size_t i;
|
||||
char *val = hh;
|
||||
int status = 0;
|
||||
for (i = 0; i < hhlen; i++) {
|
||||
if (!status) {
|
||||
hh[i] = toupper((int) hh[i]);
|
||||
if (hh[i] == '-')
|
||||
hh[i] = '_';
|
||||
if (hh[i] == ':') {
|
||||
status = 1;
|
||||
*keylen = i;
|
||||
}
|
||||
}
|
||||
else if (status == 1 && hh[i] != ' ') {
|
||||
status = 2;
|
||||
val += i;
|
||||
*vallen+=1;
|
||||
}
|
||||
else if (status == 2) {
|
||||
*vallen+=1;
|
||||
}
|
||||
}
|
||||
|
||||
if (!(*keylen))
|
||||
return NULL;
|
||||
|
||||
if (uwsgi_strncmp("CONTENT_LENGTH", 14, hh, *keylen) && uwsgi_strncmp("CONTENT_TYPE", 12, hh, *keylen)) {
|
||||
*has_prefix = 0x02;
|
||||
}
|
||||
|
||||
return val;
|
||||
}
|
||||
|
||||
static uint16_t http_add_uwsgi_header(struct wsgi_request *wsgi_req, char *hh, size_t hhlen, char *hv, size_t hvlen, int has_prefix) {
|
||||
|
||||
char *buffer = wsgi_req->buffer + wsgi_req->uh->pktsize;
|
||||
char *watermark = wsgi_req->buffer + uwsgi.buffer_size;
|
||||
|
||||
int i;
|
||||
int status = 0;
|
||||
char *val = hh;
|
||||
uint16_t keylen = 0, vallen = 0;
|
||||
int prefix = 0;
|
||||
char *ptr = buffer;
|
||||
size_t keylen = hhlen;
|
||||
|
||||
for (i = 0; i < hhlen; i++) {
|
||||
if (!status) {
|
||||
hh[i] = toupper((int) hh[i]);
|
||||
if (hh[i] == '-')
|
||||
hh[i] = '_';
|
||||
if (hh[i] == ':') {
|
||||
status = 1;
|
||||
keylen = i;
|
||||
}
|
||||
}
|
||||
else if (status == 1 && hh[i] != ' ') {
|
||||
status = 2;
|
||||
val += i;
|
||||
vallen++;
|
||||
}
|
||||
else if (status == 2) {
|
||||
vallen++;
|
||||
}
|
||||
}
|
||||
if (has_prefix) keylen += 5;
|
||||
|
||||
if (!keylen)
|
||||
return 0;
|
||||
|
||||
if (uwsgi_strncmp("CONTENT_LENGTH", 14, hh, keylen) && uwsgi_strncmp("CONTENT_TYPE", 12, hh, keylen)) {
|
||||
keylen += 5;
|
||||
prefix = 1;
|
||||
}
|
||||
|
||||
if (buffer + keylen + vallen + 2 + 2 >= watermark) {
|
||||
if (prefix) {
|
||||
uwsgi_log("[WARNING] unable to add HTTP_%.*s=%.*s to uwsgi packet, consider increasing buffer size\n", keylen, hh, vallen, val);
|
||||
if (buffer + keylen + hvlen + 2 + 2 >= watermark) {
|
||||
if (has_prefix) {
|
||||
uwsgi_log("[WARNING] unable to add HTTP_%.*s=%.*s to uwsgi packet, consider increasing buffer size\n", keylen, hh, hvlen, hv);
|
||||
}
|
||||
else {
|
||||
uwsgi_log("[WARNING] unable to add %.*s=%.*s to uwsgi packet, consider increasing buffer size\n", keylen, hh, vallen, val);
|
||||
uwsgi_log("[WARNING] unable to add %.*s=%.*s to uwsgi packet, consider increasing buffer size\n", keylen, hh, hvlen, hv);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
@@ -58,7 +61,7 @@ static uint16_t http_add_uwsgi_header(struct wsgi_request *wsgi_req, char *hh, i
|
||||
*ptr++ = (uint8_t) (keylen & 0xff);
|
||||
*ptr++ = (uint8_t) ((keylen >> 8) & 0xff);
|
||||
|
||||
if (prefix) {
|
||||
if (has_prefix) {
|
||||
memcpy(ptr, "HTTP_", 5);
|
||||
ptr += 5;
|
||||
memcpy(ptr, hh, keylen - 5);
|
||||
@@ -69,11 +72,11 @@ static uint16_t http_add_uwsgi_header(struct wsgi_request *wsgi_req, char *hh, i
|
||||
ptr += keylen;
|
||||
}
|
||||
|
||||
*ptr++ = (uint8_t) (vallen & 0xff);
|
||||
*ptr++ = (uint8_t) ((vallen >> 8) & 0xff);
|
||||
memcpy(ptr, val, vallen);
|
||||
*ptr++ = (uint8_t) (hvlen & 0xff);
|
||||
*ptr++ = (uint8_t) ((hvlen >> 8) & 0xff);
|
||||
memcpy(ptr, hv, hvlen);
|
||||
|
||||
return 2 + keylen + 2 + vallen;
|
||||
return 2 + keylen + 2 + hvlen;
|
||||
}
|
||||
|
||||
static char *proxy1_parse(char *ptr, char *watermark, char **src, uint16_t *src_len, char **dst, uint16_t *dst_len, char **src_port, uint16_t *src_port_len, char **dst_port, uint16_t *dst_port_len) {
|
||||
@@ -301,9 +304,15 @@ static int http_parse(struct wsgi_request *wsgi_req, char *watermark) {
|
||||
}
|
||||
}
|
||||
|
||||
if (wsgi_req->https_len > 0) {
|
||||
wsgi_req->uh->pktsize += proto_base_add_uwsgi_var(wsgi_req, "HTTPS", 5, wsgi_req->https, wsgi_req->https_len);
|
||||
}
|
||||
|
||||
//HEADERS
|
||||
base = ptr;
|
||||
|
||||
struct uwsgi_string_list *headers = NULL, *usl = NULL;
|
||||
|
||||
while (ptr < watermark) {
|
||||
if (*ptr == '\r') {
|
||||
if (ptr + 1 >= watermark)
|
||||
@@ -317,15 +326,71 @@ static int http_parse(struct wsgi_request *wsgi_req, char *watermark) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
wsgi_req->uh->pktsize += http_add_uwsgi_header(wsgi_req, base, ptr - base);
|
||||
size_t key_len = 0, value_len = 0;
|
||||
int has_prefix = 0;
|
||||
// last line, do not waste time
|
||||
if (ptr - base == 0) break;
|
||||
char *value = http_header_to_cgi(base, ptr - base, &key_len, &value_len, &has_prefix);
|
||||
if (!value) {
|
||||
uwsgi_log_verbose("invalid HTTP request\n");
|
||||
goto clear;
|
||||
}
|
||||
usl = uwsgi_string_list_has_item(headers, base, key_len);
|
||||
// there is already a HTTP header with the same name, let's merge them
|
||||
if (usl) {
|
||||
char *old_value = usl->custom_ptr;
|
||||
usl->custom_ptr = uwsgi_concat3n(old_value, (size_t) usl->custom, ", ", 2, value, value_len);
|
||||
usl->custom += 2 + value_len;
|
||||
if (usl->custom2 & 0x01) free(old_value);
|
||||
usl->custom2 |= 0x01;
|
||||
}
|
||||
else {
|
||||
// add an entry
|
||||
usl = uwsgi_string_new_list(&headers, NULL);
|
||||
usl->value = base;
|
||||
usl->len = key_len;
|
||||
usl->custom_ptr = value;
|
||||
usl->custom = value_len;
|
||||
usl->custom2 = has_prefix;
|
||||
}
|
||||
ptr++;
|
||||
base = ptr + 1;
|
||||
}
|
||||
ptr++;
|
||||
}
|
||||
|
||||
return 0;
|
||||
usl = headers;
|
||||
int broken = 0;
|
||||
while(usl) {
|
||||
if (!broken) {
|
||||
uint16_t old_pktsize = wsgi_req->uh->pktsize;
|
||||
wsgi_req->uh->pktsize += http_add_uwsgi_header(wsgi_req, usl->value, usl->len, usl->custom_ptr, (size_t) usl->custom, usl->custom2 & 0x02);
|
||||
// if the packet remains unchanged, the buffer is full, mark the request as broken
|
||||
if (old_pktsize == wsgi_req->uh->pktsize) {
|
||||
broken = 1;
|
||||
}
|
||||
}
|
||||
if (usl->custom2 & 0x01) {
|
||||
free(usl->custom_ptr);
|
||||
}
|
||||
struct uwsgi_string_list *tmp_usl = usl;
|
||||
usl = usl->next;
|
||||
free(tmp_usl);
|
||||
}
|
||||
|
||||
return broken;
|
||||
|
||||
clear:
|
||||
usl = headers;
|
||||
while(usl) {
|
||||
if (usl->custom2 & 0x01) {
|
||||
free(usl->custom_ptr);
|
||||
}
|
||||
struct uwsgi_string_list *tmp_usl = usl;
|
||||
usl = usl->next;
|
||||
free(tmp_usl);
|
||||
}
|
||||
return -1;
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
#!./uwsgi --http-socket :9090 --gevent 100 --module tests.websocket_chat --gevent-monkey-patch
|
||||
#!./uwsgi --http-socket :9090 --gevent 100 --module tests.websockets_chat --gevent-monkey-patch
|
||||
import uwsgi
|
||||
import time
|
||||
import gevent.select
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
#!./uwsgi --http-socket :9090 --asyncio 100 --module tests.websockets_chat_asyncio --greenlet
|
||||
import uwsgi
|
||||
import asyncio
|
||||
import asyncio_redis
|
||||
import time
|
||||
import greenlet
|
||||
|
||||
class GreenFuture(asyncio.Future):
|
||||
def __init__(self):
|
||||
super().__init__()
|
||||
self.greenlet = greenlet.getcurrent()
|
||||
self.add_done_callback(lambda f: f.greenlet.switch())
|
||||
|
||||
def result(self):
|
||||
while True:
|
||||
if self.done():
|
||||
return super().result()
|
||||
self.greenlet.parent.switch()
|
||||
|
||||
|
||||
@asyncio.coroutine
|
||||
def redis_open(f):
|
||||
connection = yield from asyncio_redis.Connection.create(host='localhost', port=6379)
|
||||
f.set_result(connection)
|
||||
f.greenlet.switch()
|
||||
|
||||
@asyncio.coroutine
|
||||
def redis_subscribe(f):
|
||||
connection = yield from asyncio_redis.Connection.create(host='localhost', port=6379)
|
||||
subscriber = yield from connection.start_subscribe()
|
||||
yield from subscriber.subscribe([ 'foobar' ])
|
||||
f.set_result(subscriber)
|
||||
f.greenlet.switch()
|
||||
|
||||
def ws_recv_msg(g):
|
||||
g.has_ws_msg = True
|
||||
g.switch()
|
||||
|
||||
@asyncio.coroutine
|
||||
def redis_wait(subscriber, f):
|
||||
reply = yield from subscriber.next_published()
|
||||
f.set_result(reply.value)
|
||||
f.greenlet.switch()
|
||||
|
||||
@asyncio.coroutine
|
||||
def redis_publish(connection, msg):
|
||||
yield from connection.publish('foobar', msg.decode('utf-8'))
|
||||
|
||||
def application(env, sr):
|
||||
|
||||
ws_scheme = 'ws'
|
||||
if 'HTTPS' in env or env['wsgi.url_scheme'] == 'https':
|
||||
ws_scheme = 'wss'
|
||||
|
||||
if env['PATH_INFO'] == '/':
|
||||
sr('200 OK', [('Content-Type','text/html')])
|
||||
return ("""
|
||||
<html>
|
||||
<head>
|
||||
<script language="Javascript">
|
||||
var s = new WebSocket("%s://%s/foobar/");
|
||||
s.onopen = function() {
|
||||
alert("connected !!!");
|
||||
s.send("ciao");
|
||||
};
|
||||
s.onmessage = function(e) {
|
||||
var bb = document.getElementById('blackboard')
|
||||
var html = bb.innerHTML;
|
||||
bb.innerHTML = html + '<br/>' + e.data;
|
||||
};
|
||||
|
||||
s.onerror = function(e) {
|
||||
alert(e);
|
||||
}
|
||||
|
||||
s.onclose = function(e) {
|
||||
alert("connection closed");
|
||||
}
|
||||
|
||||
function invia() {
|
||||
var value = document.getElementById('testo').value;
|
||||
s.send(value);
|
||||
}
|
||||
</script>
|
||||
</head>
|
||||
<body>
|
||||
<h1>WebSocket</h1>
|
||||
<input type="text" id="testo"/>
|
||||
<input type="button" value="invia" onClick="invia();"/>
|
||||
<div id="blackboard" style="width:640px;height:480px;background-color:black;color:white;border: solid 2px red;overflow:auto">
|
||||
</div>
|
||||
</body>
|
||||
</html>
|
||||
""" % (ws_scheme, env['HTTP_HOST'])).encode()
|
||||
elif env['PATH_INFO'] == '/favicon.ico':
|
||||
return b""
|
||||
elif env['PATH_INFO'] == '/foobar/':
|
||||
uwsgi.websocket_handshake()
|
||||
print("websockets...")
|
||||
# a future for waiting for redis connection
|
||||
f = GreenFuture()
|
||||
asyncio.Task(redis_subscribe(f))
|
||||
# the result() method will switch greenlets if needed
|
||||
subscriber = f.result()
|
||||
|
||||
# open another redis connection for publishing messages
|
||||
f0 = GreenFuture()
|
||||
t = asyncio.Task(redis_open(f0))
|
||||
connection = f0.result()
|
||||
|
||||
myself = greenlet.getcurrent()
|
||||
myself.has_ws_msg = False
|
||||
# start monitoring websocket events
|
||||
asyncio.get_event_loop().add_reader(uwsgi.connection_fd(), ws_recv_msg, myself)
|
||||
|
||||
# add a 4 seconds timer to manage ping/pong
|
||||
asyncio.get_event_loop().call_later(4, ws_recv_msg, myself)
|
||||
|
||||
# add a coroutine for redis messages
|
||||
f = GreenFuture()
|
||||
asyncio.Task(redis_wait(subscriber, f))
|
||||
|
||||
# switch again
|
||||
f.greenlet.parent.switch()
|
||||
|
||||
while True:
|
||||
# any redis message in the queue ?
|
||||
if f.done():
|
||||
msg = f.result()
|
||||
uwsgi.websocket_send("[%s] %s" % (time.time(), msg))
|
||||
# restart coroutine
|
||||
f = GreenFuture()
|
||||
asyncio.Task(redis_wait(subscriber, f))
|
||||
if myself.has_ws_msg:
|
||||
myself.has_ws_msg = False
|
||||
msg = uwsgi.websocket_recv_nb()
|
||||
if msg:
|
||||
asyncio.Task(redis_publish(connection, msg))
|
||||
# switch again
|
||||
f.greenlet.parent.switch()
|
||||
+1
-1
@@ -2,7 +2,7 @@ Gem::Specification.new do |s|
|
||||
s.name = 'uwsgi'
|
||||
s.license = 'GPL-2'
|
||||
s.version = `python -c "import uwsgiconfig as uc; print uc.uwsgi_version"`.sub(/-dev-.*/,'')
|
||||
s.date = '2014-03-17'
|
||||
s.date = '2014-04-22'
|
||||
s.summary = "uWSGI"
|
||||
s.description = "The uWSGI server for Ruby/Rack"
|
||||
s.authors = ["Unbit"]
|
||||
|
||||
@@ -765,6 +765,10 @@ struct uwsgi_cache_item {
|
||||
uint64_t prev;
|
||||
// next same-hash item
|
||||
uint64_t next;
|
||||
// previous lru item
|
||||
uint64_t lru_prev;
|
||||
// next lru item
|
||||
uint64_t lru_next;
|
||||
// key characters follows...
|
||||
char key[];
|
||||
} __attribute__ ((__packed__));
|
||||
@@ -819,6 +823,13 @@ struct uwsgi_cache {
|
||||
struct uwsgi_lock_item *lock;
|
||||
|
||||
struct uwsgi_cache *next;
|
||||
|
||||
int ignore_full;
|
||||
|
||||
uint64_t next_scan;
|
||||
int purge_lru;
|
||||
uint64_t lru_head;
|
||||
uint64_t lru_tail;
|
||||
};
|
||||
|
||||
struct uwsgi_option {
|
||||
@@ -4044,6 +4055,7 @@ void emperor_stop(struct uwsgi_instance *);
|
||||
void emperor_curse(struct uwsgi_instance *);
|
||||
void emperor_respawn(struct uwsgi_instance *, time_t);
|
||||
void emperor_add(struct uwsgi_emperor_scanner *, char *, time_t, char *, uint32_t, uid_t, gid_t, char *);
|
||||
void emperor_back_to_ondemand(struct uwsgi_instance *);
|
||||
|
||||
void uwsgi_exec_command_with_args(char *);
|
||||
|
||||
|
||||
+9
-9
@@ -1,6 +1,6 @@
|
||||
# uWSGI build system
|
||||
|
||||
uwsgi_version = '2.0.3'
|
||||
uwsgi_version = '2.0.4'
|
||||
|
||||
import os
|
||||
import re
|
||||
@@ -331,6 +331,13 @@ def build_uwsgi(uc, print_only=False, gcll=None):
|
||||
uwsgi_config_py = uwsgi_config_py_content.encode('hex')
|
||||
open('core/config_py.c', 'w').write('char *uwsgi_config_py = "%s";\n' % uwsgi_config_py);
|
||||
gcc_list.append('core/config_py')
|
||||
|
||||
additional_sources = os.environ.get('UWSGI_ADDITIONAL_SOURCES')
|
||||
if not additional_sources:
|
||||
additional_sources = uc.get('additional_sources')
|
||||
if additional_sources:
|
||||
for item in additional_sources.split(','):
|
||||
gcc_list.append(item)
|
||||
|
||||
cflags.append('-DUWSGI_CFLAGS=\\"%s\\"' % uwsgi_cflags)
|
||||
cflags.append('-DUWSGI_BUILD_DATE="\\"%s\\""' % time.strftime("%d %B %Y %H:%M:%S"))
|
||||
@@ -519,13 +526,6 @@ def build_uwsgi(uc, print_only=False, gcll=None):
|
||||
for ef in binary_list:
|
||||
gcc_list.append("%s.o" % ef)
|
||||
|
||||
additional_sources = os.environ.get('UWSGI_ADDITIONAL_SOURCES')
|
||||
if not additional_sources:
|
||||
additional_sources = uc.get('additional_sources')
|
||||
if additional_sources:
|
||||
for item in additional_sources.split(','):
|
||||
gcc_list.append(item)
|
||||
|
||||
if compile_queue:
|
||||
for t in thread_compilers:
|
||||
compile_queue.put((None, None))
|
||||
@@ -1058,7 +1058,7 @@ class uConf(object):
|
||||
|
||||
self.embed_config = None
|
||||
|
||||
if uwsgi_os == 'Linux':
|
||||
if uwsgi_os in ('Linux','FreeBSD'):
|
||||
self.embed_config = os.environ.get('UWSGI_EMBED_CONFIG')
|
||||
if not self.embed_config:
|
||||
self.embed_config = self.get('embed_config')
|
||||
|
||||
Reference in New Issue
Block a user