mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-09-04 04:31:44 +00:00
new emperor infrastructure
--HG-- rename : lib/amqp.c => plugins/emperor_amqp/amqp.c
This commit is contained in:
+313
-255
@@ -5,10 +5,15 @@
|
||||
extern struct uwsgi_server uwsgi;
|
||||
extern char **environ;
|
||||
|
||||
int emperor_queue;
|
||||
char *emperor_absolute_dir;
|
||||
|
||||
void emperor_send_stats(int);
|
||||
struct uwsgi_instance *emperor_get_by_fd(int);
|
||||
struct uwsgi_instance *emperor_get(char *);
|
||||
void emperor_stop(struct uwsgi_instance *);
|
||||
void emperor_respawn(struct uwsgi_instance *, time_t);
|
||||
void emperor_add(char *, time_t, char *, uint32_t, uid_t, gid_t);
|
||||
|
||||
time_t emperor_throttle;
|
||||
int emperor_throttle_level;
|
||||
|
||||
struct uwsgi_instance {
|
||||
struct uwsgi_instance *ui_prev;
|
||||
@@ -40,6 +45,148 @@ struct uwsgi_instance {
|
||||
};
|
||||
|
||||
|
||||
struct uwsgi_emperor_scanner {
|
||||
char *arg;
|
||||
struct uwsgi_imperial_monitor *monitor;
|
||||
struct uwsgi_emperor_scanner *next;
|
||||
};
|
||||
|
||||
struct uwsgi_emperor_scanner *emperor_scanners;
|
||||
|
||||
void uwsgi_imperial_monitor_directory(char *arg) {
|
||||
struct uwsgi_instance *ui_current;
|
||||
struct dirent *de;
|
||||
struct stat st;
|
||||
|
||||
if (chdir(arg)) {
|
||||
uwsgi_error("chdir()");
|
||||
return;
|
||||
}
|
||||
|
||||
DIR *dir = opendir(".");
|
||||
while ((de = readdir(dir)) != NULL) {
|
||||
if (!strcmp(de->d_name + (strlen(de->d_name) - 4), ".xml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 4), ".ini") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 4), ".yml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 5), ".yaml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 3), ".js") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 5), ".json")
|
||||
) {
|
||||
|
||||
|
||||
if (strlen(de->d_name) >= 0xff)
|
||||
continue;
|
||||
|
||||
if (stat(de->d_name, &st))
|
||||
continue;
|
||||
|
||||
if (!S_ISREG(st.st_mode))
|
||||
continue;
|
||||
|
||||
ui_current = emperor_get(de->d_name);
|
||||
|
||||
if (ui_current) {
|
||||
// check if uid or gid are changed, in such case, stop the instance
|
||||
if (uwsgi.emperor_tyrant) {
|
||||
if (st.st_uid != ui_current->uid || st.st_gid != ui_current->gid) {
|
||||
uwsgi_log("!!! permissions of file %s changed. stopping the instance... !!!\n");
|
||||
emperor_stop(ui_current);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
// check if mtime is changed and the uWSGI instance must be reloaded
|
||||
if (st.st_mtime > ui_current->last_mod) {
|
||||
emperor_respawn(ui_current, st.st_mtime);
|
||||
}
|
||||
}
|
||||
else {
|
||||
emperor_add(de->d_name, st.st_mtime, NULL, 0, st.st_uid, st.st_gid);
|
||||
}
|
||||
}
|
||||
}
|
||||
closedir(dir);
|
||||
}
|
||||
|
||||
void uwsgi_imperial_monitor_glob(char *arg) {
|
||||
|
||||
glob_t g;
|
||||
int i;
|
||||
struct stat st;
|
||||
struct uwsgi_instance *ui_current;
|
||||
|
||||
if (glob(arg, GLOB_MARK|GLOB_NOCHECK, NULL, &g)) {
|
||||
uwsgi_error("glob()");
|
||||
return;
|
||||
}
|
||||
|
||||
for (i = 0; i < (int) g.gl_pathc; i++) {
|
||||
if (!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".xml") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".ini") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".yml") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 3), ".js") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 5), ".json") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 5), ".yaml")
|
||||
) {
|
||||
|
||||
|
||||
if (strlen(g.gl_pathv[i]) >= 0xff)
|
||||
continue;
|
||||
|
||||
if (stat(g.gl_pathv[i], &st))
|
||||
continue;
|
||||
|
||||
if (!S_ISREG(st.st_mode))
|
||||
continue;
|
||||
|
||||
ui_current = emperor_get(g.gl_pathv[i]);
|
||||
|
||||
if (ui_current) {
|
||||
// check if uid or gid are changed, in such case, stop the instance
|
||||
if (uwsgi.emperor_tyrant) {
|
||||
if (st.st_uid != ui_current->uid || st.st_gid != ui_current->gid) {
|
||||
uwsgi_log("!!! permissions of file %s changed. stopping the instance... !!!\n");
|
||||
emperor_stop(ui_current);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
// check if mtime is changed and the uWSGI instance must be reloaded
|
||||
if (st.st_mtime > ui_current->last_mod) {
|
||||
emperor_respawn(ui_current, st.st_mtime);
|
||||
}
|
||||
}
|
||||
else {
|
||||
emperor_add(g.gl_pathv[i], st.st_mtime, NULL, 0, st.st_uid, st.st_gid);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
globfree(&g);
|
||||
}
|
||||
|
||||
void uwsgi_register_imperial_monitor(char *name, void (*init)(char *), void (*func)(char *)) {
|
||||
|
||||
struct uwsgi_imperial_monitor *uim = uwsgi.emperor_monitors;
|
||||
if (!uim) {
|
||||
uim = uwsgi_calloc(sizeof(struct uwsgi_imperial_monitor));
|
||||
uwsgi.emperor_monitors = uim;
|
||||
}
|
||||
else {
|
||||
while(uim) {
|
||||
if (!uim->next) {
|
||||
uim->next = uwsgi_calloc(sizeof(struct uwsgi_imperial_monitor));
|
||||
uim = uim->next;
|
||||
break;
|
||||
}
|
||||
uim = uim->next;
|
||||
}
|
||||
}
|
||||
|
||||
uim->scheme = name;
|
||||
uim->init = init;
|
||||
uim->func = func;
|
||||
uim->next = NULL;
|
||||
}
|
||||
|
||||
struct uwsgi_instance *ui;
|
||||
|
||||
static void royal_death(int signum) {
|
||||
@@ -51,9 +198,11 @@ static void royal_death(int signum) {
|
||||
|
||||
while (c_ui) {
|
||||
uwsgi_log("running vassal stop-hook: %s %s\n", uwsgi.vassals_stop_hook, c_ui->name);
|
||||
if (setenv("UWSGI_VASSALS_DIR", emperor_absolute_dir, 1)) {
|
||||
if (uwsgi.emperor_absolute_dir) {
|
||||
if (setenv("UWSGI_VASSALS_DIR", uwsgi.emperor_absolute_dir, 1)) {
|
||||
uwsgi_error("setenv()");
|
||||
}
|
||||
}
|
||||
int stop_hook_ret = uwsgi_run_command_and_wait(uwsgi.vassals_stop_hook, c_ui->name);
|
||||
uwsgi_log("%s stop-hook returned %d\n", c_ui->name, stop_hook_ret);
|
||||
c_ui = c_ui->ui_next;
|
||||
@@ -125,14 +274,16 @@ void emperor_del(struct uwsgi_instance *c_ui) {
|
||||
|
||||
if (uwsgi.vassals_stop_hook) {
|
||||
uwsgi_log("running vassal stop-hook: %s %s\n", uwsgi.vassals_stop_hook, c_ui->name);
|
||||
if (setenv("UWSGI_VASSALS_DIR", emperor_absolute_dir, 1)) {
|
||||
if (uwsgi.emperor_absolute_dir) {
|
||||
if (setenv("UWSGI_VASSALS_DIR", uwsgi.emperor_absolute_dir, 1)) {
|
||||
uwsgi_error("setenv()");
|
||||
}
|
||||
}
|
||||
int stop_hook_ret = uwsgi_run_command_and_wait(uwsgi.vassals_stop_hook, c_ui->name);
|
||||
uwsgi_log("%s stop-hook returned %d\n", c_ui->name, stop_hook_ret);
|
||||
}
|
||||
|
||||
uwsgi_log("removed uwsgi instance %s\n", c_ui->name);
|
||||
uwsgi_log("[uwsgi-emperor] removed uwsgi instance %s\n", c_ui->name);
|
||||
|
||||
if (c_ui->zerg) {
|
||||
uwsgi.emperor_broodlord_count--;
|
||||
@@ -186,7 +337,24 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
|
||||
char *colon = NULL;
|
||||
int i;
|
||||
|
||||
usleep(uwsgi.emperor_throttle*1000);
|
||||
if (uwsgi_now() - emperor_throttle < 1) {
|
||||
emperor_throttle_level = emperor_throttle_level*2;
|
||||
}
|
||||
else {
|
||||
if (emperor_throttle_level > uwsgi.emperor_throttle) {
|
||||
emperor_throttle_level = emperor_throttle_level/2;
|
||||
}
|
||||
|
||||
if (emperor_throttle_level < uwsgi.emperor_throttle) {
|
||||
emperor_throttle_level = uwsgi.emperor_throttle;
|
||||
}
|
||||
}
|
||||
|
||||
emperor_throttle = uwsgi_now();
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("emperor throttle = %d\n", emperor_throttle_level);
|
||||
#endif
|
||||
usleep(emperor_throttle_level);
|
||||
|
||||
if (uwsgi.emperor_tyrant) {
|
||||
if (uid == 0 || gid == 0) {
|
||||
@@ -233,7 +401,7 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
|
||||
goto clear;
|
||||
}
|
||||
|
||||
event_queue_add_fd_read(emperor_queue, n_ui->pipe[0]);
|
||||
event_queue_add_fd_read(uwsgi.emperor_queue, n_ui->pipe[0]);
|
||||
|
||||
if (n_ui->use_config) {
|
||||
if (socketpair(AF_UNIX, SOCK_STREAM, 0, n_ui->pipe_config)) {
|
||||
@@ -418,9 +586,11 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
|
||||
|
||||
if (uwsgi.vassals_start_hook) {
|
||||
uwsgi_log("running vassal start-hook: %s %s\n", uwsgi.vassals_start_hook, n_ui->name);
|
||||
if (setenv("UWSGI_VASSALS_DIR", emperor_absolute_dir, 1)) {
|
||||
if (uwsgi.emperor_absolute_dir) {
|
||||
if (setenv("UWSGI_VASSALS_DIR", uwsgi.emperor_absolute_dir, 1)) {
|
||||
uwsgi_error("setenv()");
|
||||
}
|
||||
}
|
||||
int start_hook_ret = uwsgi_run_command_and_wait(uwsgi.vassals_start_hook, n_ui->name);
|
||||
uwsgi_log("%s start-hook returned %d\n", n_ui->name, start_hook_ret);
|
||||
}
|
||||
@@ -441,25 +611,128 @@ void emperor_add(char *name, time_t born, char *config, uint32_t config_size, ui
|
||||
|
||||
}
|
||||
|
||||
void uwsgi_imperial_monitor_glob_init(char *arg) {
|
||||
if (chdir(uwsgi.cwd)) {
|
||||
uwsgi_error("chdir()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
uwsgi.emperor_absolute_dir = uwsgi.cwd;
|
||||
}
|
||||
|
||||
void uwsgi_imperial_monitor_directory_init(char *arg) {
|
||||
if (chdir(arg)) {
|
||||
uwsgi_error("chdir()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
uwsgi.emperor_absolute_dir = uwsgi_malloc(PATH_MAX+1);
|
||||
if (realpath(".", uwsgi.emperor_absolute_dir) == NULL) {
|
||||
uwsgi_error("realpath()");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
struct uwsgi_imperial_monitor * imperial_monitor_get_by_id(char *scheme) {
|
||||
struct uwsgi_imperial_monitor *uim = uwsgi.emperor_monitors;
|
||||
while(uim) {
|
||||
if (!strcmp(uim->scheme, scheme)) {
|
||||
return uim;
|
||||
}
|
||||
uim = uim->next;
|
||||
}
|
||||
return NULL;
|
||||
}
|
||||
|
||||
struct uwsgi_imperial_monitor * imperial_monitor_get_by_scheme(char *arg) {
|
||||
struct uwsgi_imperial_monitor *uim = uwsgi.emperor_monitors;
|
||||
while(uim) {
|
||||
char *scheme = uwsgi_concat2(uim->scheme, "://");
|
||||
if (!uwsgi_strncmp(scheme, strlen(scheme), arg, strlen(arg))) {
|
||||
free(scheme);
|
||||
return uim;
|
||||
}
|
||||
free(scheme);
|
||||
uim = uim->next;
|
||||
}
|
||||
return NULL;
|
||||
}
|
||||
|
||||
void emperor_add_scanner(struct uwsgi_imperial_monitor *monitor, char *arg) {
|
||||
struct uwsgi_emperor_scanner *ues = emperor_scanners;
|
||||
if (!ues) {
|
||||
ues = uwsgi_calloc(sizeof(struct uwsgi_emperor_scanner));
|
||||
emperor_scanners = ues;
|
||||
}
|
||||
else {
|
||||
while(ues) {
|
||||
if (!ues->next) {
|
||||
ues->next = uwsgi_calloc(sizeof(struct uwsgi_emperor_scanner));
|
||||
ues = ues->next;
|
||||
break;
|
||||
}
|
||||
ues = ues->next;
|
||||
}
|
||||
}
|
||||
|
||||
ues->arg = arg;
|
||||
ues->monitor = monitor;
|
||||
ues->next = NULL;
|
||||
|
||||
// run the init hook
|
||||
ues->monitor->init(arg);
|
||||
}
|
||||
|
||||
void uwsgi_emperor_run_scanners(void) {
|
||||
struct uwsgi_emperor_scanner *ues = emperor_scanners;
|
||||
while(ues) {
|
||||
ues->monitor->func(ues->arg);
|
||||
ues = ues->next;
|
||||
}
|
||||
}
|
||||
|
||||
void emperor_build_scanners() {
|
||||
struct uwsgi_string_list *usl = uwsgi.emperor;
|
||||
glob_t g;
|
||||
while(usl) {
|
||||
struct uwsgi_imperial_monitor *uim = imperial_monitor_get_by_scheme(usl->value);
|
||||
if (uim) {
|
||||
emperor_add_scanner(uim, usl->value);
|
||||
}
|
||||
else {
|
||||
// check for "glob" and fallback to "dir"
|
||||
if (!glob(usl->value, GLOB_MARK|GLOB_NOCHECK, NULL, &g)) {
|
||||
if (g.gl_pathc == 1 && g.gl_pathv[0][strlen(g.gl_pathv[0]) - 1] == '/') {
|
||||
globfree(&g);
|
||||
goto dir;
|
||||
}
|
||||
globfree(&g);
|
||||
uim = imperial_monitor_get_by_id("glob");
|
||||
emperor_add_scanner(uim, usl->value);
|
||||
goto next;
|
||||
}
|
||||
dir:
|
||||
uim = imperial_monitor_get_by_id("dir");
|
||||
emperor_add_scanner(uim, usl->value);
|
||||
}
|
||||
next:
|
||||
usl = usl->next;
|
||||
}
|
||||
}
|
||||
|
||||
void emperor_loop() {
|
||||
|
||||
// monitor a directory
|
||||
|
||||
struct uwsgi_instance ui_base;
|
||||
struct uwsgi_instance *ui_current;
|
||||
struct stat st;
|
||||
|
||||
pid_t diedpid;
|
||||
int waitpid_status;
|
||||
int has_children = 0;
|
||||
int i_am_alone = 0;
|
||||
int simple_mode = 0;
|
||||
glob_t g;
|
||||
int i;
|
||||
struct dirent *de;
|
||||
char *amqp_port;
|
||||
int amqp_fd = -1;
|
||||
char *amqp_routing_key;
|
||||
|
||||
void *events;
|
||||
int nevents;
|
||||
@@ -484,8 +757,14 @@ void emperor_loop() {
|
||||
|
||||
uwsgi.max_fd = rl.rlim_cur;
|
||||
|
||||
emperor_throttle_level = uwsgi.emperor_throttle;
|
||||
|
||||
emperor_queue = event_queue_init();
|
||||
uwsgi_register_imperial_monitor("dir", uwsgi_imperial_monitor_directory_init, uwsgi_imperial_monitor_directory);
|
||||
uwsgi_register_imperial_monitor("glob", uwsgi_imperial_monitor_glob_init, uwsgi_imperial_monitor_glob);
|
||||
|
||||
emperor_build_scanners();
|
||||
|
||||
uwsgi.emperor_queue = event_queue_init();
|
||||
|
||||
events = event_queue_alloc(64);
|
||||
|
||||
@@ -509,57 +788,10 @@ void emperor_loop() {
|
||||
uwsgi.emperor_stats_fd = bind_to_unix(uwsgi.emperor_stats, uwsgi.listen_queue, uwsgi.chmod_socket, uwsgi.abstract_socket);
|
||||
}
|
||||
|
||||
event_queue_add_fd_read(emperor_queue, uwsgi.emperor_stats_fd);
|
||||
event_queue_add_fd_read(uwsgi.emperor_queue, uwsgi.emperor_stats_fd);
|
||||
uwsgi_log("*** Emperor stats server enabled on %s fd: %d ***\n", uwsgi.emperor_stats, uwsgi.emperor_stats_fd);
|
||||
}
|
||||
|
||||
amqp_port = strchr(uwsgi.emperor_dir, ':');
|
||||
|
||||
if (amqp_port) {
|
||||
if (uwsgi.emperor_amqp_vhost == NULL) uwsgi.emperor_amqp_vhost = "/";
|
||||
if (uwsgi.emperor_amqp_username == NULL) uwsgi.emperor_amqp_username = "guest";
|
||||
if (uwsgi.emperor_amqp_password == NULL) uwsgi.emperor_amqp_password = "guest";
|
||||
reconnect:
|
||||
while(amqp_fd == -1) {
|
||||
uwsgi_log("connecting to AMQP server...\n");
|
||||
amqp_fd = uwsgi_connect(uwsgi.emperor_dir, -1, 0);
|
||||
if (amqp_fd < 0) {
|
||||
sleep(1);
|
||||
}
|
||||
}
|
||||
|
||||
uwsgi_log("subscribing to queue...\n");
|
||||
if (uwsgi_amqp_consume_queue(amqp_fd, uwsgi.emperor_amqp_vhost, uwsgi.emperor_amqp_username, uwsgi.emperor_amqp_password, "", "uwsgi.emperor", "fanout") < 0) {
|
||||
close(amqp_fd);
|
||||
amqp_fd = -1;
|
||||
sleep(1);
|
||||
goto reconnect;
|
||||
}
|
||||
|
||||
event_queue_add_fd_read(emperor_queue, amqp_fd);
|
||||
}
|
||||
else {
|
||||
if (!glob(uwsgi.emperor_dir, GLOB_MARK|GLOB_NOCHECK, NULL, &g)) {
|
||||
if (g.gl_pathc == 1 && g.gl_pathv[0][strlen(g.gl_pathv[0]) - 1] == '/') {
|
||||
simple_mode = 1;
|
||||
if (chdir(uwsgi.emperor_dir)) {
|
||||
uwsgi_error("chdir()");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
}
|
||||
else {
|
||||
uwsgi_error("glob()");
|
||||
exit(1);
|
||||
}
|
||||
globfree(&g);
|
||||
emperor_absolute_dir = uwsgi_malloc(PATH_MAX+1);
|
||||
if (realpath(".", emperor_absolute_dir) == NULL) {
|
||||
uwsgi_error("realpath()");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
ui = &ui_base;
|
||||
|
||||
int freq = 0 ;
|
||||
@@ -574,92 +806,13 @@ reconnect:
|
||||
}
|
||||
}
|
||||
|
||||
nevents = event_queue_wait_multi(emperor_queue, freq, events, 64);
|
||||
nevents = event_queue_wait_multi(uwsgi.emperor_queue, freq, events, 64);
|
||||
freq = 3;
|
||||
|
||||
for (i = 0; i<nevents;i++) {
|
||||
interesting_fd = event_queue_interesting_fd(events, i);
|
||||
|
||||
if (amqp_fd > -1 && interesting_fd == amqp_fd ) {
|
||||
uint64_t msgsize;
|
||||
char *config = uwsgi_amqp_consume(amqp_fd, &msgsize, &amqp_routing_key);
|
||||
|
||||
if (!config) {
|
||||
uwsgi_log("problem with RabbitMQ server, trying reconnection...\n");
|
||||
event_queue_del_fd(emperor_queue, amqp_fd, event_queue_read());
|
||||
close(amqp_fd);
|
||||
amqp_fd = -1;
|
||||
goto reconnect;
|
||||
}
|
||||
|
||||
if (amqp_routing_key) {
|
||||
uwsgi_log("AMQP routing_key = %s\n", amqp_routing_key);
|
||||
char *config_file = uwsgi_concat2("emperor://", amqp_routing_key);
|
||||
free(amqp_routing_key);
|
||||
|
||||
ui_current = emperor_get(config_file);
|
||||
|
||||
if (ui_current) {
|
||||
// ignore event if i am not loyal
|
||||
if (!ui_current->loyal) continue;
|
||||
free(ui_current->config);
|
||||
ui_current->config = config;
|
||||
ui_current->config_len = msgsize;
|
||||
if (!msgsize) {
|
||||
// SAFE
|
||||
emperor_del(ui_current);
|
||||
}
|
||||
else {
|
||||
emperor_respawn(ui_current, time(NULL));
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (msgsize > 0) {
|
||||
emperor_add(config_file, time(NULL), config, msgsize, 0, 0);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
free(config_file);
|
||||
}
|
||||
else {
|
||||
if (msgsize) {
|
||||
if (msgsize >= 0xff) { free(config); continue; }
|
||||
|
||||
#ifdef UWSGI_DEBUG
|
||||
uwsgi_log("%.*s\n", (int)msgsize, config);
|
||||
#endif
|
||||
char *config_file = uwsgi_concat2n(config, msgsize, "", 0);
|
||||
free(config);
|
||||
|
||||
if (strncmp(config_file, "http://", 7)) {
|
||||
if (stat(config_file, &st)) {
|
||||
free(config_file);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!S_ISREG(st.st_mode)) {
|
||||
free(config_file);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
ui_current = emperor_get(config_file);
|
||||
|
||||
if (ui_current) {
|
||||
// ignore event if i am not loyal
|
||||
if (!ui_current->loyal) continue;
|
||||
emperor_respawn(ui_current, time(NULL));
|
||||
}
|
||||
else {
|
||||
emperor_add(config_file, time(NULL), NULL, 0, 0, 0);
|
||||
}
|
||||
|
||||
free(config_file);
|
||||
}
|
||||
}
|
||||
}
|
||||
else if (uwsgi.emperor_stats && uwsgi.emperor_stats_fd > -1 && interesting_fd == uwsgi.emperor_stats_fd) {
|
||||
if (uwsgi.emperor_stats && uwsgi.emperor_stats_fd > -1 && interesting_fd == uwsgi.emperor_stats_fd) {
|
||||
emperor_send_stats(uwsgi.emperor_stats_fd);
|
||||
}
|
||||
else {
|
||||
@@ -693,107 +846,13 @@ reconnect:
|
||||
}
|
||||
else {
|
||||
uwsgi_log("[emperor] unrecognized vassal event on fd %d\n", interesting_fd);
|
||||
event_queue_del_fd(emperor_queue, interesting_fd, event_queue_read());
|
||||
event_queue_del_fd(uwsgi.emperor_queue, interesting_fd, event_queue_read());
|
||||
close(interesting_fd);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (amqp_fd == -1) {
|
||||
if (simple_mode) {
|
||||
DIR *dir = opendir(".");
|
||||
while ((de = readdir(dir)) != NULL) {
|
||||
if (!strcmp(de->d_name + (strlen(de->d_name) - 4), ".xml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 4), ".ini") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 4), ".yml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 5), ".yaml") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 3), ".js") ||
|
||||
!strcmp(de->d_name + (strlen(de->d_name) - 5), ".json")
|
||||
) {
|
||||
|
||||
|
||||
if (strlen(de->d_name) >= 0xff)
|
||||
continue;
|
||||
|
||||
if (stat(de->d_name, &st))
|
||||
continue;
|
||||
|
||||
if (!S_ISREG(st.st_mode))
|
||||
continue;
|
||||
|
||||
ui_current = emperor_get(de->d_name);
|
||||
|
||||
if (ui_current) {
|
||||
// check if uid or gid are changed, in such case, sotp the instance
|
||||
if (uwsgi.emperor_tyrant) {
|
||||
if (st.st_uid != ui_current->uid || st.st_gid != ui_current->gid) {
|
||||
uwsgi_log("!!! permissions of file %s changed. stopping the instance... !!!\n");
|
||||
emperor_stop(ui_current);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
// check if mtime is changed and the uWSGI instance must be reloaded
|
||||
if (st.st_mtime > ui_current->last_mod) {
|
||||
emperor_respawn(ui_current, st.st_mtime);
|
||||
}
|
||||
}
|
||||
else {
|
||||
emperor_add(de->d_name, st.st_mtime, NULL, 0, st.st_uid, st.st_gid);
|
||||
}
|
||||
}
|
||||
}
|
||||
closedir(dir);
|
||||
}
|
||||
else {
|
||||
if (glob(uwsgi.emperor_dir, GLOB_MARK|GLOB_NOCHECK, NULL, &g)) {
|
||||
uwsgi_error("glob()");
|
||||
continue;
|
||||
}
|
||||
|
||||
for (i = 0; i < (int) g.gl_pathc; i++) {
|
||||
if (!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".xml") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".ini") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 4), ".yml") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 3), ".js") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 5), ".json") ||
|
||||
!strcmp(g.gl_pathv[i] + (strlen(g.gl_pathv[i]) - 5), ".yaml")
|
||||
) {
|
||||
|
||||
|
||||
if (strlen(g.gl_pathv[i]) >= 0xff)
|
||||
continue;
|
||||
|
||||
if (stat(g.gl_pathv[i], &st))
|
||||
continue;
|
||||
|
||||
if (!S_ISREG(st.st_mode))
|
||||
continue;
|
||||
|
||||
ui_current = emperor_get(g.gl_pathv[i]);
|
||||
|
||||
if (ui_current) {
|
||||
// check if uid or gid are changed, in such case, sotp the instance
|
||||
if (uwsgi.emperor_tyrant) {
|
||||
if (st.st_uid != ui_current->uid || st.st_gid != ui_current->gid) {
|
||||
uwsgi_log("!!! permissions of file %s changed. stopping the instance... !!!\n");
|
||||
emperor_stop(ui_current);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
// check if mtime is changed and the uWSGI instance must be reloaded
|
||||
if (st.st_mtime > ui_current->last_mod) {
|
||||
emperor_respawn(ui_current, st.st_mtime);
|
||||
}
|
||||
}
|
||||
else {
|
||||
emperor_add(g.gl_pathv[i], st.st_mtime, NULL, 0, st.st_uid, st.st_gid);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
globfree(&g);
|
||||
}
|
||||
}
|
||||
uwsgi_emperor_run_scanners();
|
||||
|
||||
// check for removed instances
|
||||
|
||||
@@ -824,24 +883,12 @@ reconnect:
|
||||
ui_current = ui;
|
||||
while (ui_current->ui_next) {
|
||||
ui_current = ui_current->ui_next;
|
||||
// check loyalty if not in simple/glob mode
|
||||
if (amqp_fd > -1) {
|
||||
if (!ui_current->loyal && (time(NULL)-ui_current->last_loyal) > 60) {
|
||||
uwsgi_log("!!! no loyalty showed by instance %s !!!\n", ui_current->name);
|
||||
// SAFE
|
||||
emperor_del(ui_current);
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (ui_current->status == 1) {
|
||||
if (ui_current->config) free(ui_current->config);
|
||||
// SAFE
|
||||
emperor_del(ui_current);
|
||||
break;
|
||||
}
|
||||
else if (!ui_current->use_config && strncmp(ui_current->name, "http://",7) && stat(ui_current->name, &st) && !ui_current->zerg) {
|
||||
emperor_stop(ui_current);
|
||||
}
|
||||
else if (ui_current->pid == diedpid) {
|
||||
if (ui_current->status == 0) {
|
||||
// respawn an accidentally dead instance if its exit code is not UWSGI_EXILE_CODE
|
||||
@@ -892,7 +939,18 @@ void emperor_send_stats(int fd) {
|
||||
char *cwd = uwsgi_get_cwd();
|
||||
if (uwsgi_stats_keyval_comma(us, "cwd", cwd)) goto end0;
|
||||
|
||||
if (uwsgi_stats_keyval_comma(us, "emperor", uwsgi.emperor_dir)) goto end0;
|
||||
if (uwsgi_stats_key(us ,"emperor")) goto end0;
|
||||
if (uwsgi_stats_list_open(us)) goto end0;
|
||||
struct uwsgi_emperor_scanner *ues = emperor_scanners;
|
||||
while(ues) {
|
||||
uwsgi_stats_str(us, ues->arg);
|
||||
ues = ues->next;
|
||||
if (ues) {
|
||||
if (uwsgi_stats_comma(us)) goto end0;
|
||||
}
|
||||
}
|
||||
if (uwsgi_stats_list_close(us)) goto end0;
|
||||
|
||||
if (uwsgi_stats_keylong_comma(us, "emperor_tyrant", (unsigned long long) uwsgi.emperor_tyrant)) goto end0;
|
||||
|
||||
|
||||
|
||||
+5
-8
@@ -132,15 +132,12 @@ static struct uwsgi_option uwsgi_base_options[] = {
|
||||
{"single-interpreter", no_argument, 'i', "do not use multiple interpreters (where available)", uwsgi_opt_true, &uwsgi.single_interpreter, 0},
|
||||
{"need-app", no_argument, 0, "exit if no app can be loaded", uwsgi_opt_true, &uwsgi.need_app, 0},
|
||||
{"master", no_argument, 'M', "enable master process", uwsgi_opt_true, &uwsgi.master_process, 0},
|
||||
{"emperor", required_argument, 0, "run the Emperor", uwsgi_opt_set_str, &uwsgi.emperor_dir, 0},
|
||||
{"emperor", required_argument, 0, "run the Emperor", uwsgi_opt_add_string_list, &uwsgi.emperor, 0},
|
||||
{"emperor-tyrant", no_argument, 0, "put the Emperor in Tyrant mode", uwsgi_opt_true, &uwsgi.emperor_tyrant, 0},
|
||||
{"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},
|
||||
{"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-amqp-vhost", required_argument, 0, "set emperor amqp virtualhost", uwsgi_opt_set_str, &uwsgi.emperor_amqp_vhost, 0},
|
||||
{"emperor-amqp-username", required_argument, 0, "set emperor amqp username", uwsgi_opt_set_str, &uwsgi.emperor_amqp_username, 0},
|
||||
{"emperor-amqp-password", required_argument, 0, "set emperor amqp password", uwsgi_opt_set_str, &uwsgi.emperor_amqp_password, 0},
|
||||
{"emperor-throttle", required_argument, 0, "throttle each vassal spawn (in seconds)", uwsgi_opt_set_int, &uwsgi.emperor_throttle, 0},
|
||||
{"emperor-magic-exec", no_argument, 0, "prefix vassals config files with exec:// if they have the executable bit", uwsgi_opt_true, &uwsgi.emperor_magic_exec, 0},
|
||||
{"vassals-inherit", required_argument, 0, "add config templates to vassals config", uwsgi_opt_add_string_list, &uwsgi.vassals_templates, 0},
|
||||
@@ -1965,7 +1962,7 @@ int main(int argc, char *argv[], char *envp[]) {
|
||||
}
|
||||
|
||||
// start the Emperor if needed
|
||||
if (uwsgi.early_emperor && uwsgi.emperor_dir) {
|
||||
if (uwsgi.early_emperor && uwsgi.emperor) {
|
||||
|
||||
if (!uwsgi.sockets && !ushared->gateways_cnt && !uwsgi.master_process) {
|
||||
uwsgi_notify_ready();
|
||||
@@ -2149,7 +2146,7 @@ int uwsgi_start(void *v_argv) {
|
||||
// end of generic initialization
|
||||
|
||||
// start the Emperor if needed
|
||||
if (!uwsgi.early_emperor && uwsgi.emperor_dir) {
|
||||
if (!uwsgi.early_emperor && uwsgi.emperor) {
|
||||
|
||||
if (!uwsgi.sockets && !ushared->gateways_cnt && !uwsgi.master_process) {
|
||||
uwsgi_notify_ready();
|
||||
@@ -2595,9 +2592,9 @@ nextsock:
|
||||
#endif
|
||||
|
||||
#ifdef UWSGI_UDP
|
||||
if (!uwsgi.sockets && !ushared->gateways_cnt && !uwsgi.no_server && !uwsgi.udp_socket && !uwsgi.emperor_dir && !uwsgi.command_mode) {
|
||||
if (!uwsgi.sockets && !ushared->gateways_cnt && !uwsgi.no_server && !uwsgi.udp_socket && !uwsgi.emperor && !uwsgi.command_mode) {
|
||||
#else
|
||||
if (!uwsgi.sockets && !ushared->gateways_cnt && !uwsgi.no_server && !uwsgi.emperor_dir && !uwsgi.command_mode) {
|
||||
if (!uwsgi.sockets && !ushared->gateways_cnt && !uwsgi.no_server && !uwsgi.emperor && !uwsgi.command_mode) {
|
||||
#endif
|
||||
uwsgi_log("The -s/--socket option is missing and stdin is not a socket.\n");
|
||||
exit(1);
|
||||
|
||||
@@ -1074,6 +1074,13 @@ struct uwsgi_cheaper_algo {
|
||||
struct uwsgi_cheaper_algo *next;
|
||||
};
|
||||
|
||||
struct uwsgi_imperial_monitor {
|
||||
char *scheme;
|
||||
void (*init)(char *);
|
||||
void (*func)(char *);
|
||||
struct uwsgi_imperial_monitor *next;
|
||||
};
|
||||
|
||||
|
||||
struct uwsgi_server {
|
||||
|
||||
@@ -1168,12 +1175,15 @@ struct uwsgi_server {
|
||||
// true if run under the emperor
|
||||
int has_emperor;
|
||||
int emperor_fd;
|
||||
int emperor_queue;
|
||||
int emperor_tyrant;
|
||||
int emperor_fd_config;
|
||||
int early_emperor;
|
||||
int emperor_throttle;
|
||||
int emperor_magic_exec;
|
||||
char *emperor_dir;
|
||||
struct uwsgi_string_list *emperor;
|
||||
struct uwsgi_imperial_monitor *emperor_monitors;
|
||||
char *emperor_absolute_dir;
|
||||
pid_t emperor_pid;
|
||||
int emperor_broodlord;
|
||||
int emperor_broodlord_count;
|
||||
@@ -1183,11 +1193,6 @@ struct uwsgi_server {
|
||||
// true if loyal to the emperor
|
||||
int loyal;
|
||||
|
||||
// amqp support
|
||||
char *emperor_amqp_vhost;
|
||||
char *emperor_amqp_username;
|
||||
char *emperor_amqp_password;
|
||||
|
||||
// emperor hook (still in development)
|
||||
char *vassals_start_hook;
|
||||
char *vassals_stop_hook;
|
||||
@@ -3059,6 +3064,9 @@ void uwsgi_logit_lf_strftime(struct wsgi_request *);
|
||||
struct uwsgi_logvar *uwsgi_logvar_get(struct wsgi_request *, char *, uint8_t);
|
||||
void uwsgi_logvar_add(struct wsgi_request *, char *, uint8_t, char *, uint8_t);
|
||||
|
||||
|
||||
void uwsgi_register_imperial_monitor(char *, void (*)(char *), void (*)(char *));
|
||||
|
||||
#ifdef UWSGI_AS_SHARED_LIBRARY
|
||||
int uwsgi_init(int, char **, char **);
|
||||
#endif
|
||||
|
||||
+2
-2
@@ -347,7 +347,7 @@ class uConf(object):
|
||||
'core/setup_utils',
|
||||
'core/plugins', 'core/lock', 'core/cache',
|
||||
'core/queue', 'core/event', 'core/signal', 'core/cluster',
|
||||
'core/rpc', 'core/gateway', 'core/loop', 'lib/rbtree', 'lib/amqp', 'core/rb_timers', 'core/uwsgi']
|
||||
'core/rpc', 'core/gateway', 'core/loop', 'lib/rbtree', 'core/rb_timers', 'core/uwsgi']
|
||||
# add protocols
|
||||
self.gcc_list.append('proto/base')
|
||||
self.gcc_list.append('proto/uwsgi')
|
||||
@@ -1141,7 +1141,7 @@ if __name__ == "__main__":
|
||||
pass
|
||||
build_plugin('.', None, cflags, [], [], None)
|
||||
elif cmd == '--clean':
|
||||
os.system("rm -f *.o")
|
||||
os.system("rm -f core/*.o")
|
||||
os.system("rm -f proto/*.o")
|
||||
os.system("rm -f lib/*.o")
|
||||
os.system("rm -f plugins/*/*.o")
|
||||
|
||||
Reference in New Issue
Block a user