openrvdas.server.logger_supervisor
No module-level documentation available.
1#!/usr/bin/env python3 2""" 3""" 4import logging 5import time 6import threading 7 8from logger.utils.stderr_logging import DEFAULT_LOGGING_FORMAT # noqa: E402 9from logger.utils.read_config import read_config, expand_cruise_definition # noqa: E402 10 11from server.logger_runner import LoggerRunner # noqa: E402 12 13 14################################################################################ 15class LoggerSupervisor: 16 """Given a map of {logger:config}, start the configs and make sure 17 they keep running. 18 """ 19 20 def __init__(self, configs=None, stderr_file_pattern=None, stderr_data_server=None, 21 max_tries=3, min_uptime=10, interval=1, 22 logger_log_level=logging.WARNING): 23 """ 24 ``` 25 configs - dict of {logger_name: config} that are to be run 26 27 stderr_file_pattern - Pattern into which logger name will be interpolated 28 to create the file path/name to which the logger's stderr 29 will be written. E.g. '/var/log/openrvdas/{logger}.stderr' 30 31 stderr_data_server - If not None, host:port of cached data server to 32 send stderr messages to. 33 34 max_tries - number of times to try a dead logger config. If zero, then 35 never stop retrying. 36 37 min_uptime - How many seconds a logger must be up to count as having been 38 successfully started and reset max_tries. 39 40 interval - How many seconds between checks that a logger is still running. 41 42 logger_log_level - at what system log level the logger should log (if 43 it were a woodchuck chucking wood) 44 ``` 45 """ 46 self.configs = configs or {} 47 self.stderr_file_pattern = stderr_file_pattern 48 self.stderr_data_server = stderr_data_server 49 self.max_tries = max_tries 50 self.min_uptime = min_uptime 51 self.interval = interval 52 self.logger_log_level = logger_log_level 53 54 # Where we store the map from logger name to config actually # noqa: E402 55 # running. Also map from logger name to LoggerRunner that's doing # noqa: E402 56 # the actual work. 57 self.logger_config_map = {} 58 self.logger_runner_map = {} 59 self.logger_map_lock = threading.Lock() 60 61 # When each logger was last restarted, and how many times it has 62 # failed shortly after start so we know when to give up. 63 self.logger_last_started = {} 64 self.logger_restart_counts = {} 65 66 self.quit_flag = False 67 68 ################### 69 def run(self): 70 if self.configs: 71 self.update_configs() 72 73 while not self.quit_flag: 74 self._check_loggers() 75 time.sleep(self.interval) 76 77 ################### 78 def quit(self): 79 with self.logger_map_lock: 80 self.quit_flag = True 81 82 loggers = set(self.logger_runner_map) 83 for logger in loggers: 84 self._delete_logger(logger) 85 86 ################### 87 def _check_loggers(self): 88 logging.info('Checking loggers...') 89 with self.logger_map_lock: 90 for logger, runner in self.logger_runner_map.items(): 91 if not runner.is_runnable(): 92 logging.info('%s - okay; not runnable', logger) 93 continue 94 95 if runner.is_alive(): 96 logging.info('%s - okay; running', logger) 97 continue 98 99 if runner.is_failed(): 100 logging.info('%s - failed', logger) 101 continue 102 103 # If we're here, runner is runnable is not running, and hasn't 104 # yet been labeled as a failed logger. 105 logging.warning('%s unexpectedly dead.', logger) 106 107 # How long has it been since we've restarted? If a long enough 108 # time has passed, give it a clean slate. 109 last_started = self.logger_last_started.get(logger, 0) 110 restart_count = self.logger_restart_counts.get(logger, 0) 111 if time.time() - last_started < self.min_uptime: 112 restart_count += 1 113 else: 114 restart_count = 0 115 self.logger_restart_counts[logger] = restart_count 116 117 # If we've restarted too many times recently, declare the 118 # logger failed and move on. 119 if restart_count >= self.max_tries: 120 runner.failed = True 121 logging.warning('%s has failed %s times; not restarting', 122 logger, restart_count) 123 continue 124 125 # If here, we're going to try restarting. 126 logging.info('%s - restarting', logger) 127 self.logger_last_started[logger] = time.time() 128 runner.start() 129 130 ################### 131 def _start_logger(self, logger, config): 132 """ONLY CALL THIS FROM WITHIN update_configs for thread safety.""" 133 config_name = config.get('name', logger + '_config') 134 logging.info('Called start_logger for %s: %s', logger, config_name) 135 136 self.logger_config_map[logger] = config 137 stderr_filename = self.stderr_file_pattern.format(logger=logger) 138 139 runner = LoggerRunner(config=config, name=logger, 140 stderr_filename=stderr_filename, 141 stderr_data_server=self.stderr_data_server, 142 logger_log_level=self.logger_log_level) 143 self.logger_runner_map[logger] = runner 144 self.logger_runner_map[logger].start() 145 146 ################### 147 148 def _delete_logger(self, logger): 149 """ONLY CALL THIS FROM WITHIN update_configs for thread safety.""" 150 runner = self.logger_runner_map.get(logger) 151 if not runner: 152 logging.warning('Stale logger %s not found?!?', logger) 153 return 154 logging.info('Waiting for logger %s to complete', logger) 155 runner.quit() 156 157 del self.logger_config_map[logger] 158 del self.logger_runner_map[logger] 159 160 ################### 161 def update_configs(self, configs=None): 162 """Receive a new map of logger:config and start/stop loggers as 163 necessary. 164 """ 165 configs = configs or self.configs 166 if not configs: 167 logging.warning('No logger configs to run!') 168 169 with self.logger_map_lock: 170 # If we're in the process of quitting, go home - a different thread is 171 # already shutting things down. 172 if self.quit_flag: 173 return 174 175 stale_loggers = set(self.logger_config_map) - set(configs) 176 new_loggers = set(configs) - set(self.logger_config_map) 177 other_loggers = set(self.logger_config_map) - stale_loggers - new_loggers 178 179 logging.debug('Stale: %s', stale_loggers) 180 logging.debug('New: %s', new_loggers) 181 logging.debug('Other: %s', other_loggers) 182 183 # Find and shut down loggers that don't exist in our new configs 184 for logger in stale_loggers: 185 logging.info('Shutting down logger %s.', logger) 186 self._delete_logger(logger) 187 188 # Add new loggers that have first appeared in our new config and 189 # start them up. 190 for new_logger in new_loggers: 191 new_config = configs[new_logger] 192 logging.info('Starting new logger %s with %s.', new_logger, 193 new_config.get('name', 'no_name')) 194 self._start_logger(new_logger, new_config) 195 196 # For existing loggers, see whether their configs have 197 # changed. If so stop and restart with new config. start them 198 # up. 199 for logger in other_loggers: 200 new_config = configs[logger] 201 old_config = self.logger_config_map[logger] 202 if new_config == old_config: 203 logging.debug('Config for %s unchanged.', logger) 204 continue 205 206 logging.info('Updating %s from %s to %s', logger, 207 old_config.get('name', 'no_name'), 208 new_config.get('name', 'no_name')) 209 self._delete_logger(logger) 210 self._start_logger(logger, new_config) 211 212 ################### 213 214 def get_status(self): 215 """Return a dict of the current config name and current run status of 216 each logger in the form, e.g.: 217 218 {'s330': {'config':'s330->net', 'status':'RUNNING'}, 219 'gyr1': {'config':'gyr1->file', 'status':'FAILED'}, 220 } 221 222 Possible status are EXITED, RUNNING, FAILED and STARTING. EXITED 223 is the status when a logger is not 'runnable' - e.g. the 'off' 224 config. STARTING is used when a runner is runnable but is not 225 running and is not FAILED - i.e. we haven't given up on it. 226 """ 227 logger_status = {} 228 with self.logger_map_lock: 229 for logger, runner in self.logger_runner_map.items(): 230 231 logger_config = self.logger_config_map[logger] 232 config_name = logger_config.get('name', 'no name') 233 234 if not runner.is_runnable(): 235 status = 'EXITED' 236 elif runner.is_alive(): 237 status = 'RUNNING' 238 elif runner.is_failed(): 239 status = 'FAILED' 240 else: 241 status = 'STARTING' 242 logger_status[logger] = {'config': config_name, 'status': status} 243 return logger_status 244 245 246################################################################################ 247if __name__ == '__main__': 248 import argparse 249 parser = argparse.ArgumentParser() 250 parser.add_argument('--config', dest='config', action='store', required=True, 251 help='Initial set of configs to run.') 252 253 parser.add_argument('--stderr_file_pattern', dest='stderr_file_pattern', 254 default='/var/log/openrvdas/{logger}.stderr', 255 help='Pattern into which logger name will be ' 256 'interpolated to create the file path/name to which ' 257 'the logger\'s stderr will be written. E.g. ' 258 '\'/var/log/openrvdas/{logger}.stderr\'') 259 260 parser.add_argument('--stderr_data_server', dest='stderr_data_server', default=None, 261 help='Optional host:port of a cached data server to which ' 262 ' stderr messages should be written.') 263 264 parser.add_argument('--max_tries', dest='max_tries', action='store', 265 type=int, default=3, help='How many times to try a ' 266 'crashing config before giving up on it as failed. If ' 267 'zero, then never stop retrying.') 268 269 parser.add_argument('--min_uptime', dest='min_uptime', action='store', 270 type=float, default=60, help='How many seconds a logger ' 271 'must be up to count as having been successfully ' 272 'started and reset max_tries.') 273 274 parser.add_argument('--interval', dest='interval', action='store', 275 type=float, default=1, help='How many seconds between ' 276 'checks that a logger is still running.') 277 278 parser.add_argument('-v', '--verbosity', dest='verbosity', 279 default=0, action='count', 280 help='Increase output verbosity') 281 282 parser.add_argument('-V', '--logger_verbosity', dest='logger_verbosity', 283 default=0, action='count', 284 help='Increase output verbosity of component loggers') 285 286 parser.add_argument('--mode', dest='mode', required=True, 287 help='Cruise mode to select') 288 289 args = parser.parse_args() 290 291 # Set up logging first of all 292 293 LOG_LEVELS = {0: logging.WARNING, 1: logging.INFO, 2: logging.DEBUG} 294 log_level = LOG_LEVELS[min(args.verbosity, max(LOG_LEVELS))] 295 logging.basicConfig(format=DEFAULT_LOGGING_FORMAT, level=log_level) 296 297 # What level do we want our component loggers to write? 298 logger_log_level = LOG_LEVELS[min(args.logger_verbosity, max(LOG_LEVELS))] 299 300 config = read_config(args.config) 301 config = expand_cruise_definition(config) 302 303 mode_config_names = config.get('modes').get(args.mode) 304 all_configs = config.get('configs') 305 mode_configs = {logger: all_configs.get(mode_config_names[logger]) 306 for logger in mode_config_names} 307 sup = LoggerSupervisor(configs=mode_configs, 308 stderr_file_pattern=args.stderr_file_pattern, 309 stderr_data_server=args.stderr_data_server, 310 max_tries=args.max_tries, 311 min_uptime=args.min_uptime, 312 interval=args.interval, 313 logger_log_level=logger_log_level 314 ) 315 sup.run()
class
LoggerSupervisor:
16class LoggerSupervisor: 17 """Given a map of {logger:config}, start the configs and make sure 18 they keep running. 19 """ 20 21 def __init__(self, configs=None, stderr_file_pattern=None, stderr_data_server=None, 22 max_tries=3, min_uptime=10, interval=1, 23 logger_log_level=logging.WARNING): 24 """ 25 ``` 26 configs - dict of {logger_name: config} that are to be run 27 28 stderr_file_pattern - Pattern into which logger name will be interpolated 29 to create the file path/name to which the logger's stderr 30 will be written. E.g. '/var/log/openrvdas/{logger}.stderr' 31 32 stderr_data_server - If not None, host:port of cached data server to 33 send stderr messages to. 34 35 max_tries - number of times to try a dead logger config. If zero, then 36 never stop retrying. 37 38 min_uptime - How many seconds a logger must be up to count as having been 39 successfully started and reset max_tries. 40 41 interval - How many seconds between checks that a logger is still running. 42 43 logger_log_level - at what system log level the logger should log (if 44 it were a woodchuck chucking wood) 45 ``` 46 """ 47 self.configs = configs or {} 48 self.stderr_file_pattern = stderr_file_pattern 49 self.stderr_data_server = stderr_data_server 50 self.max_tries = max_tries 51 self.min_uptime = min_uptime 52 self.interval = interval 53 self.logger_log_level = logger_log_level 54 55 # Where we store the map from logger name to config actually # noqa: E402 56 # running. Also map from logger name to LoggerRunner that's doing # noqa: E402 57 # the actual work. 58 self.logger_config_map = {} 59 self.logger_runner_map = {} 60 self.logger_map_lock = threading.Lock() 61 62 # When each logger was last restarted, and how many times it has 63 # failed shortly after start so we know when to give up. 64 self.logger_last_started = {} 65 self.logger_restart_counts = {} 66 67 self.quit_flag = False 68 69 ################### 70 def run(self): 71 if self.configs: 72 self.update_configs() 73 74 while not self.quit_flag: 75 self._check_loggers() 76 time.sleep(self.interval) 77 78 ################### 79 def quit(self): 80 with self.logger_map_lock: 81 self.quit_flag = True 82 83 loggers = set(self.logger_runner_map) 84 for logger in loggers: 85 self._delete_logger(logger) 86 87 ################### 88 def _check_loggers(self): 89 logging.info('Checking loggers...') 90 with self.logger_map_lock: 91 for logger, runner in self.logger_runner_map.items(): 92 if not runner.is_runnable(): 93 logging.info('%s - okay; not runnable', logger) 94 continue 95 96 if runner.is_alive(): 97 logging.info('%s - okay; running', logger) 98 continue 99 100 if runner.is_failed(): 101 logging.info('%s - failed', logger) 102 continue 103 104 # If we're here, runner is runnable is not running, and hasn't 105 # yet been labeled as a failed logger. 106 logging.warning('%s unexpectedly dead.', logger) 107 108 # How long has it been since we've restarted? If a long enough 109 # time has passed, give it a clean slate. 110 last_started = self.logger_last_started.get(logger, 0) 111 restart_count = self.logger_restart_counts.get(logger, 0) 112 if time.time() - last_started < self.min_uptime: 113 restart_count += 1 114 else: 115 restart_count = 0 116 self.logger_restart_counts[logger] = restart_count 117 118 # If we've restarted too many times recently, declare the 119 # logger failed and move on. 120 if restart_count >= self.max_tries: 121 runner.failed = True 122 logging.warning('%s has failed %s times; not restarting', 123 logger, restart_count) 124 continue 125 126 # If here, we're going to try restarting. 127 logging.info('%s - restarting', logger) 128 self.logger_last_started[logger] = time.time() 129 runner.start() 130 131 ################### 132 def _start_logger(self, logger, config): 133 """ONLY CALL THIS FROM WITHIN update_configs for thread safety.""" 134 config_name = config.get('name', logger + '_config') 135 logging.info('Called start_logger for %s: %s', logger, config_name) 136 137 self.logger_config_map[logger] = config 138 stderr_filename = self.stderr_file_pattern.format(logger=logger) 139 140 runner = LoggerRunner(config=config, name=logger, 141 stderr_filename=stderr_filename, 142 stderr_data_server=self.stderr_data_server, 143 logger_log_level=self.logger_log_level) 144 self.logger_runner_map[logger] = runner 145 self.logger_runner_map[logger].start() 146 147 ################### 148 149 def _delete_logger(self, logger): 150 """ONLY CALL THIS FROM WITHIN update_configs for thread safety.""" 151 runner = self.logger_runner_map.get(logger) 152 if not runner: 153 logging.warning('Stale logger %s not found?!?', logger) 154 return 155 logging.info('Waiting for logger %s to complete', logger) 156 runner.quit() 157 158 del self.logger_config_map[logger] 159 del self.logger_runner_map[logger] 160 161 ################### 162 def update_configs(self, configs=None): 163 """Receive a new map of logger:config and start/stop loggers as 164 necessary. 165 """ 166 configs = configs or self.configs 167 if not configs: 168 logging.warning('No logger configs to run!') 169 170 with self.logger_map_lock: 171 # If we're in the process of quitting, go home - a different thread is 172 # already shutting things down. 173 if self.quit_flag: 174 return 175 176 stale_loggers = set(self.logger_config_map) - set(configs) 177 new_loggers = set(configs) - set(self.logger_config_map) 178 other_loggers = set(self.logger_config_map) - stale_loggers - new_loggers 179 180 logging.debug('Stale: %s', stale_loggers) 181 logging.debug('New: %s', new_loggers) 182 logging.debug('Other: %s', other_loggers) 183 184 # Find and shut down loggers that don't exist in our new configs 185 for logger in stale_loggers: 186 logging.info('Shutting down logger %s.', logger) 187 self._delete_logger(logger) 188 189 # Add new loggers that have first appeared in our new config and 190 # start them up. 191 for new_logger in new_loggers: 192 new_config = configs[new_logger] 193 logging.info('Starting new logger %s with %s.', new_logger, 194 new_config.get('name', 'no_name')) 195 self._start_logger(new_logger, new_config) 196 197 # For existing loggers, see whether their configs have 198 # changed. If so stop and restart with new config. start them 199 # up. 200 for logger in other_loggers: 201 new_config = configs[logger] 202 old_config = self.logger_config_map[logger] 203 if new_config == old_config: 204 logging.debug('Config for %s unchanged.', logger) 205 continue 206 207 logging.info('Updating %s from %s to %s', logger, 208 old_config.get('name', 'no_name'), 209 new_config.get('name', 'no_name')) 210 self._delete_logger(logger) 211 self._start_logger(logger, new_config) 212 213 ################### 214 215 def get_status(self): 216 """Return a dict of the current config name and current run status of 217 each logger in the form, e.g.: 218 219 {'s330': {'config':'s330->net', 'status':'RUNNING'}, 220 'gyr1': {'config':'gyr1->file', 'status':'FAILED'}, 221 } 222 223 Possible status are EXITED, RUNNING, FAILED and STARTING. EXITED 224 is the status when a logger is not 'runnable' - e.g. the 'off' 225 config. STARTING is used when a runner is runnable but is not 226 running and is not FAILED - i.e. we haven't given up on it. 227 """ 228 logger_status = {} 229 with self.logger_map_lock: 230 for logger, runner in self.logger_runner_map.items(): 231 232 logger_config = self.logger_config_map[logger] 233 config_name = logger_config.get('name', 'no name') 234 235 if not runner.is_runnable(): 236 status = 'EXITED' 237 elif runner.is_alive(): 238 status = 'RUNNING' 239 elif runner.is_failed(): 240 status = 'FAILED' 241 else: 242 status = 'STARTING' 243 logger_status[logger] = {'config': config_name, 'status': status} 244 return logger_status
Given a map of {logger:config}, start the configs and make sure they keep running.
LoggerSupervisor( configs=None, stderr_file_pattern=None, stderr_data_server=None, max_tries=3, min_uptime=10, interval=1, logger_log_level=30)
21 def __init__(self, configs=None, stderr_file_pattern=None, stderr_data_server=None, 22 max_tries=3, min_uptime=10, interval=1, 23 logger_log_level=logging.WARNING): 24 """ 25 ``` 26 configs - dict of {logger_name: config} that are to be run 27 28 stderr_file_pattern - Pattern into which logger name will be interpolated 29 to create the file path/name to which the logger's stderr 30 will be written. E.g. '/var/log/openrvdas/{logger}.stderr' 31 32 stderr_data_server - If not None, host:port of cached data server to 33 send stderr messages to. 34 35 max_tries - number of times to try a dead logger config. If zero, then 36 never stop retrying. 37 38 min_uptime - How many seconds a logger must be up to count as having been 39 successfully started and reset max_tries. 40 41 interval - How many seconds between checks that a logger is still running. 42 43 logger_log_level - at what system log level the logger should log (if 44 it were a woodchuck chucking wood) 45 ``` 46 """ 47 self.configs = configs or {} 48 self.stderr_file_pattern = stderr_file_pattern 49 self.stderr_data_server = stderr_data_server 50 self.max_tries = max_tries 51 self.min_uptime = min_uptime 52 self.interval = interval 53 self.logger_log_level = logger_log_level 54 55 # Where we store the map from logger name to config actually # noqa: E402 56 # running. Also map from logger name to LoggerRunner that's doing # noqa: E402 57 # the actual work. 58 self.logger_config_map = {} 59 self.logger_runner_map = {} 60 self.logger_map_lock = threading.Lock() 61 62 # When each logger was last restarted, and how many times it has 63 # failed shortly after start so we know when to give up. 64 self.logger_last_started = {} 65 self.logger_restart_counts = {} 66 67 self.quit_flag = False
configs - dict of {logger_name: config} that are to be run
stderr_file_pattern - Pattern into which logger name will be interpolated
to create the file path/name to which the logger's stderr
will be written. E.g. '/var/log/openrvdas/{logger}.stderr'
stderr_data_server - If not None, host:port of cached data server to
send stderr messages to.
max_tries - number of times to try a dead logger config. If zero, then
never stop retrying.
min_uptime - How many seconds a logger must be up to count as having been
successfully started and reset max_tries.
interval - How many seconds between checks that a logger is still running.
logger_log_level - at what system log level the logger should log (if
it were a woodchuck chucking wood)
def
update_configs(self, configs=None):
162 def update_configs(self, configs=None): 163 """Receive a new map of logger:config and start/stop loggers as 164 necessary. 165 """ 166 configs = configs or self.configs 167 if not configs: 168 logging.warning('No logger configs to run!') 169 170 with self.logger_map_lock: 171 # If we're in the process of quitting, go home - a different thread is 172 # already shutting things down. 173 if self.quit_flag: 174 return 175 176 stale_loggers = set(self.logger_config_map) - set(configs) 177 new_loggers = set(configs) - set(self.logger_config_map) 178 other_loggers = set(self.logger_config_map) - stale_loggers - new_loggers 179 180 logging.debug('Stale: %s', stale_loggers) 181 logging.debug('New: %s', new_loggers) 182 logging.debug('Other: %s', other_loggers) 183 184 # Find and shut down loggers that don't exist in our new configs 185 for logger in stale_loggers: 186 logging.info('Shutting down logger %s.', logger) 187 self._delete_logger(logger) 188 189 # Add new loggers that have first appeared in our new config and 190 # start them up. 191 for new_logger in new_loggers: 192 new_config = configs[new_logger] 193 logging.info('Starting new logger %s with %s.', new_logger, 194 new_config.get('name', 'no_name')) 195 self._start_logger(new_logger, new_config) 196 197 # For existing loggers, see whether their configs have 198 # changed. If so stop and restart with new config. start them 199 # up. 200 for logger in other_loggers: 201 new_config = configs[logger] 202 old_config = self.logger_config_map[logger] 203 if new_config == old_config: 204 logging.debug('Config for %s unchanged.', logger) 205 continue 206 207 logging.info('Updating %s from %s to %s', logger, 208 old_config.get('name', 'no_name'), 209 new_config.get('name', 'no_name')) 210 self._delete_logger(logger) 211 self._start_logger(logger, new_config)
Receive a new map of logger:config and start/stop loggers as necessary.
def
get_status(self):
215 def get_status(self): 216 """Return a dict of the current config name and current run status of 217 each logger in the form, e.g.: 218 219 {'s330': {'config':'s330->net', 'status':'RUNNING'}, 220 'gyr1': {'config':'gyr1->file', 'status':'FAILED'}, 221 } 222 223 Possible status are EXITED, RUNNING, FAILED and STARTING. EXITED 224 is the status when a logger is not 'runnable' - e.g. the 'off' 225 config. STARTING is used when a runner is runnable but is not 226 running and is not FAILED - i.e. we haven't given up on it. 227 """ 228 logger_status = {} 229 with self.logger_map_lock: 230 for logger, runner in self.logger_runner_map.items(): 231 232 logger_config = self.logger_config_map[logger] 233 config_name = logger_config.get('name', 'no name') 234 235 if not runner.is_runnable(): 236 status = 'EXITED' 237 elif runner.is_alive(): 238 status = 'RUNNING' 239 elif runner.is_failed(): 240 status = 'FAILED' 241 else: 242 status = 'STARTING' 243 logger_status[logger] = {'config': config_name, 'status': status} 244 return logger_status
Return a dict of the current config name and current run status of each logger in the form, e.g.:
{'s330': {'config':'s330->net', 'status':'RUNNING'}, 'gyr1': {'config':'gyr1->file', 'status':'FAILED'}, }
Possible status are EXITED, RUNNING, FAILED and STARTING. EXITED is the status when a logger is not 'runnable' - e.g. the 'off' config. STARTING is used when a runner is runnable but is not running and is not FAILED - i.e. we haven't given up on it.