implemented retries and fallback on the rawrouter

This commit is contained in:
Roberto De Ioris
2012-11-10 10:25:40 +01:00
parent 8759b66206
commit 951cbf3eba
5 changed files with 90 additions and 52 deletions
+41 -50
View File
@@ -202,19 +202,14 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se
if (cr_session->soopt) {
if (!ucr->quiet)
uwsgi_log("unable to connect() to uwsgi instance \"%.*s\": %s\n", (int) cr_session->instance_address_len, cr_session->instance_address, strerror(cr_session->soopt));
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->status == COREROUTER_STATUS_CONNECTING) {
if (cr_session->connecting) {
if (!ucr->quiet)
uwsgi_log("unable to connect() to uwsgi instance \"%.*s\": timeout\n", (int) cr_session->instance_address_len, cr_session->instance_address);
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);
}
else if (cr_session->status == COREROUTER_STATUS_RESPONSE) {
uwsgi_log("timeout waiting for instance \"%.*s\"\n", (int) cr_session->instance_address_len, cr_session->instance_address);
}
*/
}
}
@@ -246,6 +241,29 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se
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;
if (ucr->fallback) {
// ok let's try with the fallback nodes
if (!cr_session->fallback) {
@@ -259,35 +277,18 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se
cr_session->instance_address = cr_session->fallback->value;
cr_session->instance_address_len = cr_session->fallback->len;
// reset error and timeout
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->pass_fd = is_unix(cr_session->instance_address, cr_session->instance_address_len);
cr_session->instance_fd = uwsgi_connectn(cr_session->instance_address, cr_session->instance_address_len, 0, 1);
if (cr_session->instance_fd < 0) {
cr_session->instance_failed = 1;
cr_session->soopt = errno;
corerouter_close_session(ucr, cr_session);
return;
if (cr_session->retry(ucr, cr_session)) {
if (!cr_session->instance_failed) goto end;
}
ucr->cr_table[cr_session->instance_fd] = cr_session;
//cr_session->status = COREROUTER_STATUS_CONNECTING;
ucr->cr_table[cr_session->instance_fd] = cr_session;
event_queue_add_fd_write(ucr->queue, cr_session->instance_fd);
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;
}
return;
}
end:
@@ -338,23 +339,10 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) {
if (urbt->key <= current) {
cr_session = (struct corerouter_session *) urbt->data;
cr_session->timed_out = 1;
if (cr_session->retry) {
cr_session->retry = 0;
/*
TODO allows retry
ucr->switch_events(ucr, cr_session, -1);
*/
if (cr_session->retry) {
cr_del_timeout(ucr, cr_session);
cr_session->timeout = cr_add_fake_timeout(ucr, cr_session);
}
else {
cr_session->timeout = corerouter_reset_timeout(ucr, cr_session);
}
}
else {
corerouter_close_session(ucr, cr_session);
if (cr_session->connecting) {
cr_session->instance_failed = 1;
}
corerouter_close_session(ucr, cr_session);
continue;
}
@@ -907,6 +895,9 @@ int uwsgi_corerouter_init(struct uwsgi_corerouter *ucr) {
if (!ucr->nevents)
ucr->nevents = 64;
if (!ucr->max_retries)
ucr->max_retries = 3;
ucr->has_backends = uwsgi_courerouter_has_has_backends(ucr);
+5 -2
View File
@@ -38,6 +38,8 @@ struct uwsgi_corerouter {
int use_cache;
int nevents;
int max_retries;
char *magic_table[256];
int queue;
@@ -102,14 +104,13 @@ struct corerouter_session {
uint16_t hostname_len;
int has_key;
int retry;
int connecting;
char *instance_address;
uint64_t instance_address_len;
struct uwsgi_subscribe_node *un;
struct uwsgi_string_list *static_node;
int pass_fd;
int soopt;
int timed_out;
@@ -139,6 +140,8 @@ 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;
struct uwsgi_buffer *buffer;
size_t buffer_len;
+3
View File
@@ -233,6 +233,8 @@ ssize_t fr_instance_send_request_header(struct corerouter_session * cs) {
ssize_t fr_instance_connected(struct corerouter_session * cs) {
cs->connecting = 0;
socklen_t solen = sizeof(int);
// first check for errors
@@ -304,6 +306,7 @@ ssize_t fr_recv_uwsgi_vars(struct corerouter_session * cs) {
// 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);
}
+3
View File
@@ -856,6 +856,8 @@ done:
}
ssize_t hr_instance_connected(struct corerouter_session * cs) {
cs->connecting = 0;
socklen_t solen = sizeof(int);
// first check for errors
@@ -1024,6 +1026,7 @@ ssize_t hs_http_manage(struct corerouter_session * cs, ssize_t len) {
// map the instance
cs->corerouter->cr_table[cs->instance_fd] = cs;
// wait for connection
cs->connecting = 1;
uwsgi_cr_hook_instance_write(cs, hr_instance_connected);
break;
}
+38
View File
@@ -42,6 +42,7 @@ struct uwsgi_option rawrouter_options[] = {
{"rawrouter-to", required_argument, 0, "forward requests to the specified uwsgi server (you can specify it multiple times for load balancing)", uwsgi_opt_add_string_list, &urr.cr.static_nodes, 0},
{"rawrouter-gracetime", required_argument, 0, "retry connections to dead static nodes after the specified amount of seconds", uwsgi_opt_set_int, &urr.cr.static_node_gracetime, 0},
{"rawrouter-events", required_argument, 0, "set the maximum number of concurrent events", uwsgi_opt_set_int, &urr.cr.nevents, 0},
{"rawrouter-max-retries", required_argument, 0, "set the maximum number of retries/fallbacks to other nodes", uwsgi_opt_set_int, &urr.cr.max_retries, 0},
{"rawrouter-quiet", required_argument, 0, "do not report failed connections to instances", uwsgi_opt_true, &urr.cr.quiet, 0},
{"rawrouter-cheap", no_argument, 0, "run the rawrouter in cheap mode", uwsgi_opt_true, &urr.cr.cheap, 0},
{"rawrouter-subscription-server", required_argument, 0, "run the rawrouter subscription server on the spcified address", uwsgi_opt_corerouter_ss, &urr, 0},
@@ -145,6 +146,8 @@ ssize_t rr_xclient_write(struct corerouter_session * cs) {
ssize_t rr_instance_connected(struct corerouter_session * cs) {
cs->connecting = 0;
socklen_t solen = sizeof(int);
// first check for errors
@@ -196,6 +199,37 @@ ssize_t rr_read(struct corerouter_session * cs) {
return len;
}
int rr_retry(struct uwsgi_corerouter *ucr, struct corerouter_session *cs) {
if (cs->instance_address_len > 0) goto retry;
if (ucr->mapper(ucr, cs)) {
cs->instance_failed = 1;
return -1;
}
if (cs->instance_address_len == 0) {
cs->instance_failed = 1;
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);
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) {
// use the address as hostname
@@ -231,6 +265,9 @@ void rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_
return;
}
// 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) {
@@ -241,6 +278,7 @@ void rawrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_
// 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);
}