Compare commits

...
27 Commits
Author SHA1 Message Date
Unbit 6829d51f75 added fastrouter-max-retries 2015-10-09 07:46:19 +02:00
Roberto De Ioris 753c7f7343 added subscription-clear-on-shutdown 2015-09-09 11:22:06 +02:00
Roberto De Ioris 443fb3bed5 disable reading body in deferred connects 2015-09-09 09:18:35 +02:00
Roberto De Ioris 6f8cf4a06b allow updating socket address in vassal based subscriptions 2015-09-01 07:12:17 +02:00
Roberto De Ioris 1af7089613 add a note about multiple emperor subscriptions 2015-08-26 07:11:02 +02:00
Roberto De Ioris 82efdff518 added --subscription-vassal-required 2015-08-26 06:57:58 +02:00
Roberto De Ioris 4fe69c85ef improved vassal subscriptions 2015-08-25 10:45:31 +02:00
Roberto De Ioris 23c6c1af0e added fastrouter-defer-connect-timeout 2015-08-24 20:04:06 +02:00
Roberto De Ioris 4b60ecf2f2 prepare or configurable deferred connect 2015-08-24 12:10:06 +02:00
Roberto De Ioris b975026e37 fixed udp addressing 2015-08-21 20:06:04 +02:00
Roberto De Ioris b406a523c2 added --emperor-wait-for-command-ignore 2015-08-21 19:51:52 +02:00
Roberto De Ioris 3468d5b2a7 added support for multiple path elements in subscription mountpoints 2015-08-20 22:06:19 +02:00
Roberto De Ioris 77638f698d added support for subscription mountpoints in fastrouter 2015-08-20 21:28:17 +02:00
Roberto De Ioris c2441a0448 Merge branch 'intellisurvey' of https://github.com/unbit/uwsgi into intellisurvey 2015-08-18 19:48:28 +02:00
Roberto De Ioris fed2b98031 mark 2.1 as -dev 2015-08-18 19:48:13 +02:00
Unbit 110a09cf06 mark node as un-failed during fastrouter retry 2015-08-17 16:17:42 +02:00
Unbit ec937e1224 fixed sendto in fastrouter emperor socket 2015-08-17 15:33:09 +02:00
Unbit f36e66902a added --fastrouter-emperor-socket 2015-08-17 15:14:58 +02:00
Unbit cb5d4282c1 implemented corerouter spawn vassal 2015-08-17 15:12:03 +02:00
Roberto De Ioris 1df19fc1e5 fixed tiemout management 2015-08-11 17:51:31 +02:00
Roberto De Ioris c41d81cbe4 implemented deferred connect for the fastrouter 2015-08-11 16:46:51 +02:00
Unbit 57ded87fb6 improved inactive subscriptions 2015-08-04 10:15:49 +02:00
Unbit 525602d5e6 add a note about emperor subscription packet 2015-08-04 09:53:00 +02:00
Unbit 563c226466 added spawn,respawn,stop and log emperor commands 2015-08-04 09:46:39 +02:00
Unbit cf07022629 --emperor-command-socket and --emperor-wait-for-command 2015-08-04 08:37:52 +02:00
Unbit 1f288e3d99 added inactive slot concept for subscription system 2015-08-04 08:24:25 +02:00
Roberto De Ioris a4378cf484 added vassal field in subscription system 2015-08-04 07:37:21 +02:00
10 changed files with 456 additions and 24 deletions
+177 -3
View File
@@ -8,7 +8,8 @@ The uWSGI Emperor
extern struct uwsgi_server uwsgi;
extern char **environ;
void emperor_send_stats(int);
static void emperor_send_stats(int);
static void emperor_manage_command(int);
time_t emperor_throttle;
int emperor_throttle_level;
@@ -995,6 +996,11 @@ void emperor_respawn(struct uwsgi_instance *c_ui, time_t mod) {
return;
}
// check if we are in suspended mode
if (c_ui->pid == -1 && c_ui->suspended) {
return;
}
// check if we are in on_demand mode (the respawn will be ignored)
if (c_ui->pid == -1 && c_ui->on_demand_fd > -1) {
c_ui->last_mod = mod;
@@ -1135,6 +1141,7 @@ void emperor_add_with_attrs(struct uwsgi_emperor_scanner *ues, char *name, time_
// start without loyalty
n_ui->last_loyal = 0;
n_ui->loyal = 0;
n_ui->suspended = 0;
n_ui->attrs = attrs;
@@ -1152,6 +1159,15 @@ void emperor_add_with_attrs(struct uwsgi_emperor_scanner *ues, char *name, time_
n_ui->pipe_config[0] = -1;
n_ui->pipe_config[1] = -1;
// Check if the Emperor has to wait for a command before spawning a vassal
if (uwsgi.emperor_command_socket) {
if (uwsgi.emperor_wait_for_command && !uwsgi_string_list_has_item(uwsgi.emperor_wait_for_command_ignore, name, strlen(name))) {
n_ui->suspended = 1;
uwsgi_log("[uwsgi-emperor] %s -> \"wait-for-command\" instance detected, waiting for the spawn command ...\n", name);
return;
}
}
// ok here we check if we need to bind to the specified socket or continue with the activation
if (socket_name) {
n_ui->on_demand_fd = on_demand_bind(socket_name);
@@ -1314,7 +1330,7 @@ int uwsgi_emperor_vassal_start(struct uwsgi_instance *n_ui) {
if (uwsgi_hooks_run_and_return(uwsgi.hook_as_emperor_before_vassal, "as-emperor-before-vassal", NULL, 0)) {
emperor_del(n_ui);
return -1;
}
if (uwsgi.zeus) {
@@ -1444,6 +1460,8 @@ int uwsgi_emperor_vassal_start(struct uwsgi_instance *n_ui) {
#if defined(__linux__) && !defined(OBSOLETE_LINUX_KERNEL)
uwsgi_hooks_setns_run(uwsgi.hook_as_emperor_setns, n_ui->pid, n_ui->uid, n_ui->gid);
#endif
// ensure the instance is no more suspended
n_ui->suspended = 0;
return 0;
}
else {
@@ -2036,6 +2054,19 @@ void emperor_loop() {
uwsgi_log("*** Emperor trigger socket enabled on %s fd: %d ***\n", uwsgi.emperor_trigger_socket, uwsgi.emperor_trigger_socket_fd);
}
if (uwsgi.emperor_command_socket) {
char *udp_port = strchr(uwsgi.emperor_command_socket, ':');
if (udp_port) {
uwsgi.emperor_command_socket_fd = bind_to_udp(uwsgi.emperor_command_socket, 0, 0);
}
else {
uwsgi.emperor_command_socket_fd = bind_to_unix_dgram(uwsgi.emperor_command_socket);
}
event_queue_add_fd_read(uwsgi.emperor_queue, uwsgi.emperor_command_socket_fd);
uwsgi_log("*** Emperor command socket enabled on %s fd: %d ***\n", uwsgi.emperor_command_socket, uwsgi.emperor_command_socket_fd);
}
ui = &ui_base;
int freq = 0;
@@ -2101,6 +2132,11 @@ void emperor_loop() {
continue;
}
if (uwsgi.emperor_command_socket && uwsgi.emperor_command_socket_fd > -1 && interesting_fd == uwsgi.emperor_command_socket_fd) {
emperor_manage_command(uwsgi.emperor_command_socket_fd);
continue;
}
// check if a monitor is mapped to that file descriptor
if (uwsgi_emperor_scanner_event(interesting_fd)) {
continue;
@@ -2338,7 +2374,142 @@ recheck:
}
void emperor_send_stats(int fd) {
// manage commands sentto the dgram command socket
struct emperor_command {
char *cmd;
uint16_t cmd_len;
char *vassal;
uint16_t vassal_len;
char *msg;
uint16_t msg_len;
};
static void emperor_manage_command_parser(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) {
struct emperor_command *ec = (struct emperor_command *) data;
if (!uwsgi_strncmp("cmd", 3, key, keylen)) {
ec->cmd = val;
ec->cmd_len = vallen;
}
else if (!uwsgi_strncmp("command", 7, key, keylen)) {
ec->cmd = val;
ec->cmd_len = vallen;
}
else if (!uwsgi_strncmp("vassal", 6, key, keylen)) {
ec->vassal = val;
ec->vassal_len = vallen;
}
else if (!uwsgi_strncmp("msg", 3, key, keylen)) {
ec->msg = val;
ec->msg_len = vallen;
}
}
static void emperor_command_spawn(struct emperor_command *ec) {
if (ec->vassal_len == 0) {
uwsgi_log("[uwsgi-emperor] ERROR: the \"spawn\" command requires a vassal name\n");
return;
}
char *name = uwsgi_concat2n(ec->vassal, ec->vassal_len, "", 0);
struct uwsgi_instance *vassal = emperor_get(name);
if (!vassal) {
uwsgi_log("[uwsgi-emperor] ERROR: unable to find vassal \"%s\"\n", name);
free(name);
return;
}
if (!vassal->suspended) {
uwsgi_log("[uwsgi-emperor] WARNING: the vassal \"%s\" is already active\n", name);
free(name);
return;
}
uwsgi_log("[uwsgi-emperor] INFO: requested the spawn of vassal \"%s\"\n", name);
free(name);
// start the vassal
if (uwsgi_emperor_vassal_start(vassal)) {
// clear the vassal
emperor_del(vassal);
}
}
static void emperor_command_respawn(struct emperor_command *ec) {
if (ec->vassal_len == 0) {
uwsgi_log("[uwsgi-emperor] ERROR: the \"respawn\" command requires a vassal name\n");
return;
}
char *name = uwsgi_concat2n(ec->vassal, ec->vassal_len, "", 0);
struct uwsgi_instance *vassal = emperor_get(name);
if (!vassal) {
uwsgi_log("[uwsgi-emperor] ERROR: unable to find vassal \"%s\"\n", name);
free(name);
return;
}
uwsgi_log("[uwsgi-emperor] INFO: requested the respawn of vassal \"%s\"\n", name);
free(name);
// restart the vassal
emperor_respawn(vassal, uwsgi_now());
}
static void emperor_command_stop(struct emperor_command *ec) {
if (ec->vassal_len == 0) {
uwsgi_log("[uwsgi-emperor] ERROR: the \"stop\" command requires a vassal name\n");
return;
}
char *name = uwsgi_concat2n(ec->vassal, ec->vassal_len, "", 0);
struct uwsgi_instance *vassal = emperor_get(name);
if (!vassal) {
uwsgi_log("[uwsgi-emperor] ERROR: unable to find vassal \"%s\"\n", name);
free(name);
return;
}
uwsgi_log("[uwsgi-emperor] INFO: requested the stop of vassal \"%s\"\n", name);
free(name);
// stop the vassal
emperor_stop(vassal);
}
static void emperor_command_log(struct emperor_command *ec) {
if (ec->msg_len == 0) {
uwsgi_log("[uwsgi-emperor] ERROR: the \"log\" command requires a message\n");
return;
}
uwsgi_log("%.*s\n", ec->msg_len, ec->msg);
}
static void emperor_manage_command(int fd) {
char buf[4096];
ssize_t len = recv(fd, buf, 4096, 0);
if (len <= 0) return;
struct emperor_command ec;
memset(&ec, 0, sizeof(struct emperor_command));
uwsgi_hooked_parse(buf + 4, len - 4, emperor_manage_command_parser, &ec);
if (ec.cmd_len == 0) return;
// spawn a suspended vassal
if (!uwsgi_strncmp(ec.cmd, ec.cmd_len, "spawn", 5)) {
emperor_command_spawn(&ec);
return;
}
// respawn a vassal
if (!uwsgi_strncmp(ec.cmd, ec.cmd_len, "respawn", 7)) {
emperor_command_respawn(&ec);
return;
}
// stop a vassal
if (!uwsgi_strncmp(ec.cmd, ec.cmd_len, "stop", 4)) {
emperor_command_stop(&ec);
return;
}
// log a message
if (!uwsgi_strncmp(ec.cmd, ec.cmd_len, "log", 3)) {
emperor_command_log(&ec);
return;
}
}
static void emperor_send_stats(int fd) {
struct sockaddr_un client_src;
socklen_t client_src_len = 0;
@@ -2444,6 +2615,9 @@ void emperor_send_stats(int fd) {
if (uwsgi_stats_keyval_comma(us, "on_demand", c_ui->socket_name ? c_ui->socket_name : ""))
goto end0;
if (uwsgi_stats_keylong_comma(us, "suspended", (unsigned long long) c_ui->suspended))
goto end0;
if (uwsgi_stats_keylong_comma(us, "adopted", (unsigned long long) c_ui->adopted))
goto end0;
+79 -9
View File
@@ -167,8 +167,14 @@ struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slo
while (node) {
// is the node alive ?
if (now - node->last_check > uwsgi.subscription_tolerance) {
if (node->death_mark == 0)
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => marking %.*s as failed (no announce received in %d seconds)\n", (int) uwsgi.mypid, (int) keylen, key, (int) node->len, node->name, uwsgi.subscription_tolerance);
if (node->death_mark == 0) {
if (node->len > 0) {
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => marking %.*s as failed (no announce received in %d seconds)\n", (int) uwsgi.mypid, (int) keylen, key, (int) node->len, node->name, uwsgi.subscription_tolerance);
}
else if (node->vassal_len > 0) {
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => marking vassal %.*s as failed (no announce received in %d seconds)\n", (int) uwsgi.mypid, (int) keylen, key, (int) node->vassal_len, node->vassal, uwsgi.subscription_tolerance);
}
}
node->failcnt++;
node->death_mark = 1;
}
@@ -301,7 +307,10 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
struct uwsgi_subscribe_slot *current_slot = uwsgi_get_subscribe_slot(slot, usr->key, usr->keylen), *old_slot = NULL, *a_slot;
struct uwsgi_subscribe_node *node, *old_node = NULL;
if (usr->address_len > 0xff || usr->address_len == 0)
if ((usr->address_len > 0xff || usr->address_len == 0) && (usr->vassal_len > 0xff || usr->vassal_len == 0))
return NULL;
if (uwsgi.subscription_vassal_required && usr->vassal_len == 0)
return NULL;
if (current_slot) {
@@ -315,18 +324,39 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
return NULL;
}
int has_address_and_vassal = 0;
if (usr->address_len > 0 && usr->vassal_len > 0)
has_address_and_vassal = 1;
node = current_slot->nodes;
while (node) {
if (!uwsgi_strncmp(node->name, node->len, usr->address, usr->address_len)) {
if ((usr->address_len > 0 && !uwsgi_strncmp(node->name, node->len, usr->address, usr->address_len))
|| (usr->vassal_len > 0 && !uwsgi_strncmp(node->vassal, node->vassal_len, usr->vassal, usr->vassal_len))) {
#ifdef UWSGI_SSL
// this should avoid sending sniffed packets...
if (current_slot->sign_ctx && !subscription_is_safe(usr) && usr->unix_check <= node->unix_check) {
uwsgi_log("[uwsgi-subscription for pid %d] invalid (sniffed ?) packet sent for slot: %.*s node: %.*s unix_check: %lu\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, (unsigned long) usr->unix_check);
uwsgi_log("[uwsgi-subscription for pid %d] invalid (sniffed ?) packet sent for slot: %.*s node: %.*s unix_check: %lu\n", (int) uwsgi.mypid, usr->keylen, usr->key, (int) usr->address_len, usr->address, (unsigned long) usr->unix_check);
return NULL;
}
// eventually the packet could be upgraded to sni...
uwsgi_subscription_sni_check(current_slot, usr);
#endif
// only for vassal mode
if (has_address_and_vassal) {
if (usr->address_len == node->len && !memcmp(usr->address, node->name, node->len)) {
// record already exists, clear it ?
if (usr->clear) {
node->len = 0;
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => cleared address for vassal node: %.*s (weight: %d, backup: %d)\n", (int) uwsgi.mypid, usr->keylen, usr->key, (int) usr->vassal_len, usr->vassal, usr->weight, usr->backup_level);
}
}
else {
memcpy(node->name, usr->address, usr->address_len);
node->len = usr->address_len;
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => updated vassal node: %.*s with address %.*s (weight: %d, backup: %d)\n", (int) uwsgi.mypid, usr->keylen, usr->key, (int) usr->vassal_len, usr->vassal, (int) usr->address_len, usr->address, usr->weight, usr->backup_level);
}
}
// remove death mark and update cores and load
node->death_mark = 0;
node->last_check = uwsgi_now();
@@ -387,7 +417,13 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
}
node->last_check = uwsgi_now();
node->slot = current_slot;
memcpy(node->name, usr->address, usr->address_len);
node->vassal_len = usr->vassal_len;
if (node->len > 0)
memcpy(node->name, usr->address, node->len);
if (usr->vassal_len > 0)
memcpy(node->vassal, usr->vassal, node->vassal_len);
if (old_node) {
old_node->next = node;
}
@@ -460,7 +496,11 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
memcpy(current_slot->nodes->notify, usr->notify, usr->notify_len);
current_slot->nodes->notify[usr->notify_len] = 0;
}
memcpy(current_slot->nodes->name, usr->address, usr->address_len);
if (usr->address_len > 0)
memcpy(current_slot->nodes->name, usr->address, usr->address_len);
current_slot->nodes->vassal_len = usr->vassal_len;
if (current_slot->nodes->vassal_len > 0)
memcpy(current_slot->nodes->vassal, usr->vassal, usr->vassal_len);
current_slot->nodes->last_check = uwsgi_now();
current_slot->nodes->next = NULL;
@@ -482,13 +522,17 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo
current_slot->algo = usr->algo;
if (!current_slot->algo) current_slot->algo = uwsgi.subscription_algo;
if (!slot[hash_key] || current_slot->prev == NULL) {
slot[hash_key] = current_slot;
}
uwsgi_log("[uwsgi-subscription for pid %d] new pool: %.*s (hash key: %d, algo: %s)\n", (int) uwsgi.mypid, usr->keylen, usr->key, current_slot->hash, uwsgi_subscription_algo_name(current_slot->algo));
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => new node: %.*s (weight: %d, backup: %d)\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, usr->weight, usr->backup_level);
if (usr->address_len > 0) {
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => new node: %.*s (weight: %d, backup: %d)\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->address_len, usr->address, usr->weight, usr->backup_level);
}
else {
uwsgi_log("[uwsgi-subscription for pid %d] %.*s => new vassal node: %.*s (weight: %d, backup: %d)\n", (int) uwsgi.mypid, usr->keylen, usr->key, usr->vassal_len, usr->vassal, usr->weight, usr->backup_level);
}
if (current_slot->nodes->notify[0]) {
char buf[1024];
@@ -931,6 +975,8 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) {
char *s2_proto = NULL;
char *s2_algo = NULL;
char *s2_backup = NULL;
char *s2_vassal = NULL;
char *s2_inactive = NULL;
struct uwsgi_buffer *ub = NULL;
if (uwsgi_kvlist_parse(arg, strlen(arg), ',', '=',
@@ -938,6 +984,7 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) {
"key", &s2_key,
"socket", &s2_socket,
"addr", &s2_addr,
"address", &s2_addr,
"weight", &s2_weight,
"modifier1", &s2_modifier1,
"modifier2", &s2_modifier2,
@@ -949,6 +996,8 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) {
"proto", &s2_proto,
"algo", &s2_algo,
"backup", &s2_backup,
"vassal", &s2_vassal,
"inactive", &s2_inactive,
NULL)) {
return;
}
@@ -1021,6 +1070,11 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) {
if (uwsgi_buffer_append_keynum(ub, "backup", 6, backup))
goto end;
if (s2_vassal) {
if (uwsgi_buffer_append_keyval(ub, "vassal", 6, s2_vassal, strlen(s2_vassal)))
goto end;
}
if (s2_sni_key) {
if (uwsgi_buffer_append_keyval(ub, "sni_key", 7, s2_sni_key, strlen(s2_sni_key)))
goto end;
@@ -1046,6 +1100,11 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) {
goto end;
}
if (s2_inactive) {
if (uwsgi_buffer_append_keyval(ub, "inactive", 8, s2_inactive, strlen(s2_inactive)))
goto end;
}
if (uwsgi.subscription_notify_socket) {
if (uwsgi_buffer_append_keyval(ub, "notify", 6, uwsgi.subscription_notify_socket, strlen(uwsgi.subscription_notify_socket)))
goto end;
@@ -1055,6 +1114,13 @@ void uwsgi_subscribe2(char *arg, uint8_t cmd) {
goto end;
}
// clear instead of unsubscribe
if (uwsgi_instance_is_dying && cmd == 1 && uwsgi.subscription_clear_on_shutdown) {
if (uwsgi_buffer_append_keynum(ub, "clear", 5, 1))
goto end;
cmd = 0;
}
if (uwsgi_subscription_ub_fix(ub, modifier1, modifier2, cmd, s2_sign)) goto end;
send_subscription(-1, s2_server, ub->buf, ub->pos);
@@ -1091,8 +1157,12 @@ end:
free(s2_proto);
if (s2_algo)
free(s2_algo);
if (s2_inactive)
free(s2_inactive);
if (s2_backup)
free(s2_backup);
if (s2_vassal)
free(s2_vassal);
}
void uwsgi_subscribe_all(uint8_t cmd, int verbose) {
+7 -2
View File
@@ -218,6 +218,9 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"emperor-stats", required_argument, 0, "run the Emperor stats server", uwsgi_opt_set_str, &uwsgi.emperor_stats, 0},
{"emperor-stats-server", required_argument, 0, "run the Emperor stats server", uwsgi_opt_set_str, &uwsgi.emperor_stats, 0},
{"emperor-trigger-socket", required_argument, 0, "enable the Emperor trigger socket", uwsgi_opt_set_str, &uwsgi.emperor_trigger_socket, 0},
{"emperor-command-socket", required_argument, 0, "enable the Emperor command socket", uwsgi_opt_set_str, &uwsgi.emperor_command_socket, 0},
{"emperor-wait-for-command", no_argument, 0, "always wait for a 'spawn' Emperor command before starting a vassal", uwsgi_opt_true, &uwsgi.emperor_wait_for_command, 0},
{"emperor-wait-for-command-ignore", required_argument, 0, "ignore the emperor-wait-for-command directive for the specified vassal", uwsgi_opt_add_string_list, &uwsgi.emperor_wait_for_command_ignore, 0},
{"early-emperor", no_argument, 0, "spawn the emperor as soon as possibile", uwsgi_opt_true, &uwsgi.early_emperor, 0},
{"emperor-broodlord", required_argument, 0, "run the emperor in BroodLord mode", uwsgi_opt_set_int, &uwsgi.emperor_broodlord, 0},
{"emperor-throttle", required_argument, 0, "set throttling level (in milliseconds) for bad behaving vassals (default 1000)", uwsgi_opt_set_int, &uwsgi.emperor_throttle, 0},
@@ -630,8 +633,9 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"notify-socket", required_argument, 0, "enable the notification socket", uwsgi_opt_set_str, &uwsgi.notify_socket, UWSGI_OPT_MASTER},
{"subscription-notify-socket", required_argument, 0, "set the notification socket for subscriptions", uwsgi_opt_set_str, &uwsgi.subscription_notify_socket, UWSGI_OPT_MASTER},
{"subscription-mountpoints", no_argument, 0, "enable mountpoints support for subscription system", uwsgi_opt_true, &uwsgi.subscription_mountpoints, UWSGI_OPT_MASTER},
{"subscription-mountpoint", no_argument, 0, "enable mountpoints support for subscription system", uwsgi_opt_true, &uwsgi.subscription_mountpoints, UWSGI_OPT_MASTER},
{"subscription-mountpoints", required_argument, 0, "enable mountpoints support for subscription system", uwsgi_opt_set_int, &uwsgi.subscription_mountpoints, UWSGI_OPT_MASTER},
{"subscription-mountpoint", required_argument, 0, "enable mountpoints support for subscription system", uwsgi_opt_set_int, &uwsgi.subscription_mountpoints, UWSGI_OPT_MASTER},
{"subscription-vassal-required", no_argument, 0, "require a vassal field for each subscription packet", uwsgi_opt_true, &uwsgi.subscription_vassal_required, UWSGI_OPT_MASTER},
#ifdef UWSGI_SSL
{"legion", required_argument, 0, "became a member of a legion", uwsgi_opt_legion, NULL, UWSGI_OPT_MASTER},
@@ -668,6 +672,7 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"subscription-tolerance", required_argument, 0, "set tolerance for subscription servers", uwsgi_opt_set_int, &uwsgi.subscription_tolerance, 0},
{"unsubscribe-on-graceful-reload", no_argument, 0, "force unsubscribe request even during graceful reload", uwsgi_opt_true, &uwsgi.unsubscribe_on_graceful_reload, 0},
{"start-unsubscribed", no_argument, 0, "configure subscriptions but do not send them (useful with master fifo)", uwsgi_opt_true, &uwsgi.subscriptions_blocked, 0},
{"subscription-clear-on-shutdown", no_argument, 0, "force clear instead of unsubscribe during shutdown", uwsgi_opt_true, &uwsgi.subscription_clear_on_shutdown, 0},
{"subscribe-with-modifier1", required_argument, 0, "force the specififed modifier1 when subscribing", uwsgi_opt_set_str, &uwsgi.subscribe_with_modifier1, UWSGI_OPT_MASTER},
+52 -4
View File
@@ -316,22 +316,30 @@ void corerouter_manage_subscription(char *key, uint16_t keylen, char *val, uint1
usr->proto = val;
usr->proto_len = vallen;
}
else if (!uwsgi_strncmp("vassal", 6, key, keylen)) {
usr->vassal = val;
usr->vassal_len = vallen;
}
else if (!uwsgi_strncmp("clear", 5, key, keylen)) {
usr->clear = uwsgi_str_num(val, vallen);
}
}
void corerouter_close_peer(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) {
struct corerouter_session *cs = peer->session;
// 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", peer->un->len, peer->un->name, peer->un->reference);
uwsgi_log("[1] node %.*s refcnt: %llu\n", peer->un->len, peer->un->name, peer->un->reference);
#endif
peer->un->reference--;
peer->un->reference--;
#ifdef UWSGI_DEBUG
uwsgi_log("[2] node %.*s refcnt: %llu\n", peer->un->len, peer->un->name, peer->un->reference);
uwsgi_log("[2] node %.*s refcnt: %llu\n", peer->un->len, peer->un->name, peer->un->reference);
#endif
}
if (peer->failed) {
@@ -478,6 +486,24 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr, time_t now)
if (urbt->value <= current) {
peer = (struct corerouter_peer *) urbt->data;
// you can manage deferred connections upto X retry times
if (peer->defer_connect) {
peer->defer_connect = 0;
peer->retries++;
// ignore return value
if (peer->un) {
if (peer->un->reference == 0) {
uwsgi_log("[BUG] subscription reference counting is 0 !!!\n");
corerouter_close_peer(ucr, peer);
continue;
}
peer->un->reference--;
}
peer->session->retry(peer);
// increase timeout;
urbt->value += peer->current_timeout;
continue;
}
peer->timed_out = 1;
if (peer->connecting) {
peer->failed = 1;
@@ -488,6 +514,7 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr, time_t now)
break;
}
}
int uwsgi_cr_set_hooks(struct corerouter_peer *peer, ssize_t (*read_hook)(struct corerouter_peer *), ssize_t (*write_hook)(struct corerouter_peer *)) {
@@ -680,6 +707,9 @@ void uwsgi_corerouter_loop(int id, void *data) {
if (!ucr->socket_timeout)
ucr->socket_timeout = 60;
if (!ucr->defer_connect_timeout)
ucr->defer_connect_timeout = 5;
if (!ucr->static_node_gracetime)
ucr->static_node_gracetime = 30;
@@ -708,6 +738,23 @@ void uwsgi_corerouter_loop(int id, void *data) {
uwsgi_log("*** %s stats server enabled on %s fd: %d ***\n", ucr->short_name, ucr->stats_server, ucr->cr_stats_server);
}
if (ucr->emperor_socket) {
char *colon = strchr(ucr->emperor_socket, ':');
if (colon) {
ucr->emperor_socket_fd = socket(AF_INET, SOCK_DGRAM, 0);
ucr->emperor_socket_addr_len = socket_to_in_addr(ucr->emperor_socket, colon, 0, &ucr->emperor_socket_addr.sa_in);
}
else {
ucr->emperor_socket_fd = socket(AF_UNIX, SOCK_DGRAM, 0);
ucr->emperor_socket_addr_len = socket_to_un_addr(ucr->emperor_socket, &ucr->emperor_socket_addr.sa_un);
}
if (ucr->emperor_socket_fd < 0) {
uwsgi_error("error creating emperor socket client: socket()");
exit(1);
}
uwsgi_log("emperor socket mapped to: %s\n", ucr->emperor_socket);
}
if (ucr->use_socket) {
ucr->to_socket = uwsgi_get_socket_by_num(ucr->socket_num);
@@ -1086,6 +1133,7 @@ void corerouter_send_stats(struct uwsgi_corerouter *ucr) {
if (uwsgi_stats_object_open(us)) goto end0;
if (uwsgi_stats_keyvaln_comma(us, "name", s_node->name, s_node->len)) goto end0;
if (uwsgi_stats_keyvaln_comma(us, "vassal", s_node->vassal, s_node->vassal_len)) goto end0;
if (uwsgi_stats_keylong_comma(us, "modifier1", (unsigned long long) s_node->modifier1)) goto end0;
if (uwsgi_stats_keylong_comma(us, "modifier2", (unsigned long long) s_node->modifier2)) goto end0;
+14
View File
@@ -199,6 +199,11 @@ struct corerouter_peer {
int is_buffering;
int buffering_fd;
int defer_connect;
char *vassal;
uint8_t vassal_len;
};
struct uwsgi_corerouter {
@@ -281,6 +286,13 @@ struct uwsgi_corerouter {
size_t buffer_size;
int fallback_on_no_key;
char *emperor_socket;
int emperor_socket_fd;
union uwsgi_sockaddr emperor_socket_addr;
socklen_t emperor_socket_addr_len;
int defer_connect_timeout;
};
// a session is started when a client connect to the router
@@ -364,3 +376,5 @@ 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 *);
struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *, struct corerouter_peer *);
int corerouter_spawn_vassal(struct uwsgi_corerouter *, struct uwsgi_subscribe_node *, int);
+25
View File
@@ -277,3 +277,28 @@ void uwsgi_corerouter_manage_internal_subscription(struct uwsgi_corerouter *ucr,
}
}
int corerouter_spawn_vassal(struct uwsgi_corerouter *ucr, struct uwsgi_subscribe_node *node, int attempt) {
// TODO
// a vassal field could be in the form address:name where address is the emperor command socket
// to use. As we cannot know in advance if it will be a UDP or UNIX address, we need to
// re-create and close the socket every time.
if (!ucr->emperor_socket)
return -1;
int ret = -1;
uwsgi_log_verbose("spawning vassal %.*s (attempt %d)\n", node->vassal_len, node->vassal, attempt);
struct uwsgi_buffer *ub = uwsgi_buffer_new(uwsgi.page_size);
// leave space for uwsgi header
ub->pos += 4;
if (uwsgi_buffer_append_keyval(ub, "cmd", 3, "spawn", 5)) goto end;
if (uwsgi_buffer_append_keyval(ub, "vassal", 6, node->vassal, node->vassal_len)) goto end;
// TODO choose modifier1 and 2
if (uwsgi_buffer_set_uh(ub, 0, 0)) goto end;
if (sendto(ucr->emperor_socket_fd, ub->buf, ub->pos, 0, &ucr->emperor_socket_addr.sa, ucr->emperor_socket_addr_len) < 0) {
uwsgi_error("corerouter_spawn_vassal()/sendto()");
}
ret = 0;
end:
uwsgi_buffer_destroy(ub);
return ret;
}
+12 -3
View File
@@ -54,12 +54,21 @@ int uwsgi_cr_map_use_subscription(struct uwsgi_corerouter *ucr, struct coreroute
usc.cookie = NULL;
peer->un = uwsgi_get_subscribe_node(ucr->subscriptions, peer->key, peer->key_len, &usc);
if (peer->un && peer->un->len) {
peer->instance_address = peer->un->name;
peer->instance_address_len = peer->un->len;
// check if the node is ready or it requires a vassal spawn
if (peer->un && (peer->un->len || peer->un->vassal_len)) {
peer->modifier1 = peer->un->modifier1;
peer->modifier2 = peer->un->modifier2;
peer->proto = peer->un->proto;
if (peer->un->len) {
peer->instance_address = peer->un->name;
peer->instance_address_len = peer->un->len;
}
else if (peer->un->vassal_len) {
peer->vassal = peer->un->vassal;
peer->vassal_len = peer->un->vassal_len;
corerouter_spawn_vassal(ucr, peer->un, peer->retries+1);
peer->defer_connect = 1;
}
}
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);
+71 -2
View File
@@ -19,6 +19,9 @@ struct fastrouter_session {
int has_key;
uint64_t content_length;
uint64_t buffered;
char *path_info;
uint16_t path_info_len;
};
static struct uwsgi_option fastrouter_options[] = {
@@ -61,9 +64,44 @@ static struct uwsgi_option fastrouter_options[] = {
{"fastrouter-fallback-on-no-key", no_argument, 0, "move to fallback node even if a subscription key is not found", uwsgi_opt_true, &ufr.cr.fallback_on_no_key, 0},
{"fastrouter-force-key", required_argument, 0, "skip uwsgi parsing and directly set a key", uwsgi_opt_set_str, &ufr.force_key, 0},
{"fastrouter-emperor-socket", required_argument, 0, "set the emperor command socket that will receive spawn commands", uwsgi_opt_set_str, &ufr.cr.emperor_socket, 0},
{"fastrouter-defer-connect-timeout", required_argument, 0, "set fastrouter defer connect timeout", uwsgi_opt_set_int, &ufr.cr.defer_connect_timeout, 0},
{"fastrouter-max-retries", required_argument, 0, "set fastrouter max retry attempts", uwsgi_opt_set_int, &ufr.cr.max_retries, 0},
UWSGI_END_OF_OPTIONS
};
static int rebuild_key_for_mountpoint(char *path_info, uint16_t path_info_len, struct corerouter_peer *peer) {
if (path_info_len == 0) return -1;
if (path_info[0] != '/') return -1;
uint16_t len = path_info_len -1;
// is it / ?
if (len == 0) return 0;
// now find the second slash occurrence (if any)
char *second_slash = NULL;
char *last_slash = path_info;
int i;
for(i=0;i<uwsgi.subscription_mountpoints;i++) {
if (len < 1) break;
second_slash = memchr(last_slash+1, '/', len);
if (!second_slash) {
last_slash += 1+len;
break;
}
len -= second_slash - last_slash;
last_slash = second_slash;
}
char *new_key = uwsgi_concat2n(peer->key, peer->key_len, path_info, last_slash - path_info);
uint16_t new_key_len = peer->key_len + (last_slash - path_info);
if (new_key_len <= 0xff) {
memcpy(peer->key, new_key, new_key_len);
peer->key_len = new_key_len;
}
free(new_key);
return 0;
}
static void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) {
struct corerouter_peer *peer = (struct corerouter_peer *) data;
@@ -109,6 +147,14 @@ static void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vall
return;
}
if (uwsgi.subscription_mountpoints) {
if (!uwsgi_strncmp("PATH_INFO", 9, key, keylen)) {
fr->path_info = val;
fr->path_info_len = vallen;
}
}
if (ufr.cr.post_buffering > 0) {
if (!uwsgi_strncmp("CONTENT_LENGTH", 14, key, keylen)) {
fr->content_length = uwsgi_str_num(val, vallen);
@@ -225,8 +271,8 @@ static ssize_t fr_instance_connected(struct corerouter_peer *peer) {
peer->can_retry = 0;
// fix modifiers
peer->in->buf[0] = peer->modifier1;
peer->in->buf[3] = peer->modifier2;
peer->session->main_peer->in->buf[0] = peer->modifier1;
peer->session->main_peer->in->buf[3] = peer->modifier2;
// prepare to write the uwsgi packet
peer->out = peer->session->main_peer->in;
@@ -236,6 +282,7 @@ static ssize_t fr_instance_connected(struct corerouter_peer *peer) {
return fr_instance_send_request(peer);
}
// called after receaving the uwsgi header (read vars)
static ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) {
struct fastrouter_session *fr = (struct fastrouter_session *) main_peer->session;
@@ -322,12 +369,24 @@ static ssize_t fr_recv_uwsgi_vars(struct corerouter_peer *main_peer) {
if (new_peer->key_len == 0)
return -1;
if (uwsgi.subscription_mountpoints) {
if (rebuild_key_for_mountpoint(fr->path_info, fr->path_info_len, new_peer)) return -1;
}
// find an instance using the key
if (ucr->mapper(ucr, new_peer))
return -1;
// check instance
if (new_peer->instance_address_len == 0) {
// check if the connection was deferred
if (new_peer->defer_connect) {
new_peer->current_timeout = ufr.cr.defer_connect_timeout;
new_peer->timeout = corerouter_reset_timeout(&ufr.cr, new_peer);
// stop reading from the client
if (uwsgi_cr_set_hooks(main_peer, NULL, NULL)) return -1;
return len;
}
if (ufr.cr.fallback_on_no_key) {
new_peer->failed = 1;
new_peer->can_retry = 1;
@@ -386,6 +445,16 @@ static int fr_retry(struct corerouter_peer *peer) {
}
if (peer->instance_address_len == 0) {
// first retry is consumed for the first attempt
if (peer->defer_connect && (peer->retries+1) < ufr.cr.max_retries) {
peer->current_timeout = ufr.cr.defer_connect_timeout;
peer->timeout = corerouter_reset_timeout(&ufr.cr, peer);
// stop reading from the client
if (uwsgi_cr_set_hooks(peer->session->main_peer, NULL, NULL)) return -1;
return 1;
}
// ensure deferred connect is disabled
peer->defer_connect = 0;
return -1;
}
+18
View File
@@ -2808,6 +2808,13 @@ struct uwsgi_server {
char *emperor_trigger_socket;
int emperor_trigger_socket_fd;
char *emperor_command_socket;
int emperor_command_socket_fd;
int emperor_wait_for_command;
struct uwsgi_string_list *emperor_wait_for_command_ignore;
int subscription_vassal_required;
int subscription_clear_on_shutdown;
};
struct uwsgi_rpc {
@@ -3415,6 +3422,11 @@ struct uwsgi_subscribe_req {
uint16_t proto_len;
struct uwsgi_subscribe_node *(*algo) (struct uwsgi_subscribe_slot *, struct uwsgi_subscribe_node *, struct uwsgi_subscription_client *);
char *vassal;
uint16_t vassal_len;
uint8_t clear;
};
void uwsgi_nuclear_blast();
@@ -3696,6 +3708,9 @@ struct uwsgi_subscribe_node {
uint64_t backup_level;
//here the solution is a bit hacky, we take the first letter of the proto ('u','\0' -> uwsgi, 'h' -> http, 'f' -> fastcgi, 's' -> scgi)
char proto;
char vassal[0xff];
uint16_t vassal_len;
};
struct uwsgi_subscribe_slot {
@@ -4200,6 +4215,9 @@ struct uwsgi_instance {
// uWSGI 2.1 (vassal's attributes)
struct uwsgi_dyn_dict *attrs;
// when 1 the instance must be manually activated
int suspended;
};
struct uwsgi_instance *emperor_get_by_fd(int);
+1 -1
View File
@@ -1,6 +1,6 @@
# uWSGI build system
uwsgi_version = '2.1'
uwsgi_version = '2.1-dev'
import os
import re