mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-09-01 02:35:03 +00:00
preliminary Websockets and SPDY support, highly optimized corerouters
This commit is contained in:
+83
-1
@@ -48,6 +48,20 @@ int uwsgi_buffer_ensure(struct uwsgi_buffer *ub, size_t len) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_buffer_decapitate(struct uwsgi_buffer *ub, size_t len) {
|
||||
if (len > ub->pos) return -1;
|
||||
ub->buf = memmove(ub->buf, ub->buf + len, ub->pos-len);
|
||||
ub->pos = ub->pos-len;
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_buffer_byte(struct uwsgi_buffer *ub, char byte) {
|
||||
return uwsgi_buffer_append(ub, &byte, 1);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_u8(struct uwsgi_buffer *ub, uint8_t u8) {
|
||||
return uwsgi_buffer_append(ub, (char *) &u8, 1);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_append(struct uwsgi_buffer *ub, char *buf, size_t len) {
|
||||
|
||||
@@ -84,6 +98,46 @@ int uwsgi_buffer_u16le(struct uwsgi_buffer *ub, uint16_t num) {
|
||||
return uwsgi_buffer_append(ub, (char *) buf, 2);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_u16be(struct uwsgi_buffer *ub, uint16_t num) {
|
||||
uint8_t buf[2];
|
||||
buf[1] = (uint8_t) (num & 0xff);
|
||||
buf[0] = (uint8_t) ((num >> 8) & 0xff);
|
||||
return uwsgi_buffer_append(ub, (char *) buf, 2);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_u32be(struct uwsgi_buffer *ub, uint32_t num) {
|
||||
uint8_t buf[4];
|
||||
buf[3] = (uint8_t) (num & 0xff);
|
||||
buf[2] = (uint8_t) ((num >> 8) & 0xff);
|
||||
buf[1] = (uint8_t) ((num >> 16) & 0xff);
|
||||
buf[0] = (uint8_t) ((num >> 24) & 0xff);
|
||||
return uwsgi_buffer_append(ub, (char *) buf, 4);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_u64be(struct uwsgi_buffer *ub, uint64_t num) {
|
||||
uint8_t buf[8];
|
||||
buf[7] = (uint8_t) (num & 0xff);
|
||||
buf[6] = (uint8_t) ((num >> 8) & 0xff);
|
||||
buf[5] = (uint8_t) ((num >> 16) & 0xff);
|
||||
buf[4] = (uint8_t) ((num >> 24) & 0xff);
|
||||
buf[3] = (uint8_t) ((num >> 32) & 0xff);
|
||||
buf[2] = (uint8_t) ((num >> 40) & 0xff);
|
||||
buf[1] = (uint8_t) ((num >> 48) & 0xff);
|
||||
buf[0] = (uint8_t) ((num >> 56) & 0xff);
|
||||
return uwsgi_buffer_append(ub, (char *) buf, 8);
|
||||
}
|
||||
|
||||
|
||||
int uwsgi_buffer_append_ipv4(struct uwsgi_buffer *ub, void *addr) {
|
||||
char ip[INET_ADDRSTRLEN];
|
||||
if (!inet_ntop(AF_INET, addr, ip, INET_ADDRSTRLEN)) {
|
||||
uwsgi_error("uwsgi_buffer_append_ipv4() -> inet_ntop()");
|
||||
return -1;
|
||||
}
|
||||
return uwsgi_buffer_append(ub, ip, strlen(ip));
|
||||
}
|
||||
|
||||
|
||||
int uwsgi_buffer_num64(struct uwsgi_buffer *ub, int64_t num) {
|
||||
char buf[sizeof(UMAX64_STR)+1];
|
||||
int ret = snprintf(buf, sizeof(UMAX64_STR)+1, "%lld", (long long) num);
|
||||
@@ -93,13 +147,20 @@ int uwsgi_buffer_num64(struct uwsgi_buffer *ub, int64_t num) {
|
||||
return uwsgi_buffer_append(ub, buf, ret);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_append_keyval(struct uwsgi_buffer *ub, char *key, uint16_t keylen, char *val, uint64_t vallen) {
|
||||
int uwsgi_buffer_append_keyval(struct uwsgi_buffer *ub, char *key, uint16_t keylen, char *val, uint16_t vallen) {
|
||||
if (uwsgi_buffer_u16le(ub, keylen)) return -1;
|
||||
if (uwsgi_buffer_append(ub, key, keylen)) return -1;
|
||||
if (uwsgi_buffer_u16le(ub, vallen)) return -1;
|
||||
return uwsgi_buffer_append(ub, val, vallen);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_append_keyval32(struct uwsgi_buffer *ub, char *key, uint32_t keylen, char *val, uint32_t vallen) {
|
||||
if (uwsgi_buffer_u32be(ub, keylen)) return -1;
|
||||
if (uwsgi_buffer_append(ub, key, keylen)) return -1;
|
||||
if (uwsgi_buffer_u32be(ub, vallen)) return -1;
|
||||
return uwsgi_buffer_append(ub, val, vallen);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_append_keynum(struct uwsgi_buffer *ub, char *key, uint16_t keylen, int64_t num) {
|
||||
char buf[sizeof(UMAX64_STR)+1];
|
||||
int ret = snprintf(buf, (sizeof(UMAX64_STR)+1), "%lld", (long long) num);
|
||||
@@ -112,6 +173,27 @@ int uwsgi_buffer_append_keynum(struct uwsgi_buffer *ub, char *key, uint16_t keyl
|
||||
return uwsgi_buffer_append(ub, buf, ret);
|
||||
}
|
||||
|
||||
int uwsgi_buffer_append_keyipv4(struct uwsgi_buffer *ub, char *key, uint16_t keylen, void *addr) {
|
||||
if (uwsgi_buffer_u16le(ub, keylen)) return -1;
|
||||
if (uwsgi_buffer_append(ub, key, keylen)) return -1;
|
||||
if (uwsgi_buffer_u16le(ub, 15)) return -1;
|
||||
char *ptr = ub->buf + (ub->pos - 2);
|
||||
if (uwsgi_buffer_append_ipv4(ub, addr)) return -1;
|
||||
// fix the size
|
||||
*ptr = ((ub->buf+ub->pos) - (ptr+2));
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_buffer_append_base64(struct uwsgi_buffer *ub, char *s, size_t len) {
|
||||
size_t b64_len = 0;
|
||||
char *b64 = uwsgi_base64_encode(s, len, &b64_len);
|
||||
if (!b64) return -1;
|
||||
int ret = uwsgi_buffer_append(ub, b64, b64_len);
|
||||
free(b64);
|
||||
return ret;
|
||||
}
|
||||
|
||||
|
||||
void uwsgi_buffer_destroy(struct uwsgi_buffer *ub) {
|
||||
if (ub->buf)
|
||||
free(ub->buf);
|
||||
|
||||
+3
-1
@@ -90,6 +90,7 @@ void uwsgi_init_default() {
|
||||
uwsgi.multicast_loop = 1;
|
||||
#endif
|
||||
|
||||
uwsgi_websockets_init();
|
||||
}
|
||||
|
||||
void uwsgi_setup_reload() {
|
||||
@@ -265,7 +266,8 @@ void uwsgi_setup_workers() {
|
||||
}
|
||||
|
||||
total_memory *= (uwsgi.numproc + uwsgi.master_process);
|
||||
uwsgi_log("mapped %lu bytes (%lu KB) for %d cores\n", total_memory, total_memory / 1024, uwsgi.cores * uwsgi.numproc);
|
||||
if (uwsgi.numproc > 0)
|
||||
uwsgi_log("mapped %lu bytes (%lu KB) for %d cores\n", total_memory, total_memory / 1024, uwsgi.cores * uwsgi.numproc);
|
||||
|
||||
}
|
||||
|
||||
|
||||
+18
@@ -375,3 +375,21 @@ char *uwsgi_ssl_rand(size_t len) {
|
||||
}
|
||||
return (char *) buf;
|
||||
}
|
||||
|
||||
char *uwsgi_sha1(char *src, size_t len, char *dst) {
|
||||
SHA_CTX sha;
|
||||
SHA1_Init(&sha);
|
||||
SHA1_Update(&sha, src, len);
|
||||
SHA1_Final((unsigned char *)dst, &sha);
|
||||
return dst;
|
||||
}
|
||||
|
||||
char *uwsgi_sha1_2n(char *s1, size_t len1, char *s2, size_t len2, char *dst) {
|
||||
SHA_CTX sha;
|
||||
SHA1_Init(&sha);
|
||||
SHA1_Update(&sha, s1, len1);
|
||||
SHA1_Update(&sha, s2, len2);
|
||||
SHA1_Final((unsigned char *)dst, &sha);
|
||||
return dst;
|
||||
}
|
||||
|
||||
|
||||
+61
-10
@@ -734,6 +734,11 @@ void uwsgi_close_request(struct wsgi_request *wsgi_req) {
|
||||
free(ptr);
|
||||
}
|
||||
|
||||
// free websocket engine
|
||||
if (wsgi_req->websocket_buf) {
|
||||
uwsgi_buffer_destroy(wsgi_req->websocket_buf);
|
||||
}
|
||||
|
||||
|
||||
// reset request
|
||||
tmp_id = wsgi_req->async_id;
|
||||
@@ -1199,6 +1204,16 @@ int uwsgi_strncmp(char *src, int slen, char *dst, int dlen) {
|
||||
|
||||
}
|
||||
|
||||
// fast compare 2 sized strings (case insensitive)
|
||||
int uwsgi_strnicmp(char *src, int slen, char *dst, int dlen) {
|
||||
|
||||
if (slen != dlen)
|
||||
return 1;
|
||||
|
||||
return strncasecmp(src, dst, dlen);
|
||||
|
||||
}
|
||||
|
||||
// fast sized check of initial part of a string
|
||||
int uwsgi_starts_with(char *src, int slen, char *dst, int dlen) {
|
||||
|
||||
@@ -4081,33 +4096,69 @@ char *uwsgi_base64_decode(char *buf, size_t len, size_t *d_len) {
|
||||
|
||||
char *uwsgi_base64_encode(char *buf, size_t len, size_t *d_len) {
|
||||
*d_len = ((len * 4)/3) + 5;
|
||||
uint8_t *src = (uint8_t *) buf;
|
||||
char *dst = uwsgi_malloc(*d_len);
|
||||
char *ptr = dst;
|
||||
while(len >= 3) {
|
||||
*ptr++= b64_table64_2[ buf[0] >> 2];
|
||||
*ptr++= b64_table64_2[((buf[0] << 4) & 0x30) | (buf[1] >> 4)];
|
||||
*ptr++= b64_table64_2[((buf[1] << 2) & 0x3C) | (buf[2] >> 6)];
|
||||
*ptr++= b64_table64_2[buf[2] & 0x3F];
|
||||
buf += 3;
|
||||
*ptr++= b64_table64_2[ src[0] >> 2];
|
||||
*ptr++= b64_table64_2[((src[0] << 4) & 0x30) | (src[1] >> 4)];
|
||||
*ptr++= b64_table64_2[((src[1] << 2) & 0x3C) | (src[2] >> 6)];
|
||||
*ptr++= b64_table64_2[src[2] & 0x3F];
|
||||
src += 3;
|
||||
len -= 3;
|
||||
}
|
||||
|
||||
if (len > 0) {
|
||||
*ptr++= b64_table64_2[ buf[0] >> 2];
|
||||
uint8_t tmp = (buf[0] << 4) & 0x30;
|
||||
if (len > 1) tmp |= buf[1] >> 4;
|
||||
*ptr++= b64_table64_2[ src[0] >> 2];
|
||||
uint8_t tmp = (src[0] << 4) & 0x30;
|
||||
if (len > 1) tmp |= src[1] >> 4;
|
||||
*ptr++= b64_table64_2[tmp];
|
||||
if (len < 2) {
|
||||
*ptr++= '=';
|
||||
}
|
||||
else {
|
||||
*ptr++= b64_table64_2[buf[2] & 0x3F];
|
||||
*ptr++= b64_table64_2[(src[1] << 2) & 0x3C];
|
||||
}
|
||||
*ptr++= '=';
|
||||
}
|
||||
|
||||
*ptr = 0;
|
||||
*d_len = (ptr - dst);
|
||||
*d_len = ((char *)ptr - dst);
|
||||
|
||||
return dst;
|
||||
}
|
||||
|
||||
uint16_t uwsgi_be16(char *buf) {
|
||||
uint16_t *src = (uint16_t *) buf;
|
||||
uint16_t ret = 0;
|
||||
uint8_t *ptr = (uint8_t *) &ret;
|
||||
ptr[0] = (uint8_t) ((*src >> 8) & 0xff);
|
||||
ptr[1] = (uint8_t) (*src & 0xff);
|
||||
return ret;
|
||||
}
|
||||
|
||||
uint32_t uwsgi_be32(char *buf) {
|
||||
uint32_t *src = (uint32_t *) buf;
|
||||
uint32_t ret = 0;
|
||||
uint8_t *ptr = (uint8_t *) &ret;
|
||||
ptr[0] = (uint8_t) ((*src >> 24) & 0xff);
|
||||
ptr[1] = (uint8_t) ((*src >> 16) & 0xff);
|
||||
ptr[2] = (uint8_t) ((*src >> 8) & 0xff);
|
||||
ptr[3] = (uint8_t) (*src & 0xff);
|
||||
return ret;
|
||||
}
|
||||
|
||||
uint64_t uwsgi_be64(char *buf) {
|
||||
uint64_t *src = (uint64_t *) buf;
|
||||
uint64_t ret = 0;
|
||||
uint8_t *ptr = (uint8_t *) &ret;
|
||||
ptr[0] = (uint8_t) ((*src >> 56) & 0xff);
|
||||
ptr[1] = (uint8_t) ((*src >> 48) & 0xff);
|
||||
ptr[2] = (uint8_t) ((*src >> 40) & 0xff);
|
||||
ptr[3] = (uint8_t) ((*src >> 32) & 0xff);
|
||||
ptr[4] = (uint8_t) ((*src >> 24) & 0xff);
|
||||
ptr[5] = (uint8_t) ((*src >> 16) & 0xff);
|
||||
ptr[6] = (uint8_t) ((*src >> 8) & 0xff);
|
||||
ptr[7] = (uint8_t) (*src & 0xff);
|
||||
return ret;
|
||||
}
|
||||
|
||||
@@ -471,6 +471,15 @@ static struct uwsgi_option uwsgi_base_options[] = {
|
||||
{"routers-list", no_argument, 0, "list enabled routers", uwsgi_opt_true, &uwsgi.router_list, 0},
|
||||
#endif
|
||||
|
||||
{"websockets-ping-freq", required_argument, 0, "set the frequency (in seconds) of websockets automatic ping packets", uwsgi_opt_set_int, &uwsgi.websockets_ping_freq, 0},
|
||||
{"websocket-ping-freq", required_argument, 0, "set the frequency (in seconds) of websockets automatic ping packets", uwsgi_opt_set_int, &uwsgi.websockets_ping_freq, 0},
|
||||
|
||||
{"websockets-pong-freq", required_argument, 0, "set the frequency (in seconds) of websockets automatic pong/keepalive packets", uwsgi_opt_set_int, &uwsgi.websockets_pong_freq, 0},
|
||||
{"websocket-pong-freq", required_argument, 0, "set the frequency (in seconds) of websockets automatic pong/keepalive packets", uwsgi_opt_set_int, &uwsgi.websockets_pong_freq, 0},
|
||||
|
||||
{"websockets-max-size", required_argument, 0, "set the max allowed size of websocket messages (in Kbytes, default 1024)", uwsgi_opt_set_64bit, &uwsgi.websockets_max_size, 0},
|
||||
{"websocket-max-size", required_argument, 0, "set the max allowed size of websocket messages (in Kbytes, default 1024)", uwsgi_opt_set_64bit, &uwsgi.websockets_max_size, 0},
|
||||
|
||||
{"clock", required_argument, 0, "set a clock source", uwsgi_opt_set_str, &uwsgi.requested_clock, 0},
|
||||
|
||||
{"clock-list", no_argument, 0, "list enabled clocks", uwsgi_opt_true, &uwsgi.clock_list, 0},
|
||||
|
||||
+330
-379
@@ -16,6 +16,106 @@ extern struct uwsgi_server uwsgi;
|
||||
|
||||
#include "cr.h"
|
||||
|
||||
struct corerouter_peer *uwsgi_cr_peer_find_by_sid(struct corerouter_session *cs, uint32_t sid) {
|
||||
|
||||
struct corerouter_peer *peers = cs->peers;
|
||||
while(peers) {
|
||||
if (peers->sid == sid) {
|
||||
return peers;
|
||||
}
|
||||
peers = peers->next;
|
||||
}
|
||||
return NULL;
|
||||
}
|
||||
|
||||
// add a new peer to the session
|
||||
struct corerouter_peer *uwsgi_cr_peer_add(struct corerouter_session *cs) {
|
||||
struct corerouter_peer *old_peers = NULL, *peers = cs->peers;
|
||||
|
||||
while(peers) {
|
||||
old_peers = peers;
|
||||
peers = peers->next;
|
||||
}
|
||||
|
||||
peers = uwsgi_calloc(sizeof(struct corerouter_peer));
|
||||
peers->session = cs;
|
||||
peers->fd = -1;
|
||||
// create input buffer
|
||||
peers->in = uwsgi_buffer_new(uwsgi.page_size);
|
||||
// add timeout
|
||||
peers->timeout = cr_add_timeout(cs->corerouter, peers);
|
||||
// check retry
|
||||
if (cs->retry) {
|
||||
peers->can_retry = 1;
|
||||
}
|
||||
peers->prev = old_peers;
|
||||
|
||||
if (old_peers) {
|
||||
old_peers->next = peers;
|
||||
}
|
||||
else {
|
||||
cs->peers = peers;
|
||||
}
|
||||
|
||||
cs->refcnt++;
|
||||
|
||||
return peers;
|
||||
}
|
||||
|
||||
// reset a peer (allows it to connect to another backend)
|
||||
void uwsgi_cr_peer_reset(struct corerouter_peer *peer) {
|
||||
if (peer->tmp_socket_name) {
|
||||
free(peer->tmp_socket_name);
|
||||
peer->tmp_socket_name = NULL;
|
||||
}
|
||||
cr_del_timeout(peer->session->corerouter, peer);
|
||||
|
||||
if (peer->fd != -1) {
|
||||
close(peer->fd);
|
||||
peer->session->corerouter->cr_table[peer->fd] = NULL;
|
||||
peer->fd = -1;
|
||||
peer->hook_read = NULL;
|
||||
peer->hook_write = NULL;
|
||||
}
|
||||
|
||||
peer->failed = 0;
|
||||
peer->soopt = 0;
|
||||
peer->timed_out = 0;
|
||||
|
||||
peer->un = NULL;
|
||||
peer->static_node = NULL;
|
||||
}
|
||||
|
||||
// destroy a peer
|
||||
void uwsgi_cr_peer_del(struct corerouter_peer *peer) {
|
||||
struct corerouter_peer *prev = peer->prev;
|
||||
struct corerouter_peer *next = peer->next;
|
||||
|
||||
if (prev) {
|
||||
prev->next = peer->next;
|
||||
}
|
||||
|
||||
if (next) {
|
||||
next->prev = peer->prev;
|
||||
}
|
||||
|
||||
if (peer == peer->session->peers) {
|
||||
peer->session->peers = peer->next;
|
||||
}
|
||||
|
||||
uwsgi_cr_peer_reset(peer);
|
||||
|
||||
if (peer->in) {
|
||||
uwsgi_buffer_destroy(peer->in);
|
||||
}
|
||||
|
||||
// main_peer bring the output buffer from backend peers
|
||||
if (peer->out && peer->out_need_free) {
|
||||
uwsgi_buffer_destroy(peer->out);
|
||||
}
|
||||
free(peer);
|
||||
}
|
||||
|
||||
void uwsgi_opt_corerouter(char *opt, char *value, void *cr) {
|
||||
struct uwsgi_corerouter *ucr = (struct uwsgi_corerouter *) cr;
|
||||
uwsgi_new_gateway_socket(value, ucr->name);
|
||||
@@ -179,160 +279,148 @@ void corerouter_manage_subscription(char *key, uint16_t keylen, char *val, uint1
|
||||
}
|
||||
}
|
||||
|
||||
static struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
static struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
|
||||
void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
void corerouter_close_peer(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
struct corerouter_session *cs = peer->session;
|
||||
|
||||
|
||||
if (cr_session->instance_fd != -1) {
|
||||
close(cr_session->instance_fd);
|
||||
ucr->cr_table[cr_session->instance_fd] = NULL;
|
||||
}
|
||||
|
||||
if (ucr->subscriptions && cr_session->un && cr_session->un->len > 0) {
|
||||
// decrease reference count
|
||||
|
||||
// manage subscription reference count
|
||||
if (ucr->subscriptions && peer->un && peer->un->len > 0) {
|
||||
// decrease reference count
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("[1] node %.*s refcnt: %llu\n", cr_session->un->len, cr_session->un->name, cr_session->un->reference);
|
||||
uwsgi_log("[1] node %.*s refcnt: %llu\n", peer->un->len, peer->un->name, peer->un->reference);
|
||||
#endif
|
||||
cr_session->un->reference--;
|
||||
peer->un->reference--;
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("[2] node %.*s refcnt: %llu\n", cr_session->un->len, cr_session->un->name, cr_session->un->reference);
|
||||
uwsgi_log("[2] node %.*s refcnt: %llu\n", peer->un->len, peer->un->name, peer->un->reference);
|
||||
#endif
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
if (peer->failed) {
|
||||
|
||||
if (peer->soopt) {
|
||||
if (!ucr->quiet)
|
||||
uwsgi_log("[uwsgi-%s] unable to connect() to node \"%.*s\": %s\n", ucr->short_name, (int) peer->instance_address_len, peer->instance_address, strerror(peer->soopt));
|
||||
}
|
||||
else if (peer->timed_out) {
|
||||
if (peer->instance_address_len > 0) {
|
||||
if (peer->connecting) {
|
||||
if (!ucr->quiet)
|
||||
uwsgi_log("[uwsgi-%s] unable to connect() to node \"%.*s\": timeout\n", ucr->short_name, (int) peer->instance_address_len, peer->instance_address);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (cr_session->instance_failed) {
|
||||
// now check for dead nodes
|
||||
if (ucr->subscriptions && peer->un && peer->un->len > 0) {
|
||||
|
||||
if (cr_session->soopt) {
|
||||
if (!ucr->quiet)
|
||||
uwsgi_log("[uwsgi-%s] unable to connect() to node \"%.*s\": %s\n", ucr->short_name, (int) cr_session->instance_address_len, cr_session->instance_address, strerror(cr_session->soopt));
|
||||
}
|
||||
else if (cr_session->timed_out) {
|
||||
if (cr_session->instance_address_len > 0) {
|
||||
if (cr_session->connecting) {
|
||||
if (!ucr->quiet)
|
||||
uwsgi_log("[uwsgi-%s] unable to connect() to node \"%.*s\": timeout\n", ucr->short_name, (int) cr_session->instance_address_len, cr_session->instance_address);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (peer->un->death_mark == 0)
|
||||
uwsgi_log("[uwsgi-%s] %.*s => marking %.*s as failed\n", ucr->short_name, (int) peer->key_len, peer->key, (int) peer->instance_address_len, peer->instance_address);
|
||||
|
||||
// now check for dead nodes
|
||||
if (ucr->subscriptions && cr_session->un && cr_session->un->len > 0) {
|
||||
|
||||
if (cr_session->un->death_mark == 0)
|
||||
uwsgi_log("[uwsgi-%s] %.*s => marking %.*s as failed\n", ucr->short_name, (int) cr_session->hostname_len, cr_session->hostname, (int) cr_session->instance_address_len, cr_session->instance_address);
|
||||
|
||||
cr_session->un->failcnt++;
|
||||
cr_session->un->death_mark = 1;
|
||||
peer->un->failcnt++;
|
||||
peer->un->death_mark = 1;
|
||||
// check if i can remove the node
|
||||
if (cr_session->un->reference == 0) {
|
||||
uwsgi_remove_subscribe_node(ucr->subscriptions, cr_session->un);
|
||||
if (peer->un->reference == 0) {
|
||||
uwsgi_remove_subscribe_node(ucr->subscriptions, peer->un);
|
||||
}
|
||||
if (ucr->cheap && !ucr->i_am_cheap && !ucr->fallback && uwsgi_no_subscriptions(ucr->subscriptions)) {
|
||||
uwsgi_gateway_go_cheap(ucr->name, ucr->queue, &ucr->i_am_cheap);
|
||||
}
|
||||
|
||||
}
|
||||
else if (cr_session->static_node) {
|
||||
cr_session->static_node->custom = uwsgi_now();
|
||||
uwsgi_log("[uwsgi-%s] %.*s => marking %.*s as failed\n", ucr->short_name, (int) cr_session->hostname_len, cr_session->hostname, (int) cr_session->instance_address_len, cr_session->instance_address);
|
||||
}
|
||||
else if (peer->static_node) {
|
||||
peer->static_node->custom = uwsgi_now();
|
||||
uwsgi_log("[uwsgi-%s] %.*s => marking %.*s as failed\n", ucr->short_name, (int) peer->key_len, peer->key, (int) peer->instance_address_len, peer->instance_address);
|
||||
}
|
||||
|
||||
// check if the router supports the retry hook
|
||||
if (!peer->can_retry) goto end;
|
||||
if (peer->retries >= (size_t) ucr->max_retries) goto end;
|
||||
|
||||
if (cr_session->tmp_socket_name) {
|
||||
free(cr_session->tmp_socket_name);
|
||||
cr_session->tmp_socket_name = NULL;
|
||||
}
|
||||
|
||||
if (!cr_session->retry) goto end;
|
||||
// check for max retries
|
||||
if (cr_session->retries >= (size_t) ucr->max_retries) goto end;
|
||||
|
||||
cr_session->retries++;
|
||||
|
||||
// reset error and timeout
|
||||
cr_session->instance_failed = 0;
|
||||
cr_session->timeout = corerouter_reset_timeout(ucr, cr_session);
|
||||
cr_session->timed_out = 0;
|
||||
cr_session->soopt = 0;
|
||||
|
||||
// reset nodes
|
||||
cr_session->un = NULL;
|
||||
cr_session->static_node = NULL;
|
||||
cr_session->instance_fd = -1;
|
||||
|
||||
// reset hooks (safe as fd is closed)
|
||||
cr_session->event_hook_read = NULL;
|
||||
cr_session->event_hook_write = NULL;
|
||||
cr_session->event_hook_instance_read = NULL;
|
||||
cr_session->event_hook_instance_write = NULL;
|
||||
peer->retries++;
|
||||
// reset the peer
|
||||
uwsgi_cr_peer_reset(peer);
|
||||
// set new timeout
|
||||
peer->timeout = cr_add_timeout(ucr, peer);
|
||||
|
||||
if (ucr->fallback) {
|
||||
// ok let's try with the fallback nodes
|
||||
if (!cr_session->fallback) {
|
||||
cr_session->fallback = ucr->fallback;
|
||||
}
|
||||
else {
|
||||
cr_session->fallback = cr_session->fallback->next;
|
||||
if (!cr_session->fallback) goto end;
|
||||
}
|
||||
if (!cs->fallback) {
|
||||
cs->fallback = ucr->fallback;
|
||||
}
|
||||
else {
|
||||
cs->fallback = cs->fallback->next;
|
||||
if (!cs->fallback) goto end;
|
||||
}
|
||||
|
||||
cr_session->instance_address = cr_session->fallback->value;
|
||||
cr_session->instance_address_len = cr_session->fallback->len;
|
||||
peer->instance_address = cs->fallback->value;
|
||||
peer->instance_address_len = cs->fallback->len;
|
||||
|
||||
if (cr_session->retry(ucr, cr_session)) {
|
||||
if (!cr_session->instance_failed) goto end;
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
cr_session->instance_address = NULL;
|
||||
cr_session->instance_address_len = 0;
|
||||
if (cr_session->retry(ucr, cr_session)) {
|
||||
if (!cr_session->instance_failed) goto end;
|
||||
if (cs->retry(peer)) {
|
||||
if (!peer->failed) goto end;
|
||||
}
|
||||
return;
|
||||
}
|
||||
return;
|
||||
|
||||
peer->instance_address = NULL;
|
||||
peer->instance_address_len = 0;
|
||||
if (cs->retry(peer)) {
|
||||
if (!peer->failed) goto end;
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
end:
|
||||
uwsgi_cr_peer_del(peer);
|
||||
|
||||
if (cr_session->tmp_socket_name) {
|
||||
free(cr_session->tmp_socket_name);
|
||||
if (peer == cs->main_peer) {
|
||||
cs->main_peer = NULL;
|
||||
corerouter_close_session(ucr, cs);
|
||||
}
|
||||
else {
|
||||
if (cs->refcnt > 0)
|
||||
cs->refcnt--;
|
||||
if (cs->refcnt == 0) {
|
||||
corerouter_close_session(ucr, cs);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// destroy a session
|
||||
void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
|
||||
struct corerouter_peer *main_peer = cr_session->main_peer;
|
||||
if (main_peer) {
|
||||
uwsgi_cr_peer_del(main_peer);
|
||||
}
|
||||
|
||||
if (cr_session->buf_file)
|
||||
fclose(cr_session->buf_file);
|
||||
|
||||
if (cr_session->buf_file_name) {
|
||||
if (unlink(cr_session->buf_file_name)) {
|
||||
uwsgi_error("unlink()");
|
||||
}
|
||||
free(cr_session->buf_file_name);
|
||||
// free peers
|
||||
struct corerouter_peer *peers = cr_session->peers;
|
||||
while(peers) {
|
||||
struct corerouter_peer *tmp_peer = peers;
|
||||
peers = peers->next;
|
||||
uwsgi_cr_peer_del(tmp_peer);
|
||||
}
|
||||
|
||||
// could be used to free additional resources
|
||||
if (cr_session->close)
|
||||
cr_session->close(cr_session);
|
||||
|
||||
close(cr_session->fd);
|
||||
ucr->cr_table[cr_session->fd] = NULL;
|
||||
|
||||
uwsgi_buffer_destroy(cr_session->buffer);
|
||||
|
||||
cr_del_timeout(ucr, cr_session);
|
||||
free(cr_session);
|
||||
}
|
||||
|
||||
static struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
cr_del_timeout(ucr, cr_session);
|
||||
return cr_add_timeout(ucr, cr_session);
|
||||
static struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
cr_del_timeout(ucr, peer);
|
||||
return cr_add_timeout(ucr, peer);
|
||||
}
|
||||
|
||||
static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) {
|
||||
|
||||
time_t current = uwsgi_now();
|
||||
struct uwsgi_rb_timer *urbt;
|
||||
struct corerouter_session *cr_session;
|
||||
struct corerouter_peer *peer;
|
||||
|
||||
for (;;) {
|
||||
urbt = uwsgi_min_rb_timer(ucr->timeouts);
|
||||
@@ -340,12 +428,12 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) {
|
||||
return;
|
||||
|
||||
if (urbt->key <= current) {
|
||||
cr_session = (struct corerouter_session *) urbt->data;
|
||||
cr_session->timed_out = 1;
|
||||
if (cr_session->connecting) {
|
||||
cr_session->instance_failed = 1;
|
||||
peer = (struct corerouter_peer *) urbt->data;
|
||||
peer->timed_out = 1;
|
||||
if (peer->connecting) {
|
||||
peer->failed = 1;
|
||||
}
|
||||
corerouter_close_session(ucr, cr_session);
|
||||
corerouter_close_peer(ucr, peer);
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -353,237 +441,132 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
int uwsgi_cr_hook_read(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) {
|
||||
|
||||
int uwsgi_cr_set_hooks(struct corerouter_peer *peer, ssize_t (*read_hook)(struct corerouter_peer *), ssize_t (*write_hook)(struct corerouter_peer *)) {
|
||||
struct corerouter_session *cs = peer->session;
|
||||
struct uwsgi_corerouter *ucr = cs->corerouter;
|
||||
|
||||
// first check the case of event removal
|
||||
if (hook == NULL) {
|
||||
// nothing changed
|
||||
if (!cs->event_hook_read) goto unchanged;
|
||||
// if there is a write event defined, le'ts modify it
|
||||
if (cs->event_hook_write) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_readwrite_to_write() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_fd_readwrite_to_write(ucr->queue, cs->fd)) return -1;
|
||||
//uwsgi_log("uwsgi_cr_set_hooks(%d, %p, %p)\n", peer->fd, read_hook, write_hook);
|
||||
|
||||
if (read_hook) {
|
||||
peer->last_hook_read = read_hook;
|
||||
}
|
||||
|
||||
if (write_hook) {
|
||||
peer->last_hook_write = write_hook;
|
||||
}
|
||||
|
||||
int read_changed = 1;
|
||||
int write_changed = 1;
|
||||
|
||||
if (read_hook && peer->hook_read) {
|
||||
read_changed = 0;
|
||||
}
|
||||
else if (!read_hook && !peer->hook_read) {
|
||||
read_changed = 0;
|
||||
}
|
||||
|
||||
if (write_hook && peer->hook_write) {
|
||||
write_changed = 0;
|
||||
}
|
||||
else if (!write_hook && !peer->hook_write) {
|
||||
write_changed = 0;
|
||||
}
|
||||
|
||||
if (!read_changed && !write_changed) {
|
||||
goto unchanged;
|
||||
}
|
||||
|
||||
int has_read = 0;
|
||||
int has_write = 0;
|
||||
|
||||
if (peer->hook_read) {
|
||||
has_read = 1;
|
||||
}
|
||||
|
||||
if (peer->hook_write) {
|
||||
has_write = 1;
|
||||
}
|
||||
|
||||
if (!read_hook && !write_hook) {
|
||||
if (has_read) {
|
||||
if (event_queue_del_fd(ucr->queue, peer->fd, event_queue_read())) return -1;
|
||||
}
|
||||
// simply remove the read event
|
||||
else {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_del_fd() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_del_fd(ucr->queue, cs->fd, event_queue_read())) return -1;
|
||||
if (has_write) {
|
||||
if (event_queue_del_fd(ucr->queue, peer->fd, event_queue_write())) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
// set the hook
|
||||
// if write is not defined, simply add a single monitor
|
||||
if (cs->event_hook_write == NULL) {
|
||||
if (!cs->event_hook_read) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_add_fd_read() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_add_fd_read(ucr->queue, cs->fd)) return -1;
|
||||
else if (read_hook && write_hook) {
|
||||
if (has_read) {
|
||||
if (event_queue_fd_read_to_readwrite(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
else if (has_write) {
|
||||
if (event_queue_fd_write_to_readwrite(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
}
|
||||
else if (read_hook) {
|
||||
if (has_write) {
|
||||
if (write_changed) {
|
||||
if (event_queue_fd_write_to_read(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
else {
|
||||
if (event_queue_fd_write_to_readwrite(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (!cs->event_hook_read) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_write_to_readwrite() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_fd_write_to_readwrite(ucr->queue, cs->fd)) return -1;
|
||||
if (event_queue_add_fd_read(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
}
|
||||
else if (write_hook) {
|
||||
if (has_read) {
|
||||
if (read_changed) {
|
||||
if (event_queue_fd_read_to_write(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
else {
|
||||
if (event_queue_fd_read_to_readwrite(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (event_queue_add_fd_write(ucr->queue, peer->fd)) return -1;
|
||||
}
|
||||
}
|
||||
|
||||
unchanged:
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_hook_read set to %p for %d\n", hook, cs->fd);
|
||||
#endif
|
||||
cs->event_hook_read = hook;
|
||||
|
||||
peer->hook_read = read_hook;
|
||||
peer->hook_write = write_hook;
|
||||
return 0;
|
||||
|
||||
}
|
||||
|
||||
int uwsgi_cr_hook_write(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) {
|
||||
|
||||
struct uwsgi_corerouter *ucr = cs->corerouter;
|
||||
|
||||
// first check the case of event removal
|
||||
if (hook == NULL) {
|
||||
// nothing changed
|
||||
if (!cs->event_hook_write) goto unchanged;
|
||||
// if there is a read event defined, le'ts modify it
|
||||
if (cs->event_hook_read) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_readwrite_to_read() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_fd_readwrite_to_read(ucr->queue, cs->fd)) return -1;
|
||||
}
|
||||
// simply remove the write event
|
||||
else {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_del_fd() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_del_fd(ucr->queue, cs->fd, event_queue_write())) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
// set the hook
|
||||
// if read is not defined, simply add a single monitor
|
||||
if (cs->event_hook_read == NULL) {
|
||||
if (!cs->event_hook_write) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_add_fd_write() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_add_fd_write(ucr->queue, cs->fd)) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (!cs->event_hook_write) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_read_to_readwrite() for %d\n", cs->fd);
|
||||
#endif
|
||||
if (event_queue_fd_read_to_readwrite(ucr->queue, cs->fd)) return -1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
unchanged:
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_hook_write set to %p for %d\n", hook, cs->fd);
|
||||
#endif
|
||||
cs->event_hook_write = hook;
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_cr_hook_instance_read(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) {
|
||||
|
||||
struct uwsgi_corerouter *ucr = cs->corerouter;
|
||||
|
||||
// first check the case of event removal
|
||||
if (hook == NULL) {
|
||||
// nothing changed
|
||||
if (!cs->event_hook_instance_read) goto unchanged;
|
||||
// if there is a write event defined, le'ts modify it
|
||||
if (cs->event_hook_instance_write) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_readwrite_to_write() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_fd_readwrite_to_write(ucr->queue, cs->instance_fd)) return -1;
|
||||
}
|
||||
// simply remove the read event
|
||||
else {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_del_fd() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_del_fd(ucr->queue, cs->instance_fd, event_queue_read())) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
// set the hook
|
||||
// if write is not defined, simply add a single monitor
|
||||
if (cs->event_hook_instance_write == NULL) {
|
||||
if (!cs->event_hook_instance_read) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_add_fd_read() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_add_fd_read(ucr->queue, cs->instance_fd)) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (!cs->event_hook_instance_read) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_write_to_readwrite() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_fd_write_to_readwrite(ucr->queue, cs->instance_fd)) return -1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
unchanged:
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_hook_instance_read set to %p for %d\n", hook, cs->instance_fd);
|
||||
#endif
|
||||
cs->event_hook_instance_read = hook;
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_cr_hook_instance_write(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) {
|
||||
|
||||
struct uwsgi_corerouter *ucr = cs->corerouter;
|
||||
|
||||
// first check the case of event removal
|
||||
if (hook == NULL) {
|
||||
// nothing changed
|
||||
if (!cs->event_hook_instance_write) goto unchanged;
|
||||
// if there is a read event defined, le'ts modify it
|
||||
if (cs->event_hook_instance_read) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_readwrite_to_read() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_fd_readwrite_to_read(ucr->queue, cs->instance_fd)) return -1;
|
||||
}
|
||||
// simply remove the write event
|
||||
else {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_del_fd() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_del_fd(ucr->queue, cs->instance_fd, event_queue_write())) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
// set the hook
|
||||
// if read is not defined, simply add a single monitor
|
||||
if (cs->event_hook_instance_read == NULL) {
|
||||
if (!cs->event_hook_instance_write) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_add_fd_write() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_add_fd_write(ucr->queue, cs->instance_fd)) return -1;
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (!cs->event_hook_instance_write) {
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_queue_fd_read_to_readwrite() for %d\n", cs->instance_fd);
|
||||
#endif
|
||||
if (event_queue_fd_read_to_readwrite(ucr->queue, cs->instance_fd)) return -1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
unchanged:
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("event_hook_instance_write set to %p for %d\n", hook, cs->instance_fd);
|
||||
#endif
|
||||
cs->event_hook_instance_write = hook;
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
|
||||
struct corerouter_session *corerouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, int new_connection, struct sockaddr *cr_addr, socklen_t cr_addr_len) {
|
||||
|
||||
ucr->cr_table[new_connection] = uwsgi_calloc(ucr->session_size);
|
||||
ucr->cr_table[new_connection]->fd = new_connection;
|
||||
ucr->cr_table[new_connection]->instance_fd = -1;
|
||||
struct corerouter_session *cs = uwsgi_calloc(ucr->session_size);
|
||||
|
||||
struct corerouter_peer *peer = uwsgi_calloc(sizeof(struct corerouter_peer));
|
||||
// main_peer has only input buffer as output buffer is taken from backend peers
|
||||
peer->in = uwsgi_buffer_new(uwsgi.page_size);
|
||||
|
||||
ucr->cr_table[new_connection] = peer;
|
||||
cs->main_peer = peer;
|
||||
|
||||
peer->fd = new_connection;
|
||||
peer->session = cs;
|
||||
|
||||
// map corerouter and socket
|
||||
ucr->cr_table[new_connection]->corerouter = ucr;
|
||||
ucr->cr_table[new_connection]->ugs = ugs;
|
||||
cs->corerouter = ucr;
|
||||
cs->ugs = ugs;
|
||||
|
||||
// set initial timeout
|
||||
ucr->cr_table[new_connection]->timeout = cr_add_timeout(ucr, ucr->cr_table[new_connection]);
|
||||
|
||||
// create dynamic buffer
|
||||
ucr->cr_table[new_connection]->buffer = uwsgi_buffer_new(uwsgi.page_size);
|
||||
peer->timeout = cr_add_timeout(ucr, ucr->cr_table[new_connection]);
|
||||
|
||||
// here we prepare the real session and set the hooks
|
||||
ucr->alloc_session(ucr, ugs, ucr->cr_table[new_connection], cr_addr, cr_addr_len);
|
||||
if (ucr->alloc_session(ucr, ugs, cs, cr_addr, cr_addr_len)) {
|
||||
uwsgi_cr_peer_del(cs->main_peer);
|
||||
free(cs);
|
||||
cs = NULL;
|
||||
}
|
||||
|
||||
return ucr->cr_table[new_connection];
|
||||
return cs;
|
||||
}
|
||||
|
||||
void uwsgi_corerouter_loop(int id, void *data) {
|
||||
@@ -662,7 +645,6 @@ void uwsgi_corerouter_loop(int id, void *data) {
|
||||
|
||||
struct uwsgi_rb_timer *min_timeout;
|
||||
|
||||
int interesting_fd;
|
||||
int new_connection;
|
||||
|
||||
|
||||
@@ -673,8 +655,6 @@ void uwsgi_corerouter_loop(int id, void *data) {
|
||||
union uwsgi_sockaddr cr_addr;
|
||||
socklen_t cr_addr_len = sizeof(struct sockaddr_un);
|
||||
|
||||
struct corerouter_session *cr_session;
|
||||
|
||||
ucr->mapper = uwsgi_cr_map_use_void;
|
||||
|
||||
if (ucr->use_cache) {
|
||||
@@ -741,24 +721,24 @@ void uwsgi_corerouter_loop(int id, void *data) {
|
||||
for (i = 0; i < nevents; i++) {
|
||||
|
||||
// get the interesting fd
|
||||
interesting_fd = event_queue_interesting_fd(events, i);
|
||||
ucr->interesting_fd = event_queue_interesting_fd(events, i);
|
||||
// something bad happened
|
||||
if (interesting_fd < 0) continue;
|
||||
if (ucr->interesting_fd < 0) continue;
|
||||
|
||||
// check if the interesting_fd matches a gateway socket
|
||||
// check if the ucr->interesting_fd matches a gateway socket
|
||||
struct uwsgi_gateway_socket *ugs = uwsgi.gateway_sockets;
|
||||
int taken = 0;
|
||||
while (ugs) {
|
||||
if (ugs->gateway == &ushared->gateways[id] && interesting_fd == ugs->fd) {
|
||||
if (ugs->gateway == &ushared->gateways[id] && ucr->interesting_fd == ugs->fd) {
|
||||
if (!ugs->subscription) {
|
||||
#if defined(__linux__) && defined(SOCK_NONBLOCK) && !defined(OBSOLETE_LINUX_KERNEL)
|
||||
new_connection = accept4(interesting_fd, (struct sockaddr *) &cr_addr, &cr_addr_len, SOCK_NONBLOCK);
|
||||
new_connection = accept4(ucr->interesting_fd, (struct sockaddr *) &cr_addr, &cr_addr_len, SOCK_NONBLOCK);
|
||||
if (new_connection < 0) {
|
||||
taken = 1;
|
||||
break;
|
||||
}
|
||||
#else
|
||||
new_connection = accept(interesting_fd, (struct sockaddr *) &cr_addr, &cr_addr_len);
|
||||
new_connection = accept(ucr->interesting_fd, (struct sockaddr *) &cr_addr, &cr_addr_len);
|
||||
if (new_connection < 0) {
|
||||
taken = 1;
|
||||
break;
|
||||
@@ -768,12 +748,9 @@ void uwsgi_corerouter_loop(int id, void *data) {
|
||||
uwsgi_socket_nb(new_connection);
|
||||
#endif
|
||||
#endif
|
||||
|
||||
struct corerouter_session *cr = corerouter_alloc_session(ucr, ugs, new_connection, (struct sockaddr *) &cr_addr, cr_addr_len);
|
||||
//something wrong in the allocation
|
||||
if (cr->instance_failed) {
|
||||
corerouter_close_session(ucr, cr);
|
||||
}
|
||||
if (!cr) break;
|
||||
}
|
||||
else if (ugs->subscription) {
|
||||
uwsgi_corerouter_manage_subscription(ucr, id, ugs);
|
||||
@@ -792,79 +769,53 @@ void uwsgi_corerouter_loop(int id, void *data) {
|
||||
}
|
||||
|
||||
// manage internal subscription
|
||||
if (interesting_fd == ushared->gateways[id].internal_subscription_pipe[1]) {
|
||||
uwsgi_corerouter_manage_internal_subscription(ucr, interesting_fd);
|
||||
if (ucr->interesting_fd == ushared->gateways[id].internal_subscription_pipe[1]) {
|
||||
uwsgi_corerouter_manage_internal_subscription(ucr, ucr->interesting_fd);
|
||||
}
|
||||
// manage a stats request
|
||||
else if (interesting_fd == ucr->cr_stats_server) {
|
||||
else if (ucr->interesting_fd == ucr->cr_stats_server) {
|
||||
corerouter_send_stats(ucr);
|
||||
}
|
||||
else {
|
||||
cr_session = ucr->cr_table[interesting_fd];
|
||||
struct corerouter_peer *peer = ucr->cr_table[ucr->interesting_fd];
|
||||
|
||||
// something is going wrong...
|
||||
if (cr_session == NULL)
|
||||
if (peer == NULL)
|
||||
continue;
|
||||
|
||||
// on error, destroy the session
|
||||
if (event_queue_interesting_fd_has_error(events, i)) {
|
||||
if (interesting_fd == cr_session->instance_fd) {
|
||||
cr_session->instance_failed = 1;
|
||||
}
|
||||
corerouter_close_session(ucr, cr_session);
|
||||
peer->failed = 1;
|
||||
corerouter_close_peer(ucr, peer);
|
||||
continue;
|
||||
}
|
||||
|
||||
// set timeout
|
||||
cr_session->timeout = corerouter_reset_timeout(ucr, cr_session);
|
||||
// set timeout (in main_peer too)
|
||||
peer->timeout = corerouter_reset_timeout(ucr, peer);
|
||||
peer->session->main_peer->timeout = corerouter_reset_timeout(ucr, peer->session->main_peer);
|
||||
|
||||
ssize_t (*hook)(struct corerouter_peer *) = NULL;
|
||||
|
||||
// call event hook
|
||||
ssize_t (*hook)(struct corerouter_session *) = NULL;
|
||||
if (interesting_fd == cr_session->fd) {
|
||||
if (event_queue_interesting_fd_is_read(events, i)) {
|
||||
hook = cr_session->event_hook_read;
|
||||
}
|
||||
else if (event_queue_interesting_fd_is_write(events, i)) {
|
||||
hook = cr_session->event_hook_write;
|
||||
}
|
||||
if (event_queue_interesting_fd_is_read(events, i)) {
|
||||
hook = peer->hook_read;
|
||||
}
|
||||
else if (interesting_fd == cr_session->instance_fd) {
|
||||
if (event_queue_interesting_fd_is_read(events, i)) {
|
||||
hook = cr_session->event_hook_instance_read;
|
||||
}
|
||||
else if (event_queue_interesting_fd_is_write(events, i)) {
|
||||
hook = cr_session->event_hook_instance_write;
|
||||
}
|
||||
}
|
||||
|
||||
// not having a hook could mean a previous event in the loop cleared it...
|
||||
if (!hook) {
|
||||
// a single event cannot be unexpected..
|
||||
if (nevents == 1) {
|
||||
if (interesting_fd == cr_session->instance_fd) {
|
||||
uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received from backend instance (fd: %d nevents: %d) !!!\n", interesting_fd, nevents);
|
||||
}
|
||||
else if (interesting_fd == cr_session->fd) {
|
||||
uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received from client (fd: %d nevents: %d)!!!\n", interesting_fd, nevents);
|
||||
}
|
||||
else {
|
||||
uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received !!!\n");
|
||||
}
|
||||
corerouter_close_session(ucr, cr_session);
|
||||
}
|
||||
continue;
|
||||
else if (event_queue_interesting_fd_is_write(events, i)) {
|
||||
hook = peer->hook_write;
|
||||
}
|
||||
|
||||
if (!hook) continue;
|
||||
// reset errno (as we use it for internal signalling)
|
||||
errno = 0;
|
||||
ssize_t ret = hook(cr_session);
|
||||
ssize_t ret = hook(peer);
|
||||
// connection closed
|
||||
if (ret == 0) {
|
||||
corerouter_close_session(ucr, cr_session);
|
||||
corerouter_close_peer(ucr, peer);
|
||||
continue;
|
||||
}
|
||||
else if (ret < 0) {
|
||||
if (errno == EINPROGRESS) continue;
|
||||
corerouter_close_session(ucr, cr_session);
|
||||
corerouter_close_peer(ucr, peer);
|
||||
continue;
|
||||
}
|
||||
|
||||
|
||||
+185
-67
@@ -4,9 +4,6 @@
|
||||
#define COREROUTER_STATUS_RESPONSE 3
|
||||
|
||||
#define cr_add_timeout(u, x) uwsgi_add_rb_timer(u->timeouts, time(NULL)+u->socket_timeout, x)
|
||||
#define cr_add_fake_timeout(u, x) uwsgi_add_rb_timer(u->timeouts, time(NULL)+1, x)
|
||||
#define cr_add_check_timeout(x) uwsgi_add_rb_timer(timeouts, time(NULL)+x, NULL)
|
||||
#define cr_del_check_timeout(x) rb_erase(&x->rbt, timeouts);
|
||||
#define cr_del_timeout(u, x) rb_erase(&x->timeout->rbt, u->timeouts); free(x->timeout);
|
||||
|
||||
#define cr_try_again if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) {\
|
||||
@@ -14,17 +11,174 @@
|
||||
return -1;\
|
||||
}
|
||||
|
||||
#define cr_write(peer, f) write(peer->fd, peer->out->buf + peer->out_pos, peer->out->pos - peer->out_pos);\
|
||||
if (len < 0) {\
|
||||
cr_try_again;\
|
||||
uwsgi_error(f);\
|
||||
return -1;\
|
||||
}\
|
||||
peer->out_pos += len;
|
||||
|
||||
#define cr_write_buf(peer, buf, f) write(peer->fd, buf + buf##_pos, buf->pos - buf##_pos);\
|
||||
if (len < 0) {\
|
||||
cr_try_again;\
|
||||
uwsgi_error(f);\
|
||||
return -1;\
|
||||
}\
|
||||
buf##_pos += len;
|
||||
|
||||
#define cr_write_complete(peer) peer->out_pos == peer->out->pos
|
||||
|
||||
#define cr_write_complete_buf(peer, buf) buf##_pos == buf->pos
|
||||
|
||||
#define cr_connect(peer, f) peer->fd = uwsgi_connectn(peer->instance_address, peer->instance_address_len, 0, 1);\
|
||||
if (peer->fd < 0) {\
|
||||
peer->failed = 1;\
|
||||
peer->soopt = errno;\
|
||||
return -1;\
|
||||
}\
|
||||
peer->session->corerouter->cr_table[peer->fd] = peer;\
|
||||
peer->connecting = 1;\
|
||||
cr_write_to_backend(peer, f);
|
||||
|
||||
#define cr_read(peer, f) read(peer->fd, peer->in->buf + peer->in->pos, peer->in->len - peer->in->pos);\
|
||||
if (len < 0) {\
|
||||
cr_try_again;\
|
||||
uwsgi_error(f);\
|
||||
return -1;\
|
||||
}\
|
||||
peer->in->pos += len;\
|
||||
|
||||
#define cr_read_exact(peer, l, f) read(peer->fd, peer->in->buf + peer->in->pos,l - peer->in->pos);\
|
||||
if (len < 0) {\
|
||||
cr_try_again;\
|
||||
uwsgi_error(f);\
|
||||
return -1;\
|
||||
}\
|
||||
peer->in->pos += len;\
|
||||
|
||||
#define cr_reset_hooks(peer) if (uwsgi_cr_set_hooks(peer->session->main_peer, peer->session->main_peer->last_hook_read, NULL)) return -1;\
|
||||
struct corerouter_peer *peers = peer->session->peers;\
|
||||
while(peers) {\
|
||||
if (uwsgi_cr_set_hooks(peers, peers->last_hook_read, NULL)) return -1;\
|
||||
peers = peers->next;\
|
||||
}
|
||||
|
||||
#define cr_reset_hooks_and_read(peer, f) if (uwsgi_cr_set_hooks(peer->session->main_peer, peer->session->main_peer->last_hook_read, NULL)) return -1;\
|
||||
peer->last_hook_read = f;\
|
||||
struct corerouter_peer *peers = peer->session->peers;\
|
||||
while(peers) {\
|
||||
if (uwsgi_cr_set_hooks(peers, peers->last_hook_read, NULL)) return -1;\
|
||||
peers = peers->next;\
|
||||
}
|
||||
|
||||
#define cr_write_to_main(peer, f) if (uwsgi_cr_set_hooks(peer->session->main_peer, NULL, f)) return -1;\
|
||||
struct corerouter_peer *peers = peer->session->peers;\
|
||||
while(peers) {\
|
||||
if (uwsgi_cr_set_hooks(peers, NULL, NULL)) return -1;\
|
||||
peers = peers->next;\
|
||||
}
|
||||
|
||||
#define cr_write_to_backend(peer, f) if (uwsgi_cr_set_hooks(peer->session->main_peer, NULL, NULL)) return -1;\
|
||||
if (uwsgi_cr_set_hooks(peer, NULL, f)) return -1;\
|
||||
struct corerouter_peer *peers = peer->session->peers;\
|
||||
while(peers) {\
|
||||
if (peers != peer) {\
|
||||
if (uwsgi_cr_set_hooks(peers, NULL, NULL)) return -1;\
|
||||
}\
|
||||
peers = peers->next;\
|
||||
}
|
||||
|
||||
#define cr_peer_connected(peer, f) socklen_t solen = sizeof(int);\
|
||||
if (getsockopt(peer->fd, SOL_SOCKET, SO_ERROR, (void *) (&peer->soopt), &solen) < 0) {\
|
||||
uwsgi_error(f "/getsockopt()");\
|
||||
peer->failed = 1;\
|
||||
return -1;\
|
||||
}\
|
||||
if (peer->soopt) {\
|
||||
peer->failed = 1;\
|
||||
return -1;\
|
||||
}\
|
||||
peer->connecting = 0;\
|
||||
peer->can_retry = 0;\
|
||||
if (peer->static_node) peer->static_node->custom2++;\
|
||||
if (peer->un) peer->un->requests++;\
|
||||
|
||||
|
||||
struct corerouter_session;
|
||||
|
||||
// a peer is a connection to a socket (a client or a backend) and can be monitored for events.
|
||||
struct corerouter_peer {
|
||||
// the file descriptor
|
||||
int fd;
|
||||
// the session
|
||||
struct corerouter_session *session;
|
||||
|
||||
// hook to run on a read event
|
||||
ssize_t (*hook_read)(struct corerouter_peer *);
|
||||
ssize_t (*last_hook_read)(struct corerouter_peer *);
|
||||
// hook to run on a write event
|
||||
ssize_t (*hook_write)(struct corerouter_peer *);
|
||||
ssize_t (*last_hook_write)(struct corerouter_peer *);
|
||||
|
||||
// has the peer failed ?
|
||||
int failed;
|
||||
// is the peer connecting ?
|
||||
int connecting;
|
||||
// is there a connection error ?
|
||||
int soopt;
|
||||
// has the peer timed out ?
|
||||
int timed_out;
|
||||
// the timeout rb_tree
|
||||
struct uwsgi_rb_timer *timeout;
|
||||
|
||||
// each peer can map to a different instance
|
||||
char *tmp_socket_name;
|
||||
char *instance_address;
|
||||
uint64_t instance_address_len;
|
||||
|
||||
// backend info
|
||||
struct uwsgi_subscribe_node *un;
|
||||
struct uwsgi_string_list *static_node;
|
||||
|
||||
// incoming data
|
||||
struct uwsgi_buffer *in;
|
||||
// data to send
|
||||
struct uwsgi_buffer *out;
|
||||
// amount of sent data (partial write management)
|
||||
size_t out_pos;
|
||||
int out_need_free;
|
||||
|
||||
// stream id (could have various use)
|
||||
uint32_t sid;
|
||||
|
||||
// internal parser status
|
||||
int r_parser_status;
|
||||
|
||||
// can retry ?
|
||||
int can_retry;
|
||||
// how many retries ?
|
||||
uint16_t retries;
|
||||
|
||||
// parsed key
|
||||
char *key;
|
||||
uint16_t key_len;
|
||||
|
||||
uint8_t modifier1;
|
||||
uint8_t modifier2;
|
||||
|
||||
struct corerouter_peer *prev;
|
||||
struct corerouter_peer *next;
|
||||
};
|
||||
|
||||
struct uwsgi_corerouter {
|
||||
|
||||
char *name;
|
||||
char *short_name;
|
||||
size_t session_size;
|
||||
|
||||
void (*alloc_session)(struct uwsgi_corerouter *, struct uwsgi_gateway_socket *, struct corerouter_session *, struct sockaddr *, socklen_t);
|
||||
int (*mapper)(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int (*alloc_session)(struct uwsgi_corerouter *, struct uwsgi_gateway_socket *, struct corerouter_session *, struct sockaddr *, socklen_t);
|
||||
int (*mapper)(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
|
||||
int has_sockets;
|
||||
int has_backends;
|
||||
@@ -76,7 +230,6 @@ struct uwsgi_corerouter {
|
||||
char *code_string_code;
|
||||
char *code_string_function;
|
||||
|
||||
|
||||
struct uwsgi_rb_timer *subscriptions_check;
|
||||
|
||||
int cheap;
|
||||
@@ -85,72 +238,37 @@ struct uwsgi_corerouter {
|
||||
int tolerance;
|
||||
int harakiri;
|
||||
|
||||
struct corerouter_session **cr_table;
|
||||
struct corerouter_peer **cr_table;
|
||||
|
||||
int interesting_fd;
|
||||
|
||||
};
|
||||
|
||||
// a session is started when a client connect to the router
|
||||
struct corerouter_session {
|
||||
|
||||
int fd;
|
||||
int instance_fd;
|
||||
|
||||
// corerouter related to this session
|
||||
struct uwsgi_corerouter *corerouter;
|
||||
// gateway socket related to this session
|
||||
struct uwsgi_gateway_socket *ugs;
|
||||
|
||||
// parsed hostname
|
||||
char *hostname;
|
||||
uint16_t hostname_len;
|
||||
|
||||
int has_key;
|
||||
int connecting;
|
||||
|
||||
char *instance_address;
|
||||
uint64_t instance_address_len;
|
||||
|
||||
struct uwsgi_subscribe_node *un;
|
||||
struct uwsgi_string_list *static_node;
|
||||
int soopt;
|
||||
int timed_out;
|
||||
|
||||
struct uwsgi_rb_timer *timeout;
|
||||
int instance_failed;
|
||||
|
||||
// check content_length
|
||||
size_t post_cl;
|
||||
size_t post_remains;
|
||||
|
||||
// the list of fallback instances
|
||||
struct uwsgi_string_list *fallback;
|
||||
|
||||
char *buf_file_name;
|
||||
FILE *buf_file;
|
||||
|
||||
char *tmp_socket_name;
|
||||
|
||||
// store the client address
|
||||
struct sockaddr_un addr;
|
||||
socklen_t addr_len;
|
||||
|
||||
// async hooks:
|
||||
// the session is watiting for this fd
|
||||
ssize_t (*event_hook_read)(struct corerouter_session *);
|
||||
ssize_t (*event_hook_write)(struct corerouter_session *);
|
||||
ssize_t (*event_hook_instance_read)(struct corerouter_session *);
|
||||
ssize_t (*event_hook_instance_write)(struct corerouter_session *);
|
||||
|
||||
void (*close)(struct corerouter_session *);
|
||||
int (*retry)(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
size_t retries;
|
||||
int (*retry)(struct corerouter_peer *);
|
||||
|
||||
struct uwsgi_buffer *buffer;
|
||||
size_t buffer_len;
|
||||
off_t buffer_pos;
|
||||
// this is the peer of the client
|
||||
struct corerouter_peer *main_peer;
|
||||
// this is the linked list of backends
|
||||
struct corerouter_peer *peers;
|
||||
|
||||
struct uwsgi_header uh;
|
||||
|
||||
uint8_t modifier1;
|
||||
uint8_t modifier2;
|
||||
// when it reaches 0 the session can be destroyed
|
||||
uint64_t refcnt;
|
||||
};
|
||||
|
||||
void uwsgi_opt_corerouter(char *, char *, void *);
|
||||
@@ -174,20 +292,20 @@ int uwsgi_corerouter_init(struct uwsgi_corerouter *);
|
||||
struct corerouter_session *corerouter_alloc_session(struct uwsgi_corerouter *, struct uwsgi_gateway_socket *, int, struct sockaddr *, socklen_t);
|
||||
void corerouter_close_session(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
|
||||
int uwsgi_cr_map_use_void(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_cache(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_pattern(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_cluster(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_subscription_dotsplit(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_base(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_cs(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_to(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *, struct corerouter_session *);
|
||||
int uwsgi_cr_map_use_void(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_cache(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_pattern(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_cluster(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_subscription_dotsplit(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_base(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_cs(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_to(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
|
||||
int uwsgi_corerouter_has_backends(struct uwsgi_corerouter *);
|
||||
|
||||
int uwsgi_cr_hook_read(struct corerouter_session *, ssize_t (*)(struct corerouter_session *));
|
||||
int uwsgi_cr_hook_write(struct corerouter_session *, ssize_t (*)(struct corerouter_session *));
|
||||
int uwsgi_cr_hook_instance_read(struct corerouter_session *, ssize_t (*)(struct corerouter_session *));
|
||||
int uwsgi_cr_hook_instance_write(struct corerouter_session *, ssize_t (*)(struct corerouter_session *));
|
||||
int uwsgi_cr_set_hooks(struct corerouter_peer *, ssize_t (*)(struct corerouter_peer *), ssize_t (*)(struct corerouter_peer *));
|
||||
struct corerouter_peer *uwsgi_cr_peer_add(struct corerouter_session *);
|
||||
struct corerouter_peer *uwsgi_cr_peer_find_by_sid(struct corerouter_session *, uint32_t);
|
||||
void corerouter_close_peer(struct uwsgi_corerouter *, struct corerouter_peer *);
|
||||
|
||||
+57
-57
@@ -4,38 +4,38 @@
|
||||
|
||||
extern struct uwsgi_server uwsgi;
|
||||
|
||||
int uwsgi_cr_map_use_void(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_void(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_cr_map_use_cache(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
cr_session->instance_address = uwsgi_cache_get(cr_session->hostname, cr_session->hostname_len, &cr_session->instance_address_len);
|
||||
char *cs_mod = uwsgi_str_contains(cr_session->instance_address, cr_session->instance_address_len, ',');
|
||||
int uwsgi_cr_map_use_cache(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
peer->instance_address = uwsgi_cache_get(peer->key, peer->key_len, &peer->instance_address_len);
|
||||
char *cs_mod = uwsgi_str_contains(peer->instance_address, peer->instance_address_len, ',');
|
||||
if (cs_mod) {
|
||||
cr_session->modifier1 = uwsgi_str_num(cs_mod + 1, (cr_session->instance_address_len - (cs_mod - cr_session->instance_address)) - 1);
|
||||
cr_session->instance_address_len = (cs_mod - cr_session->instance_address);
|
||||
peer->modifier1 = uwsgi_str_num(cs_mod + 1, (peer->instance_address_len - (cs_mod - peer->instance_address)) - 1);
|
||||
peer->instance_address_len = (cs_mod - peer->instance_address);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_cr_map_use_pattern(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_pattern(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
int tmp_socket_name_len = 0;
|
||||
ucr->magic_table['s'] = uwsgi_concat2n(cr_session->hostname, cr_session->hostname_len, "", 0);
|
||||
cr_session->tmp_socket_name = magic_sub(ucr->pattern, ucr->pattern_len, &tmp_socket_name_len, ucr->magic_table);
|
||||
ucr->magic_table['s'] = uwsgi_concat2n(peer->key, peer->key_len, "", 0);
|
||||
peer->tmp_socket_name = magic_sub(ucr->pattern, ucr->pattern_len, &tmp_socket_name_len, ucr->magic_table);
|
||||
free(ucr->magic_table['s']);
|
||||
cr_session->instance_address_len = tmp_socket_name_len;
|
||||
cr_session->instance_address = cr_session->tmp_socket_name;
|
||||
peer->instance_address_len = tmp_socket_name_len;
|
||||
peer->instance_address = peer->tmp_socket_name;
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
|
||||
cr_session->un = uwsgi_get_subscribe_node(ucr->subscriptions, cr_session->hostname, cr_session->hostname_len);
|
||||
if (cr_session->un && cr_session->un->len) {
|
||||
cr_session->instance_address = cr_session->un->name;
|
||||
cr_session->instance_address_len = cr_session->un->len;
|
||||
cr_session->modifier1 = cr_session->un->modifier1;
|
||||
peer->un = uwsgi_get_subscribe_node(ucr->subscriptions, peer->key, peer->key_len);
|
||||
if (peer->un && peer->un->len) {
|
||||
peer->instance_address = peer->un->name;
|
||||
peer->instance_address_len = peer->un->len;
|
||||
peer->modifier1 = peer->un->modifier1;
|
||||
}
|
||||
else if (ucr->cheap && !ucr->i_am_cheap && uwsgi_no_subscriptions(ucr->subscriptions)) {
|
||||
uwsgi_gateway_go_cheap(ucr->name, ucr->queue, &ucr->i_am_cheap);
|
||||
@@ -44,17 +44,17 @@ int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *ucr, struct coreroute
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_cr_map_use_subscription_dotsplit(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_subscription_dotsplit(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
|
||||
char *name = cr_session->hostname;
|
||||
uint16_t name_len = cr_session->hostname_len;
|
||||
char *name = peer->key;
|
||||
uint16_t name_len = peer->key_len;
|
||||
|
||||
split:
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("trying with %.*s\n", name_len, name);
|
||||
#endif
|
||||
cr_session->un = uwsgi_get_subscribe_node(ucr->subscriptions, name, name_len);
|
||||
if (!cr_session->un) {
|
||||
peer->un = uwsgi_get_subscribe_node(ucr->subscriptions, name, name_len);
|
||||
if (!peer->un) {
|
||||
char *next = memchr(name+1, '.', name_len-1);
|
||||
if (next) {
|
||||
name_len -= next - name;
|
||||
@@ -63,10 +63,10 @@ split:
|
||||
}
|
||||
}
|
||||
|
||||
if (cr_session->un && cr_session->un->len) {
|
||||
cr_session->instance_address = cr_session->un->name;
|
||||
cr_session->instance_address_len = cr_session->un->len;
|
||||
cr_session->modifier1 = cr_session->un->modifier1;
|
||||
if (peer->un && peer->un->len) {
|
||||
peer->instance_address = peer->un->name;
|
||||
peer->instance_address_len = peer->un->len;
|
||||
peer->modifier1 = peer->un->modifier1;
|
||||
}
|
||||
else if (ucr->cheap && !ucr->i_am_cheap && uwsgi_no_subscriptions(ucr->subscriptions)) {
|
||||
uwsgi_gateway_go_cheap(ucr->name, ucr->queue, &ucr->i_am_cheap);
|
||||
@@ -76,46 +76,46 @@ split:
|
||||
}
|
||||
|
||||
|
||||
int uwsgi_cr_map_use_base(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_base(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
|
||||
int tmp_socket_name_len = 0;
|
||||
|
||||
cr_session->tmp_socket_name = uwsgi_concat2nn(ucr->base, ucr->base_len, cr_session->hostname, cr_session->hostname_len, &tmp_socket_name_len);
|
||||
cr_session->instance_address_len = tmp_socket_name_len;
|
||||
cr_session->instance_address = cr_session->tmp_socket_name;
|
||||
peer->tmp_socket_name = uwsgi_concat2nn(ucr->base, ucr->base_len, peer->key, peer->key_len, &tmp_socket_name_len);
|
||||
peer->instance_address_len = tmp_socket_name_len;
|
||||
peer->instance_address = peer->tmp_socket_name;
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
int uwsgi_cr_map_use_cs(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_cs(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
if (uwsgi.p[ucr->code_string_modifier1]->code_string) {
|
||||
char *name = uwsgi_concat2("uwsgi_", ucr->short_name);
|
||||
cr_session->instance_address = uwsgi.p[ucr->code_string_modifier1]->code_string(name, ucr->code_string_code, ucr->code_string_function, cr_session->hostname, cr_session->hostname_len);
|
||||
peer->instance_address = uwsgi.p[ucr->code_string_modifier1]->code_string(name, ucr->code_string_code, ucr->code_string_function, peer->key, peer->key_len);
|
||||
free(name);
|
||||
if (cr_session->instance_address) {
|
||||
cr_session->instance_address_len = strlen(cr_session->instance_address);
|
||||
char *cs_mod = uwsgi_str_contains(cr_session->instance_address, cr_session->instance_address_len, ',');
|
||||
if (peer->instance_address) {
|
||||
peer->instance_address_len = strlen(peer->instance_address);
|
||||
char *cs_mod = uwsgi_str_contains(peer->instance_address, peer->instance_address_len, ',');
|
||||
if (cs_mod) {
|
||||
cr_session->modifier1 = uwsgi_str_num(cs_mod + 1, (cr_session->instance_address_len - (cs_mod - cr_session->instance_address)) - 1);
|
||||
cr_session->instance_address_len = (cs_mod - cr_session->instance_address);
|
||||
peer->modifier1 = uwsgi_str_num(cs_mod + 1, (peer->instance_address_len - (cs_mod - peer->instance_address)) - 1);
|
||||
peer->instance_address_len = (cs_mod - peer->instance_address);
|
||||
}
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_cr_map_use_to(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
cr_session->instance_address = ucr->to_socket->name;
|
||||
cr_session->instance_address_len = ucr->to_socket->name_len;
|
||||
int uwsgi_cr_map_use_to(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
peer->instance_address = ucr->to_socket->name;
|
||||
peer->instance_address_len = ucr->to_socket->name_len;
|
||||
return 0;
|
||||
}
|
||||
|
||||
int uwsgi_cr_map_use_cluster(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_cluster(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
#ifdef UWSGI_MULTICAST
|
||||
cr_session->instance_address = uwsgi_cluster_best_node();
|
||||
if (cr_session->instance_address) {
|
||||
cr_session->instance_address_len = strlen(cr_session->instance_address);
|
||||
peer->instance_address = uwsgi_cluster_best_node();
|
||||
if (peer->instance_address) {
|
||||
peer->instance_address_len = strlen(peer->instance_address);
|
||||
}
|
||||
#else
|
||||
uwsgi_log("uWSGI has been built without multicast/clustering support !!!\n");
|
||||
@@ -124,24 +124,24 @@ int uwsgi_cr_map_use_cluster(struct uwsgi_corerouter *ucr, struct corerouter_ses
|
||||
}
|
||||
|
||||
|
||||
int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *ucr, struct corerouter_session *cr_session) {
|
||||
int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
|
||||
if (!ucr->current_static_node) {
|
||||
ucr->current_static_node = ucr->static_nodes;
|
||||
}
|
||||
|
||||
cr_session->static_node = ucr->current_static_node;
|
||||
peer->static_node = ucr->current_static_node;
|
||||
|
||||
// is it a dead node ?
|
||||
if (cr_session->static_node->custom > 0) {
|
||||
if (peer->static_node->custom > 0) {
|
||||
|
||||
// gracetime passed ?
|
||||
if (cr_session->static_node->custom + ucr->static_node_gracetime <= (uint64_t) uwsgi_now()) {
|
||||
cr_session->static_node->custom = 0;
|
||||
if (peer->static_node->custom + ucr->static_node_gracetime <= (uint64_t) uwsgi_now()) {
|
||||
peer->static_node->custom = 0;
|
||||
}
|
||||
else {
|
||||
struct uwsgi_string_list *tmp_node = cr_session->static_node;
|
||||
struct uwsgi_string_list *next_node = cr_session->static_node->next;
|
||||
cr_session->static_node = NULL;
|
||||
struct uwsgi_string_list *tmp_node = peer->static_node;
|
||||
struct uwsgi_string_list *next_node = peer->static_node->next;
|
||||
peer->static_node = NULL;
|
||||
// needed for 1-node only setups
|
||||
if (!next_node)
|
||||
next_node = ucr->static_nodes;
|
||||
@@ -155,7 +155,7 @@ int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *ucr, struct coreroute
|
||||
break;
|
||||
|
||||
if (next_node->custom == 0) {
|
||||
cr_session->static_node = next_node;
|
||||
peer->static_node = next_node;
|
||||
break;
|
||||
}
|
||||
next_node = next_node->next;
|
||||
@@ -163,12 +163,12 @@ int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *ucr, struct coreroute
|
||||
}
|
||||
}
|
||||
|
||||
if (cr_session->static_node) {
|
||||
if (peer->static_node) {
|
||||
|
||||
cr_session->instance_address = cr_session->static_node->value;
|
||||
cr_session->instance_address_len = cr_session->static_node->len;
|
||||
peer->instance_address = peer->static_node->value;
|
||||
peer->instance_address_len = peer->static_node->len;
|
||||
// set the next one
|
||||
ucr->current_static_node = cr_session->static_node->next;
|
||||
ucr->current_static_node = peer->static_node->next;
|
||||
}
|
||||
else {
|
||||
// set the next one
|
||||
|
||||
+147
-239
@@ -14,11 +14,8 @@ struct uwsgi_fastrouter {
|
||||
extern struct uwsgi_server uwsgi;
|
||||
|
||||
struct fastrouter_session {
|
||||
struct corerouter_session crs;
|
||||
struct uwsgi_buffer *post_buf;
|
||||
size_t post_buf_max;
|
||||
size_t post_buf_len;
|
||||
off_t post_buf_pos;
|
||||
struct corerouter_session session;
|
||||
int has_key;
|
||||
};
|
||||
|
||||
struct uwsgi_option fastrouter_options[] = {
|
||||
@@ -56,295 +53,206 @@ struct uwsgi_option fastrouter_options[] = {
|
||||
{0, 0, 0, 0, 0, 0, 0},
|
||||
};
|
||||
|
||||
ssize_t fr_recv_uwsgi_header(struct corerouter_session *);
|
||||
ssize_t fr_instance_read_response(struct corerouter_session *);
|
||||
ssize_t fr_read_body(struct corerouter_session *);
|
||||
|
||||
void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) {
|
||||
|
||||
// here i use directly corerouter_session
|
||||
struct corerouter_session *cs = (struct corerouter_session *) data;
|
||||
struct corerouter_peer *peer = (struct corerouter_peer *) data;
|
||||
struct fastrouter_session *fr = (struct fastrouter_session *) peer->session;
|
||||
|
||||
//uwsgi_log("%.*s = %.*s\n", keylen, key, vallen, val);
|
||||
if (!uwsgi_strncmp("SERVER_NAME", 11, key, keylen) && !cs->hostname_len) {
|
||||
cs->hostname = val;
|
||||
cs->hostname_len = vallen;
|
||||
if (!uwsgi_strncmp("SERVER_NAME", 11, key, keylen) && !peer->key_len) {
|
||||
peer->key = val;
|
||||
peer->key_len = vallen;
|
||||
return;
|
||||
}
|
||||
|
||||
if (!uwsgi_strncmp("HTTP_HOST", 9, key, keylen) && !cs->has_key) {
|
||||
cs->hostname = val;
|
||||
cs->hostname_len = vallen;
|
||||
if (!uwsgi_strncmp("HTTP_HOST", 9, key, keylen) && !fr->has_key) {
|
||||
peer->key = val;
|
||||
peer->key_len = vallen;
|
||||
return;
|
||||
}
|
||||
|
||||
if (!uwsgi_strncmp("UWSGI_FASTROUTER_KEY", 20, key, keylen)) {
|
||||
cs->has_key = 1;
|
||||
cs->hostname = val;
|
||||
cs->hostname_len = vallen;
|
||||
return;
|
||||
}
|
||||
|
||||
if (!uwsgi_strncmp("CONTENT_LENGTH", 14, key, keylen)) {
|
||||
cs->post_cl = uwsgi_str_num(val, vallen);
|
||||
fr->has_key = 1;
|
||||
peer->key = val;
|
||||
peer->key_len = vallen;
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
ssize_t fr_write_body(struct corerouter_session * cs) {
|
||||
struct fastrouter_session *fs = (struct fastrouter_session *) cs;
|
||||
ssize_t len = write(cs->instance_fd, fs->post_buf->buf + fs->post_buf_pos, fs->post_buf_len - fs->post_buf_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_write_body()");
|
||||
return -1;
|
||||
}
|
||||
// writing client body to the instance
|
||||
ssize_t fr_instance_write_body(struct corerouter_peer *peer) {
|
||||
ssize_t len = cr_write(peer, "fr_instance_write_body()");
|
||||
// end on empty write
|
||||
if (!len) return 0;
|
||||
|
||||
fs->post_buf_pos += len;
|
||||
// the chunk has been sent, start (again) reading from client and instances
|
||||
if (cr_write_complete(peer)) {
|
||||
// reset the original read buffer
|
||||
peer->out->pos = 0;
|
||||
cr_reset_hooks(peer);
|
||||
}
|
||||
|
||||
// the body chunk has been sent, start again reading from client and instance
|
||||
if (fs->post_buf_pos == (ssize_t) fs->post_buf_len) {
|
||||
uwsgi_cr_hook_instance_write(cs, NULL);
|
||||
uwsgi_cr_hook_instance_read(cs, fr_instance_read_response);
|
||||
uwsgi_cr_hook_read(cs, fr_read_body);
|
||||
}
|
||||
return len;
|
||||
}
|
||||
|
||||
|
||||
// read client body
|
||||
ssize_t fr_read_body(struct corerouter_peer *main_peer) {
|
||||
ssize_t len = cr_read(main_peer, "fr_read_body()");
|
||||
if (!len) return 0;
|
||||
|
||||
main_peer->session->peers->out = main_peer->in;
|
||||
main_peer->session->peers->out_pos = 0;
|
||||
|
||||
cr_write_to_backend(main_peer->session->peers, fr_instance_write_body);
|
||||
return len;
|
||||
}
|
||||
|
||||
// write to the client
|
||||
ssize_t fr_write(struct corerouter_peer *main_peer) {
|
||||
ssize_t len = cr_write(main_peer, "fr_write()");
|
||||
// end on empty write
|
||||
if (!len) return 0;
|
||||
|
||||
// ok this response chunk is sent, let's start reading again
|
||||
if (cr_write_complete(main_peer)) {
|
||||
// reset the original read buffer
|
||||
main_peer->out->pos = 0;
|
||||
cr_reset_hooks(main_peer);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
// data from instance
|
||||
ssize_t fr_instance_read(struct corerouter_peer *peer) {
|
||||
ssize_t len = cr_read(peer, "fr_instance_read()");
|
||||
if (!len) return 0;
|
||||
|
||||
// set the input buffer as the main output one
|
||||
peer->session->main_peer->out = peer->in;
|
||||
peer->session->main_peer->out_pos = 0;
|
||||
|
||||
cr_write_to_main(peer, fr_write);
|
||||
return len;
|
||||
}
|
||||
|
||||
// send the uwsgi request header and vars
|
||||
ssize_t fr_instance_send_request(struct corerouter_peer *peer) {
|
||||
ssize_t len = cr_write(peer, "fr_instance_send_request()");
|
||||
// end on empty write
|
||||
if (!len) return 0;
|
||||
|
||||
// the chunk has been sent, start (again) reading from client and instances
|
||||
if (cr_write_complete(peer)) {
|
||||
// reset the original read buffer
|
||||
peer->out->pos = 0;
|
||||
cr_reset_hooks(peer);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
// instance is connected
|
||||
ssize_t fr_instance_connected(struct corerouter_peer *peer) {
|
||||
|
||||
ssize_t fr_read_body(struct corerouter_session * cs) {
|
||||
struct fastrouter_session *fs = (struct fastrouter_session *) cs;
|
||||
ssize_t len = read(cs->fd, fs->post_buf->buf, fs->post_buf_max);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_read_body()");
|
||||
return -1;
|
||||
}
|
||||
cr_peer_connected(peer, "fr_instance_connected()");
|
||||
|
||||
// connection closed
|
||||
if (len == 0)
|
||||
return 0;
|
||||
// fix modifiers
|
||||
peer->in->buf[0] = peer->session->main_peer->modifier1;
|
||||
peer->in->buf[3] = peer->session->main_peer->modifier2;
|
||||
|
||||
fs->post_buf_len = len;
|
||||
fs->post_buf_pos = 0;
|
||||
// prepare to write the uwsgi packet
|
||||
peer->out = peer->session->main_peer->in;
|
||||
peer->out_pos = 0;
|
||||
|
||||
// ok we have a body, stop reading from the client and the instance and start writing to the instance
|
||||
uwsgi_cr_hook_read(cs, NULL);
|
||||
uwsgi_cr_hook_instance_read(cs, NULL);
|
||||
uwsgi_cr_hook_instance_write(cs, fr_write_body);
|
||||
|
||||
return len;
|
||||
peer->last_hook_write = fr_instance_send_request;
|
||||
return fr_instance_send_request(peer);
|
||||
}
|
||||
|
||||
ssize_t fr_write_response(struct corerouter_session * cs) {
|
||||
ssize_t len = write(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_write_response()");
|
||||
return -1;
|
||||
}
|
||||
|
||||
cs->buffer_pos += len;
|
||||
|
||||
// ok this response chunk is sent, let's wait for another one
|
||||
if (cs->buffer_pos == (ssize_t) cs->buffer_len) {
|
||||
uwsgi_cr_hook_write(cs, NULL);
|
||||
uwsgi_cr_hook_instance_read(cs, fr_instance_read_response);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t fr_instance_read_response(struct corerouter_session * cs) {
|
||||
ssize_t len = read(cs->instance_fd, cs->buffer->buf, cs->buffer->len);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_instance_read_response()");
|
||||
return -1;
|
||||
}
|
||||
|
||||
// end of the response
|
||||
if (len == 0) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
cs->buffer_pos = 0;
|
||||
cs->buffer_len = len;
|
||||
// ok stop reading from the instance, and start writing to the client
|
||||
uwsgi_cr_hook_instance_read(cs, NULL);
|
||||
uwsgi_cr_hook_write(cs, fr_write_response);
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t fr_instance_send_request(struct corerouter_session * cs) {
|
||||
ssize_t len = write(cs->instance_fd, cs->buffer->buf + cs->buffer_pos, cs->uh.pktsize - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_instance_send_request()");
|
||||
return -1;
|
||||
}
|
||||
|
||||
cs->buffer_pos += len;
|
||||
|
||||
// ok the request is sent, we can start sending client body (if any) and we can start waiting
|
||||
// for response
|
||||
if (cs->buffer_pos == cs->uh.pktsize) {
|
||||
cs->buffer_pos = 0;
|
||||
// stop writing to the instance
|
||||
uwsgi_cr_hook_instance_write(cs, NULL);
|
||||
// start reading from the instance
|
||||
uwsgi_cr_hook_instance_read(cs, fr_instance_read_response);
|
||||
// re-start reading from the client (for body or connection close)
|
||||
struct fastrouter_session *fs = (struct fastrouter_session *) cs;
|
||||
// allocate a buffer for client body (could be delimited or dynamic)
|
||||
fs->post_buf_max = UMAX16;
|
||||
if (cs->post_cl > 0) {
|
||||
fs->post_buf_max = UMIN(UMAX16, cs->post_cl);
|
||||
}
|
||||
fs->post_buf = uwsgi_buffer_new(fs->post_buf_max);
|
||||
if (!fs->post_buf)
|
||||
return -1;
|
||||
uwsgi_cr_hook_read(cs, fr_read_body);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t fr_instance_send_request_header(struct corerouter_session * cs) {
|
||||
ssize_t len = write(cs->instance_fd, &cs->uh + cs->buffer_pos, 4 - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_instance_send_request_header()");
|
||||
return -1;
|
||||
}
|
||||
|
||||
cs->buffer_pos += len;
|
||||
|
||||
// ok the request is sent, we can start sending client body (if any) and we can start waiting
|
||||
// for response
|
||||
if (cs->buffer_pos == 4) {
|
||||
cs->buffer_pos = 0;
|
||||
uwsgi_cr_hook_instance_write(cs, fr_instance_send_request);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t fr_instance_connected(struct corerouter_session * cs) {
|
||||
|
||||
cs->connecting = 0;
|
||||
|
||||
socklen_t solen = sizeof(int);
|
||||
|
||||
// first check for errors
|
||||
if (getsockopt(cs->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&cs->soopt), &solen) < 0) {
|
||||
uwsgi_error("fr_instance_connected()/getsockopt()");
|
||||
cs->instance_failed = 1;
|
||||
return -1;
|
||||
}
|
||||
|
||||
if (cs->soopt) {
|
||||
cs->instance_failed = 1;
|
||||
return -1;
|
||||
}
|
||||
|
||||
cs->buffer_pos = 0;
|
||||
|
||||
// ok instance is connected, wait for write again
|
||||
if (cs->static_node) cs->static_node->custom2++;
|
||||
if (cs->un) cs->un->requests++;
|
||||
uwsgi_cr_hook_instance_write(cs, fr_instance_send_request_header);
|
||||
// return a value > 0
|
||||
return 1;
|
||||
}
|
||||
|
||||
ssize_t fr_recv_uwsgi_vars(struct corerouter_session * cs) {
|
||||
// called after receaving the uwsgi header (read vars)
|
||||
ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) {
|
||||
struct uwsgi_header *uh = (struct uwsgi_header *) main_peer->in->buf;
|
||||
// increase buffer if needed
|
||||
if (uwsgi_buffer_fix(cs->buffer, cs->uh.pktsize))
|
||||
if (uwsgi_buffer_fix(main_peer->in, uh->pktsize+4))
|
||||
return -1;
|
||||
ssize_t len = read(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->uh.pktsize - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_recv_uwsgi_vars()");
|
||||
return -1;
|
||||
}
|
||||
|
||||
cs->buffer_pos += len;
|
||||
ssize_t len = cr_read_exact(main_peer, uh->pktsize+4, "fr_recv_uwsgi_vars()");
|
||||
if (!len) return 0;
|
||||
|
||||
// headers received, ready to choose the instance
|
||||
if (cs->buffer_pos == cs->uh.pktsize) {
|
||||
struct uwsgi_corerouter *ucr = cs->corerouter;
|
||||
if (main_peer->in->pos == (size_t)(uh->pktsize+4)) {
|
||||
struct uwsgi_corerouter *ucr = main_peer->session->corerouter;
|
||||
|
||||
struct corerouter_peer *new_peer = uwsgi_cr_peer_add(main_peer->session);
|
||||
new_peer->last_hook_read = fr_instance_read;
|
||||
|
||||
// find the hostname
|
||||
if (uwsgi_hooked_parse(cs->buffer->buf, cs->uh.pktsize, fr_get_hostname, (void *) cs)) {
|
||||
if (uwsgi_hooked_parse(main_peer->in->buf+4, uh->pktsize, fr_get_hostname, (void *) new_peer)) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
// check the hostname;
|
||||
if (cs->hostname_len == 0)
|
||||
if (new_peer->key_len == 0)
|
||||
return -1;
|
||||
|
||||
// find an instance using the key
|
||||
if (cs->corerouter->mapper(cs->corerouter, cs))
|
||||
if (ucr->mapper(ucr, new_peer))
|
||||
return -1;
|
||||
|
||||
// check instance
|
||||
if (cs->instance_address_len == 0) {
|
||||
// if fallback nodes are configured, trigger them
|
||||
if (ucr->fallback) {
|
||||
cs->instance_failed = 1;
|
||||
}
|
||||
if (new_peer->instance_address_len == 0)
|
||||
return -1;
|
||||
}
|
||||
|
||||
// stop receiving from the client
|
||||
uwsgi_cr_hook_read(cs, NULL);
|
||||
|
||||
// start async connect
|
||||
cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1);
|
||||
if (cs->instance_fd < 0) {
|
||||
cs->instance_failed = 1;
|
||||
cs->soopt = errno;
|
||||
return -1;
|
||||
}
|
||||
// map the instance
|
||||
cs->corerouter->cr_table[cs->instance_fd] = cs;
|
||||
// wait for connection
|
||||
cs->connecting = 1;
|
||||
uwsgi_cr_hook_instance_write(cs, fr_instance_connected);
|
||||
cr_connect(new_peer, fr_instance_connected);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t fr_recv_uwsgi_header(struct corerouter_session * cs) {
|
||||
ssize_t len = read(cs->fd, cs->buffer->buf + cs->buffer_pos, 4 - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_recv_uwsgi_header()");
|
||||
return -1;
|
||||
}
|
||||
|
||||
cs->buffer_pos += len;
|
||||
// called soon after accept()
|
||||
ssize_t fr_recv_uwsgi_header(struct corerouter_peer *main_peer) {
|
||||
ssize_t len = cr_read_exact(main_peer, 4, "fr_recv_uwsgi_header()");
|
||||
if (!len) return 0;
|
||||
|
||||
// header ready
|
||||
if (cs->buffer_pos == 4) {
|
||||
memcpy(&cs->uh, cs->buffer->buf, 4);
|
||||
cs->buffer_pos = 0;
|
||||
uwsgi_cr_hook_read(cs, fr_recv_uwsgi_vars);
|
||||
if (main_peer->in->pos == 4) {
|
||||
// change the reading default hook
|
||||
main_peer->last_hook_read = fr_recv_uwsgi_vars;
|
||||
return fr_recv_uwsgi_vars(main_peer);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
void fr_session_close(struct corerouter_session *cs) {
|
||||
struct fastrouter_session *fr = (struct fastrouter_session *) cs;
|
||||
if (fr->post_buf) {
|
||||
uwsgi_buffer_destroy(fr->post_buf);
|
||||
}
|
||||
// retry connection to the backend
|
||||
int fr_retry(struct corerouter_peer *peer) {
|
||||
|
||||
struct uwsgi_corerouter *ucr = peer->session->corerouter;
|
||||
|
||||
if (peer->instance_address_len > 0) goto retry;
|
||||
|
||||
if (ucr->mapper(ucr, peer)) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
if (peer->instance_address_len == 0) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
retry:
|
||||
// start async connect (again)
|
||||
cr_connect(peer, fr_instance_connected);
|
||||
return 0;
|
||||
}
|
||||
|
||||
void fastrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) {
|
||||
cs->close = fr_session_close;
|
||||
// set the first hook
|
||||
uwsgi_cr_hook_read(cs, fr_recv_uwsgi_header);
|
||||
|
||||
// called when a new session is created
|
||||
int fastrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) {
|
||||
// set the retry hook
|
||||
cs->retry = fr_retry;
|
||||
// wait for requests...
|
||||
if (uwsgi_cr_set_hooks(cs->main_peer, fr_recv_uwsgi_header, NULL)) return -1;
|
||||
return 0;
|
||||
}
|
||||
|
||||
int fastrouter_init() {
|
||||
|
||||
+359
-882
File diff suppressed because it is too large
Load Diff
@@ -6,4 +6,4 @@ LIBS = []
|
||||
|
||||
REQUIRES = ['corerouter']
|
||||
|
||||
GCC_LIST = ['http']
|
||||
GCC_LIST = ['http', 'https', 'spdy3']
|
||||
|
||||
@@ -1089,6 +1089,40 @@ PyObject *py_uwsgi_connection_fd(PyObject * self, PyObject * args) {
|
||||
return PyInt_FromLong(wsgi_req->poll.fd);
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_websocket_send(PyObject * self, PyObject * args) {
|
||||
char *message = NULL;
|
||||
Py_ssize_t message_len = 0;
|
||||
|
||||
if (!PyArg_ParseTuple(args, "s#:websocket_send", &message, &message_len)) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
struct wsgi_request *wsgi_req = current_wsgi_req();
|
||||
|
||||
UWSGI_RELEASE_GIL
|
||||
ssize_t len = uwsgi_websocket_send(wsgi_req, message, message_len);
|
||||
UWSGI_GET_GIL
|
||||
if (len <= 0) {
|
||||
return PyErr_Format(PyExc_IOError, "Unable to send websocket message");
|
||||
}
|
||||
Py_INCREF(Py_None);
|
||||
return Py_None;
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_websocket_recv(PyObject * self, PyObject * args) {
|
||||
struct wsgi_request *wsgi_req = current_wsgi_req();
|
||||
UWSGI_RELEASE_GIL
|
||||
struct uwsgi_buffer *ub = uwsgi_websocket_recv(wsgi_req);
|
||||
UWSGI_GET_GIL
|
||||
if (!ub) {
|
||||
return PyErr_Format(PyExc_IOError, "Unable to receive websocket message");
|
||||
}
|
||||
|
||||
PyObject *ret = PyString_FromStringAndSize(ub->buf, ub->pos);
|
||||
uwsgi_buffer_destroy(ub);
|
||||
return ret;
|
||||
}
|
||||
|
||||
PyObject *py_uwsgi_embedded_data(PyObject * self, PyObject * args) {
|
||||
|
||||
char *name;
|
||||
@@ -3270,6 +3304,9 @@ static PyMethodDef uwsgi_advanced_methods[] = {
|
||||
{"set_user_harakiri", py_uwsgi_set_user_harakiri, METH_VARARGS, ""},
|
||||
//{"call_hook", py_uwsgi_call_hook, METH_VARARGS, ""},
|
||||
|
||||
{"websocket_recv", py_uwsgi_websocket_recv, METH_VARARGS, ""},
|
||||
{"websocket_send", py_uwsgi_websocket_send, METH_VARARGS, ""},
|
||||
|
||||
{NULL, NULL},
|
||||
};
|
||||
|
||||
|
||||
+131
-198
@@ -15,13 +15,13 @@ struct uwsgi_rawrouter {
|
||||
extern struct uwsgi_server uwsgi;
|
||||
|
||||
struct rawrouter_session {
|
||||
struct corerouter_session crs;
|
||||
struct corerouter_session session;
|
||||
|
||||
in_addr_t ip_addr;
|
||||
|
||||
// XCLIENT ADDR=xxx\r\n
|
||||
char xclient[13+INET_ADDRSTRLEN+2];
|
||||
size_t xclient_len;
|
||||
off_t xclient_pos;
|
||||
size_t xclient_remains;
|
||||
struct uwsgi_buffer *xclient;
|
||||
size_t xclient_pos;
|
||||
// placeholder for \r\n
|
||||
size_t xclient_rn;
|
||||
};
|
||||
@@ -63,95 +63,93 @@ struct uwsgi_option rawrouter_options[] = {
|
||||
{0, 0, 0, 0, 0, 0, 0},
|
||||
};
|
||||
|
||||
ssize_t rr_instance_read(struct corerouter_session *);
|
||||
ssize_t rr_read(struct corerouter_session *);
|
||||
|
||||
// write to backend
|
||||
ssize_t rr_instance_write(struct corerouter_session * cs) {
|
||||
ssize_t len = write(cs->instance_fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("fr_instance_write()");
|
||||
return -1;
|
||||
}
|
||||
ssize_t rr_instance_write(struct corerouter_peer *peer) {
|
||||
ssize_t len = cr_write(peer, "rr_instance_write()");
|
||||
// end on empty write
|
||||
if (!len) return 0;
|
||||
|
||||
cs->buffer_pos += len;
|
||||
|
||||
// the chunk has been sent, start (again) reading from client and instance
|
||||
if (cs->buffer_pos == (ssize_t) cs->buffer_len) {
|
||||
uwsgi_cr_hook_instance_write(cs, NULL);
|
||||
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
|
||||
uwsgi_cr_hook_read(cs, rr_read);
|
||||
// the chunk has been sent, start (again) reading from client and instances
|
||||
if (cr_write_complete(peer)) {
|
||||
// reset the buffer
|
||||
peer->out->pos = 0;
|
||||
cr_reset_hooks(peer);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
// write to client
|
||||
ssize_t rr_write(struct corerouter_session * cs) {
|
||||
ssize_t len = write(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("rr_write()");
|
||||
return -1;
|
||||
}
|
||||
ssize_t rr_write(struct corerouter_peer *main_peer) {
|
||||
ssize_t len = cr_write(main_peer, "rr_write()");
|
||||
// end on empty write
|
||||
if (!len) return 0;
|
||||
|
||||
cs->buffer_pos += len;
|
||||
|
||||
// ok this response chunk is sent, let's wait for another one
|
||||
if (cs->buffer_pos == (ssize_t) cs->buffer_len) {
|
||||
uwsgi_cr_hook_write(cs, NULL);
|
||||
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
|
||||
}
|
||||
// ok this response chunk is sent, let's start reading again
|
||||
if (cr_write_complete(main_peer)) {
|
||||
// reset the buffer
|
||||
main_peer->out->pos = 0;
|
||||
cr_reset_hooks(main_peer);
|
||||
}
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t rr_instance_read(struct corerouter_session * cs) {
|
||||
ssize_t len = read(cs->instance_fd, cs->buffer->buf, cs->buffer->len);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("rr_instance_read()");
|
||||
return -1;
|
||||
}
|
||||
// read from backend
|
||||
ssize_t rr_instance_read(struct corerouter_peer *peer) {
|
||||
ssize_t len = cr_read(peer, "rr_instance_read()");
|
||||
if (!len) return 0;
|
||||
|
||||
// end of the response
|
||||
if (len == 0) {
|
||||
return 0;
|
||||
}
|
||||
// set the input buffer as the main output one
|
||||
peer->session->main_peer->out = peer->in;
|
||||
peer->session->main_peer->out_pos = 0;
|
||||
|
||||
cs->buffer_pos = 0;
|
||||
cs->buffer_len = len;
|
||||
// ok stop reading from the instance, and start writing to the client
|
||||
uwsgi_cr_hook_instance_read(cs, NULL);
|
||||
uwsgi_cr_hook_write(cs, rr_write);
|
||||
cr_write_to_main(peer, rr_write);
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t rr_xclient_write(struct corerouter_session *);
|
||||
// write the xclient banner
|
||||
ssize_t rr_xclient_write(struct corerouter_peer *peer) {
|
||||
struct corerouter_session *cs = peer->session;
|
||||
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
|
||||
ssize_t len = cr_write_buf(peer, rr->xclient, "rr_xclient_write()");
|
||||
if (!len) return 0;
|
||||
|
||||
ssize_t rr_xclient_read(struct corerouter_session * cs) {
|
||||
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
|
||||
cs->buffer_len = cs->buffer->len;
|
||||
ssize_t len = read(cs->instance_fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("rr_xclient_read()");
|
||||
return -1;
|
||||
if (cr_write_complete_buf(peer, rr->xclient)) {
|
||||
if (peer->session->main_peer->out) {
|
||||
// (eventually) send previous data
|
||||
cr_write_to_main(peer, rr_write);
|
||||
}
|
||||
else {
|
||||
// reset to standard behaviour
|
||||
cr_reset_hooks(peer);
|
||||
}
|
||||
}
|
||||
if (len == 0) return 0;
|
||||
|
||||
char *ptr = cs->buffer->buf + cs->buffer_pos;
|
||||
return len;
|
||||
}
|
||||
|
||||
// read the first line from the backend and skip it
|
||||
ssize_t rr_xclient_read(struct corerouter_peer *peer) {
|
||||
struct corerouter_session *cs = peer->session;
|
||||
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
|
||||
ssize_t len = cr_read(peer, "rr_xclient_read()");
|
||||
if (!len) return 0;
|
||||
|
||||
char *ptr = (peer->in->buf + peer->in->pos) - len;
|
||||
ssize_t i;
|
||||
for(i=0;i<len;i++) {
|
||||
if (rr->xclient_rn == 1) {
|
||||
if (ptr[i] != '\n') {
|
||||
return -1;
|
||||
}
|
||||
// banner received
|
||||
cs->buffer_pos = len - (i+1);
|
||||
uwsgi_cr_hook_instance_read(cs, NULL);
|
||||
uwsgi_cr_hook_instance_write(cs, rr_xclient_write);
|
||||
// banner received (skip it, will be sent later)
|
||||
size_t remains = peer->in->pos - (len - (i+1));
|
||||
if (remains > 0) {
|
||||
peer->session->main_peer->out = peer->in;
|
||||
peer->session->main_peer->out_pos = remains;
|
||||
}
|
||||
cr_write_to_backend(peer, rr_xclient_write);
|
||||
return len;
|
||||
}
|
||||
else if (ptr[i] == '\r') {
|
||||
@@ -159,171 +157,107 @@ ssize_t rr_xclient_read(struct corerouter_session * cs) {
|
||||
}
|
||||
}
|
||||
|
||||
cs->buffer_pos += len;
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t rr_xclient_write(struct corerouter_session * cs) {
|
||||
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
|
||||
ssize_t len = write(cs->instance_fd, rr->xclient + rr->xclient_pos, rr->xclient_len - rr->xclient_pos);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("rr_xclient_write()");
|
||||
return -1;
|
||||
}
|
||||
// the instance is connected now we cannot retry connections
|
||||
ssize_t rr_instance_connected(struct corerouter_peer *peer) {
|
||||
|
||||
rr->xclient_pos += len;
|
||||
if (rr->xclient_pos == (ssize_t) rr->xclient_len) {
|
||||
uwsgi_cr_hook_instance_write(cs, NULL);
|
||||
if (cs->buffer_pos > 0) {
|
||||
// send remaining data...
|
||||
uwsgi_cr_hook_write(cs, rr_write);
|
||||
}
|
||||
else {
|
||||
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
|
||||
uwsgi_cr_hook_read(cs, rr_read);
|
||||
}
|
||||
}
|
||||
struct corerouter_session *cs = peer->session;
|
||||
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
|
||||
cr_peer_connected(peer, "rr_instance_connected()");
|
||||
|
||||
return len;
|
||||
}
|
||||
|
||||
ssize_t rr_instance_connected(struct corerouter_session * cs) {
|
||||
|
||||
cs->connecting = 0;
|
||||
|
||||
socklen_t solen = sizeof(int);
|
||||
|
||||
// first check for errors
|
||||
if (getsockopt(cs->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&cs->soopt), &solen) < 0) {
|
||||
uwsgi_error("rr_instance_connected()/getsockopt()");
|
||||
cs->instance_failed = 1;
|
||||
return -1;
|
||||
}
|
||||
|
||||
if (cs->soopt) {
|
||||
cs->instance_failed = 1;
|
||||
return -1;
|
||||
}
|
||||
|
||||
cs->buffer_pos = 0;
|
||||
|
||||
// ok instance is connected, begin...
|
||||
if (cs->static_node) cs->static_node->custom2++;
|
||||
if (cs->un) cs->un->requests++;
|
||||
|
||||
uwsgi_cr_hook_instance_write(cs, NULL);
|
||||
if (urr.xclient) {
|
||||
uwsgi_cr_hook_instance_read(cs, rr_xclient_read);
|
||||
if (rr->xclient) {
|
||||
cr_reset_hooks_and_read(peer, rr_xclient_read);
|
||||
return 1;
|
||||
}
|
||||
uwsgi_cr_hook_instance_read(cs, rr_instance_read);
|
||||
uwsgi_cr_hook_read(cs, rr_read);
|
||||
// return a value > 0
|
||||
cr_reset_hooks_and_read(peer, rr_instance_read);
|
||||
return 1;
|
||||
}
|
||||
|
||||
ssize_t rr_read(struct corerouter_session * cs) {
|
||||
ssize_t len = read(cs->fd, cs->buffer->buf, cs->buffer->len);
|
||||
if (len < 0) {
|
||||
cr_try_again;
|
||||
uwsgi_error("rr_recv()");
|
||||
return -1;
|
||||
}
|
||||
// read from client
|
||||
ssize_t rr_read(struct corerouter_peer *main_peer) {
|
||||
ssize_t len = cr_read(main_peer, "rr_read()");
|
||||
if (!len) return 0;
|
||||
|
||||
if (len == 0) return 0;
|
||||
|
||||
cs->buffer_pos = 0;
|
||||
cs->buffer_len = len;
|
||||
|
||||
uwsgi_cr_hook_read(cs, NULL);
|
||||
uwsgi_cr_hook_instance_read(cs, NULL);
|
||||
uwsgi_cr_hook_instance_write(cs, rr_instance_write);
|
||||
main_peer->session->peers->out = main_peer->in;
|
||||
main_peer->session->peers->out_pos = 0;
|
||||
|
||||
cr_write_to_backend(main_peer->session->peers, rr_instance_write);
|
||||
return len;
|
||||
}
|
||||
|
||||
int rr_retry(struct uwsgi_corerouter *ucr, struct corerouter_session *cs) {
|
||||
// retry the connection
|
||||
int rr_retry(struct corerouter_peer *peer) {
|
||||
|
||||
if (cs->instance_address_len > 0) goto retry;
|
||||
struct corerouter_session *cs = peer->session;
|
||||
struct uwsgi_corerouter *ucr = cs->corerouter;
|
||||
|
||||
if (ucr->mapper(ucr, cs)) {
|
||||
cs->instance_failed = 1;
|
||||
return -1;
|
||||
}
|
||||
if (peer->instance_address_len > 0) goto retry;
|
||||
|
||||
if (cs->instance_address_len == 0) {
|
||||
cs->instance_failed = 1;
|
||||
return -1;
|
||||
}
|
||||
if (ucr->mapper(ucr, peer)) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
if (peer->instance_address_len == 0) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
retry:
|
||||
// start async connect
|
||||
cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1);
|
||||
if (cs->instance_fd < 0) {
|
||||
cs->instance_failed = 1;
|
||||
cs->soopt = errno;
|
||||
return -1;
|
||||
}
|
||||
// map the instance
|
||||
cs->corerouter->cr_table[cs->instance_fd] = cs;
|
||||
// wait for connection
|
||||
cs->connecting = 1;
|
||||
// wait for connection
|
||||
uwsgi_cr_hook_instance_write(cs, rr_instance_connected);
|
||||
// start async connect (again)
|
||||
cr_connect(peer, rr_instance_connected);
|
||||
return 0;
|
||||
}
|
||||
|
||||
void rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) {
|
||||
void rr_session_close(struct corerouter_session *cs) {
|
||||
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
|
||||
if (rr->xclient) {
|
||||
uwsgi_buffer_destroy(rr->xclient);
|
||||
}
|
||||
}
|
||||
|
||||
// use the address as hostname
|
||||
cs->hostname = cs->ugs->name;
|
||||
cs->hostname_len = cs->ugs->name_len;
|
||||
// allocate a new session
|
||||
int rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) {
|
||||
|
||||
// set default read hook
|
||||
cs->main_peer->last_hook_read = rr_read;
|
||||
// set close hook
|
||||
cs->close = rr_session_close;
|
||||
// set retry hook
|
||||
cs->retry = rr_retry;
|
||||
|
||||
if (sa && sa->sa_family == AF_INET) {
|
||||
struct rawrouter_session *rr = (struct rawrouter_session *) cs;
|
||||
rr->ip_addr = ((struct sockaddr_in *) sa)->sin_addr.s_addr;
|
||||
if (urr.xclient) {
|
||||
if (!inet_ntop(AF_INET, &rr->ip_addr, rr->xclient+13, INET_ADDRSTRLEN)) {
|
||||
uwsgi_error("rawrouter_alloc_session() -> inet_ntop()");
|
||||
cs->instance_failed = 1;
|
||||
return;
|
||||
}
|
||||
// fix string
|
||||
size_t ip_addr_len = strlen(rr->xclient+13);
|
||||
memcpy(rr->xclient,"XCLIENT ADDR=", 13);
|
||||
rr->xclient[13+ip_addr_len] = '\r';
|
||||
rr->xclient[13+ip_addr_len+1] = '\n';
|
||||
rr->xclient_len = 13 + ip_addr_len + 2;
|
||||
rr->xclient = uwsgi_buffer_new(13+INET_ADDRSTRLEN+2);
|
||||
if (uwsgi_buffer_append(rr->xclient, "XCLIENT ADDR=", 13)) return -1;
|
||||
if (uwsgi_buffer_append_ipv4(rr->xclient, &rr->ip_addr)) return -1;
|
||||
if (uwsgi_buffer_append(rr->xclient, "\r\n", 2)) return -1;
|
||||
}
|
||||
}
|
||||
|
||||
// add a new peer
|
||||
struct corerouter_peer *peer = uwsgi_cr_peer_add(cs);
|
||||
|
||||
// set default peer hook
|
||||
peer->last_hook_read = rr_instance_read;
|
||||
|
||||
// use the address as hostname
|
||||
peer->key = cs->ugs->name;
|
||||
peer->key_len = cs->ugs->name_len;
|
||||
|
||||
// the mapper hook
|
||||
if (ucr->mapper(ucr, cs)) {
|
||||
cs->instance_failed = 1;
|
||||
return;
|
||||
}
|
||||
if (ucr->mapper(ucr, peer)) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
if (cs->instance_address_len == 0) {
|
||||
cs->instance_failed = 1;
|
||||
return;
|
||||
}
|
||||
if (peer->instance_address_len == 0) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
// ok, now we could retry
|
||||
cs->retry = rr_retry;
|
||||
|
||||
// start async connect
|
||||
cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1);
|
||||
if (cs->instance_fd < 0) {
|
||||
cs->instance_failed = 1;
|
||||
cs->soopt = errno;
|
||||
return;
|
||||
}
|
||||
// map the instance
|
||||
cs->corerouter->cr_table[cs->instance_fd] = cs;
|
||||
// wait for connection
|
||||
cs->connecting = 1;
|
||||
uwsgi_cr_hook_instance_write(cs, rr_instance_connected);
|
||||
cr_connect(peer, rr_instance_connected);
|
||||
return 0;
|
||||
}
|
||||
|
||||
int rawrouter_init() {
|
||||
@@ -340,7 +274,6 @@ void rawrouter_setup() {
|
||||
urr.cr.short_name = uwsgi_str("rawrouter");
|
||||
}
|
||||
|
||||
|
||||
struct uwsgi_plugin rawrouter_plugin = {
|
||||
|
||||
.name = "rawrouter",
|
||||
|
||||
@@ -310,7 +310,7 @@ extern int pivot_root(const char *new_root, const char *put_old);
|
||||
|
||||
struct uwsgi_buffer {
|
||||
char *buf;
|
||||
off_t pos;
|
||||
size_t pos;
|
||||
size_t len;
|
||||
size_t limit;
|
||||
};
|
||||
@@ -1192,7 +1192,18 @@ struct wsgi_request {
|
||||
struct uwsgi_logvar *logvars;
|
||||
struct uwsgi_string_list *additional_headers;
|
||||
|
||||
uint16_t stream_id;
|
||||
struct uwsgi_buffer *websocket_buf;
|
||||
size_t websocket_need;
|
||||
int websocket_phase;
|
||||
uint8_t websocket_opcode;
|
||||
size_t websocket_has_mask;
|
||||
size_t websocket_size;
|
||||
size_t websocket_pktsize;
|
||||
time_t websocket_last_ping;
|
||||
time_t websocket_last_pong;
|
||||
int websocket_closed;
|
||||
|
||||
uint64_t stream_id;
|
||||
|
||||
// avoid routing loops
|
||||
int is_routing;
|
||||
@@ -2032,6 +2043,14 @@ struct uwsgi_server {
|
||||
#endif
|
||||
#endif
|
||||
|
||||
ssize_t (*websockets_hook_send)(struct wsgi_request *, struct uwsgi_buffer *);
|
||||
ssize_t (*websockets_hook_recv)(struct wsgi_request *);
|
||||
struct uwsgi_buffer *websockets_ping;
|
||||
struct uwsgi_buffer *websockets_pong;
|
||||
int websockets_ping_freq;
|
||||
int websockets_pong_freq;
|
||||
uint64_t websockets_max_size;
|
||||
|
||||
};
|
||||
|
||||
struct uwsgi_rpc {
|
||||
@@ -2474,6 +2493,7 @@ void uwsgi_ldap_config(char *);
|
||||
#endif
|
||||
|
||||
int uwsgi_strncmp(char *, int, char *, int);
|
||||
int uwsgi_strnicmp(char *, int, char *, int);
|
||||
int uwsgi_startswith(char *, char *, int);
|
||||
|
||||
|
||||
@@ -3383,6 +3403,9 @@ char *uwsgi_rsa_sign(char *, char *, size_t, unsigned int *);
|
||||
char *uwsgi_sanitize_cert_filename(char *, char *, uint16_t);
|
||||
void uwsgi_opt_scd(char *, char *, void *);
|
||||
int uwsgi_subscription_sign_check(struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_req *);
|
||||
|
||||
char *uwsgi_sha1(char *, size_t, char *);
|
||||
char *uwsgi_sha1_2n(char *, size_t, char *, size_t, char *);
|
||||
#endif
|
||||
|
||||
void uwsgi_opt_ssa(char *, char *, void *);
|
||||
@@ -3513,10 +3536,20 @@ int uwsgi_buffer_append(struct uwsgi_buffer *, char *, size_t);
|
||||
int uwsgi_buffer_fix(struct uwsgi_buffer *, size_t);
|
||||
int uwsgi_buffer_ensure(struct uwsgi_buffer *, size_t);
|
||||
void uwsgi_buffer_destroy(struct uwsgi_buffer *);
|
||||
int uwsgi_buffer_u8(struct uwsgi_buffer *, uint8_t);
|
||||
int uwsgi_buffer_byte(struct uwsgi_buffer *, char);
|
||||
int uwsgi_buffer_u16le(struct uwsgi_buffer *, uint16_t);
|
||||
int uwsgi_buffer_u16be(struct uwsgi_buffer *, uint16_t);
|
||||
int uwsgi_buffer_u32be(struct uwsgi_buffer *, uint32_t);
|
||||
int uwsgi_buffer_u64be(struct uwsgi_buffer *, uint64_t);
|
||||
int uwsgi_buffer_num64(struct uwsgi_buffer *, int64_t);
|
||||
int uwsgi_buffer_append_keyval(struct uwsgi_buffer *, char *, uint16_t, char *, uint64_t);
|
||||
int uwsgi_buffer_append_keyval(struct uwsgi_buffer *, char *, uint16_t, char *, uint16_t);
|
||||
int uwsgi_buffer_append_keyval32(struct uwsgi_buffer *, char *, uint32_t, char *, uint32_t);
|
||||
int uwsgi_buffer_append_keynum(struct uwsgi_buffer *, char *, uint16_t, int64_t);
|
||||
int uwsgi_buffer_append_ipv4(struct uwsgi_buffer *, void *);
|
||||
int uwsgi_buffer_append_keyipv4(struct uwsgi_buffer *, char *, uint16_t, void *);
|
||||
int uwsgi_buffer_decapitate(struct uwsgi_buffer *, size_t);
|
||||
int uwsgi_buffer_append_base64(struct uwsgi_buffer *, char *, size_t);
|
||||
|
||||
void uwsgi_httpize_var(char *, size_t);
|
||||
struct uwsgi_buffer *uwsgi_to_http(struct wsgi_request *, char *, uint16_t, char *, uint16_t);
|
||||
@@ -3673,6 +3706,16 @@ char *uwsgi_base64_encode(char *, size_t, size_t *);
|
||||
void uwsgi_subscribe_all(uint8_t, int);
|
||||
#define uwsgi_unsubscribe_all() uwsgi_subscribe_all(1, 1)
|
||||
|
||||
void uwsgi_websockets_init(void);
|
||||
ssize_t uwsgi_websocket_send(struct wsgi_request *, char *, size_t);
|
||||
struct uwsgi_buffer *uwsgi_websocket_recv(struct wsgi_request *);
|
||||
ssize_t uwsgi_websockets_simple_send(struct wsgi_request *, struct uwsgi_buffer *);
|
||||
ssize_t uwsgi_websockets_simple_recv(struct wsgi_request *);
|
||||
|
||||
uint16_t uwsgi_be16(char *);
|
||||
uint32_t uwsgi_be32(char *);
|
||||
uint64_t uwsgi_be64(char *);
|
||||
|
||||
void uwsgi_check_emperor(void);
|
||||
#ifdef UWSGI_AS_SHARED_LIBRARY
|
||||
int uwsgi_init(int, char **, char **);
|
||||
|
||||
+1
-1
@@ -453,7 +453,7 @@ class uConf(object):
|
||||
self.config.read(filename)
|
||||
self.gcc_list = ['core/utils', 'core/protocol', 'core/socket', 'core/logging', 'core/master', 'core/master_utils', 'core/emperor',
|
||||
'core/notify', 'core/mule', 'core/subscription', 'core/stats', 'core/sendfile',
|
||||
'core/offload', 'core/io', 'core/static',
|
||||
'core/offload', 'core/io', 'core/static', 'core/websockets',
|
||||
'core/setup_utils', 'core/clock', 'core/init', 'core/buffer',
|
||||
'core/plugins', 'core/lock', 'core/cache', 'core/daemons',
|
||||
'core/queue', 'core/event', 'core/signal', 'core/cluster',
|
||||
|
||||
Reference in New Issue
Block a user