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
30 changed files with 1005 additions and 2675 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;
+1 -2
View File
@@ -396,10 +396,9 @@ static int uwsgi_hook_chown2(char *arg) {
}
#if defined(UWSGI_SUNOS_EXTERN_SETHOSTNAME)
#ifdef __sun__
extern int sethostname(char *, int);
#endif
static int uwsgi_hook_hostname(char *arg) {
#ifdef __CYGWIN__
return -1;
+1 -27
View File
@@ -62,7 +62,6 @@ void uwsgi_init_default() {
uwsgi.cpus = 1;
uwsgi.new_argc = -1;
uwsgi.binary_argc = 1;
uwsgi.backtrace_depth = 64;
uwsgi.max_apps = 64;
@@ -262,33 +261,11 @@ void uwsgi_commandline_config() {
int argc = uwsgi.argc;
char **argv = uwsgi.argv;
// we might want to ignore some arguments not meant for us
char binary_argv0_pretty[256] = {'\0'};
char *binary_argv0_actual = NULL;
if (uwsgi.new_argc > -1 && uwsgi.new_argv) {
argc = uwsgi.new_argc;
argv = uwsgi.new_argv;
}
if (uwsgi.binary_argc > 1 && argc >= uwsgi.binary_argc) {
char *pretty = (char *)binary_argv0_pretty;
strncat(pretty, argv[0], 255);
for (i = 1; i < uwsgi.binary_argc; i++) {
if (strlen(pretty) + 1 + strlen(argv[i]) + 1 > 256)
break;
strcat(pretty, " ");
strcat(pretty, argv[i]);
}
argc -= uwsgi.binary_argc - 1;
argv += uwsgi.binary_argc - 1;
binary_argv0_actual = argv[0];
argv[0] = (char *)binary_argv0_pretty;
}
char *optname;
while ((i = getopt_long(argc, argv, uwsgi.short_options, uwsgi.long_options, &uwsgi.option_index)) != -1) {
@@ -312,8 +289,6 @@ void uwsgi_commandline_config() {
add_exported_option(optname, optarg, 0);
}
if (binary_argv0_actual != NULL)
argv[0] = binary_argv0_actual;
#ifdef UWSGI_DEBUG
uwsgi_log("optind:%d argc:%d\n", optind, uwsgi.argc);
@@ -438,8 +413,7 @@ pid_t uwsgi_daemonize2() {
uwsgi_write_pidfile(uwsgi.pidfile2);
}
if (uwsgi.log_master)
uwsgi_setup_log_master();
if (uwsgi.log_master) uwsgi_setup_log_master();
return uwsgi.mypid;
}
+2 -16
View File
@@ -99,7 +99,7 @@ retry:
uwsgi_log("unable to set PTHREAD_PRIO_INHERIT\n");
exit(1);
}
if (pthread_mutexattr_setrobust(&attr, PTHREAD_MUTEX_ROBUST)) {
if (pthread_mutexattr_setrobust_np(&attr, PTHREAD_MUTEX_ROBUST)) {
uwsgi_log("unable to make the mutex 'robust'\n");
exit(1);
}
@@ -161,7 +161,7 @@ void uwsgi_lock_fast(struct uwsgi_lock_item *uli) {
#ifdef EOWNERDEAD
if (pthread_mutex_lock((pthread_mutex_t *) uli->lock_ptr) == EOWNERDEAD) {
uwsgi_log("[deadlock-detector] a process holding a robust mutex died. recovering...\n");
pthread_mutex_consistent((pthread_mutex_t *) uli->lock_ptr);
pthread_mutex_consistent_np((pthread_mutex_t *) uli->lock_ptr);
}
#else
pthread_mutex_lock((pthread_mutex_t *) uli->lock_ptr);
@@ -476,20 +476,6 @@ void uwsgi_rwunlock_fast(struct uwsgi_lock_item *uli) {
#endif
#ifdef __RUMP__
int semctl(int _0, int _1, int _2, ...) {
return 0;
}
int semget(key_t _0, int _1, int _2) {
return 0;
}
int semop(int _0, struct sembuf * _1, size_t _2) {
return 0;
}
#undef UWSGI_LOCK_ENGINE_NAME
#define UWSGI_LOCK_ENGINE_NAME "fake"
#endif
struct uwsgi_lock_item *uwsgi_lock_ipcsem_init(char *id) {
// used by ftok
-14
View File
@@ -1512,7 +1512,6 @@ static void *logger_thread_loop(void *noarg) {
if (uwsgi.req_log_master) {
logpoll[1].events = POLLIN;
logpoll[1].fd = uwsgi.shared->worker_req_log_pipe[0];
logpolls++;
}
@@ -1552,19 +1551,6 @@ void uwsgi_threaded_logger_spawn() {
}
}
void uwsgi_threaded_logger_worker_spawn() {
pthread_t logger_thread;
pthread_mutex_init(&uwsgi.threaded_logger_lock, NULL);
uwsgi.log_master_buf = uwsgi_malloc(uwsgi.log_master_bufsize);
if (pthread_create(&logger_thread, NULL, logger_thread_loop, NULL)) {
uwsgi_error_safe("uwsgi_threaded_logger_worker_spawn()/pthread_create()");
exit(1);
}
}
void uwsgi_register_log_encoder(char *name, char *(*func)(struct uwsgi_log_encoder *, char *, size_t, size_t *)) {
struct uwsgi_log_encoder *old_ule = NULL, *ule = uwsgi.log_encoders;
-1
View File
@@ -772,7 +772,6 @@ int uwsgi_respawn_worker(int wid) {
for(i=0;i<uwsgi.cores;i++) {
uwsgi.workers[uwsgi.mywid].cores[i].in_request = 0;
memset(&uwsgi.workers[uwsgi.mywid].cores[i].req, 0, sizeof(struct wsgi_request));
memset(uwsgi.workers[uwsgi.mywid].cores[i].buffer, 0, sizeof(struct uwsgi_header));
}
uwsgi_fixup_fds(wid, 0, NULL);
+1 -1
View File
@@ -715,7 +715,7 @@ next:
}
if (uwsgi.manage_script_name) {
if (uwsgi_apps_cnt > 0 && wsgi_req->path_info_len >= 1 && wsgi_req->path_info_pos != -1) {
if (uwsgi_apps_cnt > 0 && wsgi_req->path_info_len > 1 && wsgi_req->path_info_pos != -1) {
// starts with 1 as the 0 app is the default (/) one
int best_found = 0;
char *orig_path_info = wsgi_req->path_info;
+4 -5
View File
@@ -467,7 +467,8 @@ cycle:
return received_signal;
}
int uwsgi_receive_signal(struct wsgi_request *wsgi_req, int fd, char *name, int id) {
void uwsgi_receive_signal(struct wsgi_request *wsgi_req, int fd, char *name, int id) {
uint8_t uwsgi_signal;
ssize_t ret = read(fd, &uwsgi_signal, 1);
@@ -486,15 +487,13 @@ int uwsgi_receive_signal(struct wsgi_request *wsgi_req, int fd, char *name, int
if (uwsgi_signal_handler(wsgi_req, uwsgi_signal)) {
uwsgi_log_verbose("error managing signal %d on %s %d\n", uwsgi_signal, name, id);
}
return 1;
}
return 0;
return;
destroy:
// better to kill the whole worker...
uwsgi_log_verbose("uWSGI %s %d screams: UAAAAAAH my master disconnected: I will kill myself!!!\n", name, id);
end_me(0);
// never here
return 0;
}
+1 -13
View File
@@ -465,15 +465,7 @@ void spooler(struct uwsgi_spooler *uspool) {
if (event_queue_wait(spooler_event_queue, timeout, &interesting_fd) > 0) {
if (uwsgi.master_process) {
if (interesting_fd == uwsgi.shared->spooler_signal_pipe[1]) {
if (uwsgi_receive_signal(NULL, interesting_fd, "spooler", (int) getpid())) {
if (uwsgi.spooler_signal_as_task) {
uspool->tasks++;
if (uwsgi.spooler_max_tasks > 0 && uspool->tasks >= (uint64_t) uwsgi.spooler_max_tasks) {
uwsgi_log("[spooler %s pid: %d] maximum number of tasks reached (%d) recycling ...\n", uspool->dir, (int) uwsgi.mypid, uwsgi.spooler_max_tasks);
end_me(0);
}
}
}
uwsgi_receive_signal(NULL, interesting_fd, "spooler", (int) getpid());
}
}
}
@@ -496,11 +488,7 @@ static void spooler_scandir(struct uwsgi_spooler *uspool, char *dir) {
if (!dir)
dir = uspool->dir;
#ifdef __NetBSD__
n = scandir(dir, &tasklist, NULL, (void *)uwsgi_versionsort);
#else
n = scandir(dir, &tasklist, NULL, uwsgi_versionsort);
#endif
if (n < 0) {
uwsgi_error("scandir()");
return;
+1 -1
View File
@@ -459,7 +459,7 @@ int uwsgi_real_file_serve(struct wsgi_request *wsgi_req, char *real_filename, si
}
}
#ifdef UWSGI_DEBUG
uwsgi_log("[uwsgi-fileserve] file %s found, mimetype %s\n", real_filename, mime_type);
uwsgi_log("[uwsgi-fileserve] file %s found\n", real_filename);
#endif
// static file - don't update avg_rt after request
+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) {
-4
View File
@@ -333,11 +333,9 @@ void uwsgi_as_root() {
if (getuid() > 0)
goto nonroot;
#ifndef __RUMP__
if (!uwsgi.master_as_root && !uwsgi.uidname) {
uwsgi_log_initial("uWSGI running as root, you can use --uid/--gid/--chroot options\n");
}
#endif
int in_jail = 0;
@@ -907,11 +905,9 @@ void uwsgi_as_root() {
}
}
#ifndef __RUMP__
if (!getuid()) {
uwsgi_log_initial("*** WARNING: you are running uWSGI as root !!! (use the --uid flag) *** \n");
}
#endif
#ifdef UWSGI_CAP
+11 -45
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},
@@ -326,7 +329,6 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"spooler-processes", required_argument, 0, "set the number of processes for spoolers", uwsgi_opt_set_int, &uwsgi.spooler_numproc, UWSGI_OPT_IMMEDIATE},
{"spooler-quiet", no_argument, 0, "do not be verbose with spooler tasks", uwsgi_opt_true, &uwsgi.spooler_quiet, 0},
{"spooler-max-tasks", required_argument, 0, "set the maximum number of tasks to run before recycling a spooler", uwsgi_opt_set_int, &uwsgi.spooler_max_tasks, 0},
{"spooler-signal-as-task", no_argument, 0, "treat signal events as tasks in spooler, combine used with spooler-max-tasks", uwsgi_opt_true, &uwsgi.spooler_signal_as_task, 0},
{"spooler-harakiri", required_argument, 0, "set harakiri timeout for spooler tasks", uwsgi_opt_set_int, &uwsgi.harakiri_options.spoolers, 0},
{"spooler-frequency", required_argument, 0, "set spooler frequency", uwsgi_opt_set_int, &uwsgi.spooler_frequency, 0},
{"spooler-freq", required_argument, 0, "set spooler frequency", uwsgi_opt_set_int, &uwsgi.spooler_frequency, 0},
@@ -631,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},
@@ -669,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},
@@ -714,8 +718,6 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"req-logger", required_argument, 0, "set/append a request logger", uwsgi_opt_set_req_logger, NULL, UWSGI_OPT_REQ_LOG_MASTER},
{"logger-req", required_argument, 0, "set/append a request logger", uwsgi_opt_set_req_logger, NULL, UWSGI_OPT_REQ_LOG_MASTER},
{"logger", required_argument, 0, "set/append a logger", uwsgi_opt_set_logger, NULL, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
{"worker-logger", required_argument, 0, "set/append a logger in single-worker setup", uwsgi_opt_set_worker_logger, NULL, 0},
{"worker-logger-req", required_argument, 0, "set/append a request logger in single-worker setup", uwsgi_opt_set_req_logger, NULL, 0},
{"logger-list", no_argument, 0, "list enabled loggers", uwsgi_opt_true, &uwsgi.loggers_list, 0},
{"loggers-list", no_argument, 0, "list enabled loggers", uwsgi_opt_true, &uwsgi.loggers_list, 0},
{"threaded-logger", no_argument, 0, "offload log writing to a thread", uwsgi_opt_true, &uwsgi.threaded_logger, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
@@ -723,9 +725,6 @@ static struct uwsgi_option uwsgi_base_options[] = {
{"log-encoder", required_argument, 0, "add an item in the log encoder chain", uwsgi_opt_add_string_list, &uwsgi.requested_log_encoders, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
{"log-req-encoder", required_argument, 0, "add an item in the log req encoder chain", uwsgi_opt_add_string_list, &uwsgi.requested_log_req_encoders, UWSGI_OPT_MASTER | UWSGI_OPT_LOG_MASTER},
{"worker-log-encoder", required_argument, 0, "add an item in the log encoder chain", uwsgi_opt_add_string_list, &uwsgi.requested_log_encoders, 0},
{"worker-log-req-encoder", required_argument, 0, "add an item in the log req encoder chain", uwsgi_opt_add_string_list, &uwsgi.requested_log_req_encoders, 0},
#ifdef UWSGI_PCRE
@@ -1695,18 +1694,6 @@ next:
}
}
}
if (uwsgi.master_fifo) {
// also remove all master fifo
struct uwsgi_string_list *usl;
uwsgi_foreach(usl, uwsgi.master_fifo) {
if (unlink(usl->value)) {
uwsgi_error("unlink()");
}
else {
uwsgi_log("VACUUM: master fifo %s removed.\n", usl->value);
}
}
}
}
}
}
@@ -1850,7 +1837,7 @@ void uwsgi_plugins_atexit(void) {
void uwsgi_backtrace(int depth) {
#if (defined(__GLIBC__) && !defined(__UCLIBC__))|| (defined(__APPLE__) && !defined(NO_EXECINFO)) || defined(UWSGI_HAS_EXECINFO)
#if (defined(__linux__) && !defined(__UCLIBC__))|| (defined(__APPLE__) && !defined(NO_EXECINFO)) || defined(UWSGI_HAS_EXECINFO)
#include <execinfo.h>
@@ -2098,6 +2085,7 @@ void uwsgi_setup(int argc, char *argv[], char *envp[]) {
int i;
struct utsname uuts;
// signal mask is inherited, and sme process manager could make a real mess...
sigset_t smask;
@@ -2370,13 +2358,6 @@ configure:
// setup master logging
if (uwsgi.log_master)
uwsgi_setup_log_master();
else if (uwsgi.numproc == 1 && uwsgi.log_worker && uwsgi.master_process == 0) {
// hack for allowing request loggers
if (uwsgi.requested_req_logger)
uwsgi.req_log_master = 1;
uwsgi_setup_log_master();
uwsgi_threaded_logger_worker_spawn();
}
// setup offload engines
uwsgi_offload_engines_register_all();
@@ -2459,10 +2440,6 @@ configure:
uwsgi_log_initial("compiled with version: %s on %s\n", __VERSION__, UWSGI_BUILD_DATE);
#ifdef __RUMP__
uwsgi_log_initial("Rump system detected\n");
#else
struct utsname uuts;
#ifdef __sun__
if (uname(&uuts) < 0) {
#else
@@ -2475,7 +2452,6 @@ configure:
uwsgi_log_initial("nodename: %s\n", uuts.nodename);
uwsgi_log_initial("machine: %s\n", uuts.machine);
}
#endif
uwsgi_log_initial("clock source: %s\n", uwsgi.clock->name);
#ifdef UWSGI_PCRE
@@ -2688,13 +2664,10 @@ int uwsgi_start(void *v_argv) {
uwsgi_write_pidfile(uwsgi.pidfile2);
}
#ifndef __RUMP__
if (!uwsgi.master_process && !uwsgi.command_mode) {
uwsgi_log_initial("*** WARNING: you are running uWSGI without its master process manager ***\n");
}
#endif
#ifndef __RUMP__
#ifdef RLIMIT_NPROC
if (uwsgi.rl_nproc.rlim_max > 0) {
uwsgi.rl_nproc.rlim_cur = uwsgi.rl_nproc.rlim_max;
@@ -2714,8 +2687,6 @@ int uwsgi_start(void *v_argv) {
}
}
#endif
#endif
#ifndef __OpenBSD__
if (uwsgi.rl.rlim_max > 0) {
@@ -3890,8 +3861,8 @@ void uwsgi_init_all_apps() {
what++;
for (j = 0; j < 256; j++) {
if (uwsgi.p[j]->mount_app) {
uwsgi_log("mounting %s on %s\n", what, app_mps->value[0] == 0 ? "/" : app_mps->value);
if (uwsgi.p[j]->mount_app(app_mps->value[0] == 0 ? "/" : app_mps->value, what) != -1)
uwsgi_log("mounting %s on %s\n", what, app_mps->value);
if (uwsgi.p[j]->mount_app(app_mps->value, what) != -1)
break;
}
}
@@ -4186,11 +4157,6 @@ void uwsgi_opt_set_logger(char *opt, char *value, void *prefix) {
}
}
void uwsgi_opt_set_worker_logger(char *opt, char *value, void *prefix) {
uwsgi_opt_set_logger(opt, value, prefix);
uwsgi.log_worker = 1;
}
void uwsgi_opt_set_req_logger(char *opt, char *value, void *prefix) {
if (!value)
+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;
}
+504 -2442
View File
File diff suppressed because it is too large Load Diff
+1 -12
View File
@@ -711,11 +711,7 @@ PyObject *uwsgi_paste_loader(void *arg1) {
exit(UWSGI_FAILED_APP_CODE);
}
if (up.paste_name) {
paste_arg = PyTuple_New(2);
} else {
paste_arg = PyTuple_New(1);
}
paste_arg = PyTuple_New(1);
if (!paste_arg) {
PyErr_Print();
exit(UWSGI_FAILED_APP_CODE);
@@ -726,13 +722,6 @@ PyObject *uwsgi_paste_loader(void *arg1) {
exit(UWSGI_FAILED_APP_CODE);
}
if (up.paste_name) {
if (PyTuple_SetItem(paste_arg, 1, UWSGI_PYFROMSTRING(up.paste_name))) {
PyErr_Print();
exit(UWSGI_FAILED_APP_CODE);
}
}
paste_app = PyEval_CallObject(paste_loadapp, paste_arg);
if (!paste_app) {
PyErr_Print();
+3 -7
View File
@@ -149,7 +149,6 @@ struct uwsgi_option uwsgi_python_options[] = {
{"paste", required_argument, 0, "load a paste.deploy config file", uwsgi_opt_set_str, &up.paste, 0},
{"paste-logger", no_argument, 0, "enable paste fileConfig logger", uwsgi_opt_true, &up.paste_logger, 0},
{"paste-name", required_argument, 0, "specify the name of the paste section", uwsgi_opt_set_str, &up.paste_name, 0},
{"web3", required_argument, 0, "load a web3 app", uwsgi_opt_set_str, &up.web3, 0},
@@ -205,18 +204,15 @@ struct uwsgi_option uwsgi_python_options[] = {
/* this routine will be called after each fork to reinitialize the various locks */
void uwsgi_python_pthread_prepare(void) {
if (!up.is_dynamically_loading_an_app)
pthread_mutex_lock(&up.lock_pyloaders);
pthread_mutex_lock(&up.lock_pyloaders);
}
void uwsgi_python_pthread_parent(void) {
if (!up.is_dynamically_loading_an_app)
pthread_mutex_unlock(&up.lock_pyloaders);
pthread_mutex_unlock(&up.lock_pyloaders);
}
void uwsgi_python_pthread_child(void) {
if (!up.is_dynamically_loading_an_app)
pthread_mutex_init(&up.lock_pyloaders, NULL);
pthread_mutex_init(&up.lock_pyloaders, NULL);
}
PyMethodDef uwsgi_spit_method[] = { {"uwsgi_spit", py_uwsgi_spit, METH_VARARGS, ""} };
+2 -1
View File
@@ -11,7 +11,8 @@ extern struct uwsgi_plugin python_plugin;
PyTypeObject uwsgi_RequestContextType = {
PyVarObject_HEAD_INIT(NULL, 0)
PyObject_HEAD_INIT(NULL)
0,
"uwsgi.RequestContext",
sizeof(uwsgi_RequestContext),
0,
-4
View File
@@ -154,7 +154,6 @@ struct uwsgi_python {
char *file_config;
char *paste;
int paste_logger;
char *paste_name;
char *eval;
char *web3;
@@ -214,9 +213,6 @@ struct uwsgi_python {
int call_osafterfork;
int pre_initialized;
// when 1 we have the app-loading lock held
int is_dynamically_loading_an_app;
};
+4 -7
View File
@@ -71,14 +71,11 @@ if 'UWSGI_PYTHON_NOLIB' not in os.environ:
LIBS.append('-lutil')
else:
try:
libdir = sysconfig.get_config_var('LIBDIR')
LDFLAGS.append("-L%s" % sysconfig.get_config_var('LIBDIR'))
os.environ['LD_RUN_PATH'] = "%s" % (sysconfig.get_config_var('LIBDIR'))
except:
libdir = "%s/lib" % sysconfig.PREFIX
LDFLAGS.append("-L%s" % libdir)
LDFLAGS.append("-Wl,-rpath=%s" % libdir)
os.environ['LD_RUN_PATH'] = "%s" % libdir
LDFLAGS.append("-L%s/lib" % sysconfig.PREFIX)
os.environ['LD_RUN_PATH'] = "%s/lib" % sysconfig.PREFIX
LIBS.append('-lpython%s' % get_python_version())
else:
-2
View File
@@ -339,7 +339,6 @@ int uwsgi_request_wsgi(struct wsgi_request *wsgi_req) {
// this part must be heavy locked in threaded modes
if (uwsgi.threads > 1) {
pthread_mutex_lock(&up.lock_pyloaders);
up.is_dynamically_loading_an_app = 1;
}
}
@@ -364,7 +363,6 @@ int uwsgi_request_wsgi(struct wsgi_request *wsgi_req) {
if (wsgi_req->dynamic) {
if (uwsgi.threads > 1) {
up.is_dynamically_loading_an_app = 0;
pthread_mutex_unlock(&up.lock_pyloaders);
}
}
+1 -3
View File
@@ -47,7 +47,6 @@ struct uwsgi_rados_mountpoint {
char *allow_delete;
char *allow_mkcol;
char *allow_propfind;
char *username;
};
static struct uwsgi_option uwsgi_rados_options[] = {
@@ -406,7 +405,6 @@ static void uwsgi_rados_add_mountpoint(char *arg, size_t arg_len) {
"allow_delete", &urmp->allow_delete,
"allow_mkcol", &urmp->allow_mkcol,
"allow_propfind", &urmp->allow_propfind,
"username", &urmp->username,
NULL)) {
uwsgi_log("unable to parse rados mountpoint definition\n");
exit(1);
@@ -425,7 +423,7 @@ static void uwsgi_rados_add_mountpoint(char *arg, size_t arg_len) {
uwsgi_log("[rados] mounting %s ...\n", urmp->mountpoint);
rados_t cluster;
if (rados_create(&cluster, urmp->username) < 0) {
if (rados_create(&cluster, NULL) < 0) {
uwsgi_error("can't create Ceph cluster handle");
exit(1);
}
-13
View File
@@ -16,14 +16,6 @@ it exports values exposed by the metric subsystem
extern struct uwsgi_server uwsgi;
struct uwsgi_stats_pusher_statsd {
int no_workers;
} u_stats_pusher_statsd;
static struct uwsgi_option stats_pusher_statsd_options[] = {
{"statsd-no-workers", no_argument, 0, "disable generation of single worker metrics", uwsgi_opt_true, &u_stats_pusher_statsd.no_workers, 0}
};
// configuration of a statsd node
struct statsd_node {
int fd;
@@ -94,9 +86,6 @@ static void stats_pusher_statsd(struct uwsgi_stats_pusher_instance *uspi, time_t
struct uwsgi_buffer *ub = uwsgi_buffer_new(uwsgi.page_size);
struct uwsgi_metric *um = uwsgi.metrics;
while(um) {
if (u_stats_pusher_statsd.no_workers && !uwsgi_starts_with(um->name, um->name_len, "worker.", 7)) {
goto next;
}
uwsgi_rlock(uwsgi.metrics_lock);
// ignore return value
if (um->type == UWSGI_METRIC_GAUGE) {
@@ -111,7 +100,6 @@ static void stats_pusher_statsd(struct uwsgi_stats_pusher_instance *uspi, time_t
*um->value = um->initial_value;
uwsgi_rwunlock(uwsgi.metrics_lock);
}
next:
um = um->next;
}
uwsgi_buffer_destroy(ub);
@@ -126,7 +114,6 @@ static void stats_pusher_statsd_init(void) {
struct uwsgi_plugin stats_pusher_statsd_plugin = {
.name = "stats_pusher_statsd",
.options = stats_pusher_statsd_options,
.on_load = stats_pusher_statsd_init,
};
-2
View File
@@ -3,8 +3,6 @@ import os
NAME = 'systemd_logger'
CFLAGS = os.popen('pkg-config --cflags libsystemd-journal').read().rstrip().split()
CFLAGS += os.popen('pkg-config --cflags libsystemd').read().rstrip().split()
LDFLAGS = []
LIBS = os.popen('pkg-config --libs libsystemd-journal').read().rstrip().split()
LIBS += os.popen('pkg-config --libs libsystemd').read().rstrip().split()
GCC_LIST = ['systemd_logger']
+35 -25
View File
@@ -153,7 +153,22 @@ extern "C" {
#endif
#endif
#if defined(__linux__) || defined(__GNUC__)
#ifndef _GNU_SOURCE
#define _GNU_SOURCE
#endif
#include <stdio.h>
#ifdef __UCLIBC__
#include <sched.h>
#endif
#undef _GNU_SOURCE
#include <stdlib.h>
#include <stddef.h>
#include <signal.h>
#include <math.h>
#include <sys/types.h>
#ifdef __linux__
#ifndef _GNU_SOURCE
#define _GNU_SOURCE
#endif
@@ -161,14 +176,6 @@ extern "C" {
#define __USE_GNU
#endif
#endif
#include <stdio.h>
#include <stdlib.h>
#include <stddef.h>
#include <signal.h>
#include <math.h>
#include <sys/types.h>
#include <sys/socket.h>
#include <net/if.h>
#ifdef __linux__
@@ -176,6 +183,7 @@ extern "C" {
#define MSG_FASTOPEN 0x20000000
#endif
#endif
#undef _GNU_SOURCE
#include <netinet/in.h>
#include <termios.h>
@@ -261,9 +269,6 @@ extern int pivot_root(const char *new_root, const char *put_old);
#include <stdint.h>
#include <sys/wait.h>
#ifndef WAIT_ANY
#define WAIT_ANY (-1)
#endif
#ifdef __APPLE__
#ifndef MAC_OS_X_VERSION_MIN_REQUIRED
@@ -339,9 +344,6 @@ extern int pivot_root(const char *new_root, const char *put_old);
#ifdef __sun__
#undef __EXTENSIONS__
#endif
#ifdef _GNU_SOURCE
#undef _GNU_SOURCE
#endif
#define UWSGI_CACHE_FLAG_UNGETTABLE 0x01
#define UWSGI_CACHE_FLAG_UPDATE 1 << 1
@@ -2373,7 +2375,6 @@ struct uwsgi_server {
int spooler_quiet;
int spooler_frequency;
int snmp;
char *snmp_addr;
char *snmp_community;
@@ -2808,12 +2809,12 @@ struct uwsgi_server {
char *emperor_trigger_socket;
int emperor_trigger_socket_fd;
int spooler_signal_as_task;
int log_worker;
// number of args to consider part of binary_path
int binary_argc;
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 {
@@ -3421,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();
@@ -3702,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 {
@@ -3806,7 +3815,6 @@ void uwsgi_opt_set_str(char *, char *, void *);
void uwsgi_opt_custom(char *, char *, void *);
void uwsgi_opt_set_null(char *, char *, void *);
void uwsgi_opt_set_logger(char *, char *, void *);
void uwsgi_opt_set_worker_logger(char *, char *, void *);
void uwsgi_opt_set_req_logger(char *, char *, void *);
void uwsgi_opt_set_str_spaced(char *, char *, void *);
void uwsgi_opt_add_string_list(char *, char *, void *);
@@ -4007,7 +4015,7 @@ int uwsgi_is_file2(char *, struct stat *);
int uwsgi_is_dir(char *);
int uwsgi_is_link(char *);
int uwsgi_receive_signal(struct wsgi_request *, int, char *, int);
void uwsgi_receive_signal(struct wsgi_request *, int, char *, int);
void uwsgi_exec_atexit(void);
struct uwsgi_stats {
@@ -4207,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);
@@ -4553,7 +4564,6 @@ void uwsgi_master_manage_emperor(void);
void uwsgi_master_manage_udp(int);
void uwsgi_threaded_logger_spawn(void);
void uwsgi_threaded_logger_worker_spawn(void);
void uwsgi_master_check_idle(void);
int uwsgi_master_check_workers_deadline(void);
+3 -7
View File
@@ -5,7 +5,7 @@ uwsgi_version = '2.1-dev'
import os
import re
import time
uwsgi_os = os.environ.get('UWSGI_FORCE_OS', os.uname()[0])
uwsgi_os = os.uname()[0]
uwsgi_os_k = re.split('[-+]', os.uname()[2])[0]
uwsgi_os_v = os.uname()[3]
uwsgi_cpu = os.uname()[4]
@@ -31,6 +31,7 @@ GCC = os.environ.get('CC', sysconfig.get_config_var('CC'))
if not GCC:
GCC = 'gcc'
def get_preprocessor():
if 'clang' in GCC:
return 'clang -xc core/clang_fake.c'
@@ -809,12 +810,7 @@ class uConf(object):
self.libs.append('-lsendfile')
self.libs.append('-lrt')
self.gcc_list.append('lib/sun_fixes')
sunos_major = int(uwsgi_os_k.split('.')[0])
sunos_minor = int(uwsgi_os_k.split('.')[1])
# solaris < 11 does not have sethostname declared in unistd
if not (sunos_major == 5 and sunos_minor > 10):
self.cflags.append('-DUWSGI_SUNOS_EXTERN_SETHOSTNAME')
self.ldflags.append('-L/lib')
self.ldflags.append('-L/lib')
if not uwsgi_os_v.startswith('Nexenta'):
self.libs.remove('-rdynamic')