diff --git a/plugins/corerouter/corerouter.c b/plugins/corerouter/corerouter.c index 81838ce4..cb2fac91 100644 --- a/plugins/corerouter/corerouter.c +++ b/plugins/corerouter/corerouter.c @@ -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); diff --git a/plugins/corerouter/cr.h b/plugins/corerouter/cr.h index 161a20bc..ab28e291 100644 --- a/plugins/corerouter/cr.h +++ b/plugins/corerouter/cr.h @@ -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; diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c index c462001d..a5865d85 100644 --- a/plugins/fastrouter/fastrouter.c +++ b/plugins/fastrouter/fastrouter.c @@ -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); } diff --git a/plugins/http/http.c b/plugins/http/http.c index bcbfd55a..ff6b2b34 100644 --- a/plugins/http/http.c +++ b/plugins/http/http.c @@ -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; } diff --git a/plugins/rawrouter/rawrouter.c b/plugins/rawrouter/rawrouter.c index ef0932a3..b02417a2 100644 --- a/plugins/rawrouter/rawrouter.c +++ b/plugins/rawrouter/rawrouter.c @@ -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); }