openrvdas.server.logger_manager
No module-level documentation available.
1#! /usr/bin/env python3 2""" 3""" 4import datetime 5import getpass # to get username 6import logging 7import multiprocessing 8import os 9import signal 10import socket # to get hostname 11import threading 12import time 13 14from importlib import reload 15 16# Imports for running CachedDataServer 17from server.cached_data_server import CachedDataServer # noqa: E402 18 19from server.logger_supervisor import LoggerSupervisor # noqa: E402 20from server.server_api import ServerAPI # noqa: E402 21from logger.transforms.to_das_record_transform import ToDASRecordTransform # noqa: E402 22from logger.utils.stderr_logging import DEFAULT_LOGGING_FORMAT # noqa: E402 23from logger.utils.stderr_logging import StdErrLoggingHandler # noqa: E402 24from logger.utils.read_config import read_config, expand_cruise_definition # noqa: E402 25 26# For sending stderr to CachedDataServer 27from logger.utils.das_record import DASRecord # noqa: E402 28from logger.writers.cached_data_writer import CachedDataWriter # noqa: E402 29from logger.writers.composed_writer import ComposedWriter # noqa: E402 30 31try: 32 from server.sqlite_server_api import SQLiteServerAPI # noqa: E402 33 SQLITE_API_DEFINED = True 34except ImportError: 35 SQLITE_API_DEFINED = False 36 37DEFAULT_MAX_TRIES = 3 38 39SOURCE_NAME = 'LoggerManager' 40USER = getpass.getuser() 41HOSTNAME = socket.gethostname() 42 43DEFAULT_DATA_SERVER_WEBSOCKET = 'localhost:8766' 44 45############################ 46 47 48def kill_handler(self, signum): 49 """Translate an external signal (such as we'd get from os.kill) into a 50 KeyboardInterrupt, which will signal the start() loop to exit nicely.""" 51 raise KeyboardInterrupt('Received external kill signal') 52 53################################################################################ 54################################################################################ 55 56 57class LoggerManager: 58 ############################ 59 def __init__(self, 60 api, supervisor, data_server_websocket=None, 61 stderr_file_pattern='/var/log/openrvdas/{logger}.stderr', 62 interval=0.25, log_level=logging.info, logger_log_level=logging.WARNING): 63 """Read desired/current logger configs from Django DB and try to run the 64 loggers specified in those configs. 65 ``` 66 api - ServerAPI (or subclass) instance by which LoggerManager will get 67 its data store updates 68 69 supervisor - a LoggerSupervisor object to use to manage logger 70 processes. 71 72 data_server_websocket - cached data server host:port to which we are 73 going to send our status updates. 74 75 stderr_file_pattern - Pattern into which logger name will be 76 interpolated to create the file path/name to which the 77 logger's stderr will be written. E.g. 78 '/var/log/openrvdas/{logger}.stderr' If 79 data_server_websocket is defined, will write logger 80 stderr to it. 81 82 interval - number of seconds to sleep between checking/updating loggers 83 84 log_level - LoggerManager's log level 85 86 logger_log_level - At what logging level our component loggers 87 should operate. 88 ``` 89 90 """ 91 # Set signal to catch SIGTERM and convert it into a 92 # KeyboardInterrupt so we can shut things down gracefully. 93 try: 94 signal.signal(signal.SIGTERM, kill_handler) 95 except ValueError: 96 logging.warning('LoggerManager not running in main thread; ' 97 'shutting down with Ctl-C may not work.') 98 99 # api class must be subclass of ServerAPI 100 if not issubclass(type(api), ServerAPI): 101 raise ValueError('Passed api "%s" must be subclass of ServerAPI' % api) 102 self.api = api 103 self.supervisor = supervisor 104 105 # Data server to which we're going to send status updates 106 if data_server_websocket: 107 self.data_server_writer = CachedDataWriter(data_server_websocket) 108 else: 109 self.data_server_writer = None 110 111 self.stderr_file_pattern = stderr_file_pattern 112 self.interval = interval 113 self.logger_log_level = logger_log_level 114 115 # Try to set up logging, right off the bat: reset logging to its 116 # freshly-imported state and add handler that also sends logged 117 # messages to the cached data server. 118 reload(logging) 119 logging.basicConfig(format=DEFAULT_LOGGING_FORMAT, level=log_level) 120 121 if self.data_server_writer: 122 cds_writer = ComposedWriter( 123 transforms=ToDASRecordTransform(data_id='stderr', 124 field_name='stderr:logger_manager'), 125 writers=self.data_server_writer) 126 logging.getLogger().addHandler(StdErrLoggingHandler(cds_writer)) 127 128 # How our various loops and threads will know it's time to quit 129 self.quit_flag = False 130 131 # Where we store the latest cruise definition and status reports. 132 self.cruise = None 133 self.cruise_filename = None 134 self.cruise_loaded_time = 0 135 136 self.loggers = {} 137 self.config_to_logger = {} 138 139 self.logger_status = None 140 self.status_time = 0 141 142 # We loop to check the logger status and pass it off to the cached 143 # data server. Do this in a separate thread. 144 self.check_logger_status_thread = None 145 146 # We'll loop to check the API for updates to our desired 147 # configs. Do this in a separate thread. Also keep track of 148 # currently active configs so that we know when an update is 149 # actually needed. 150 self.update_configs_thread = None 151 self.config_lock = threading.Lock() 152 153 self.active_mode = None # which mode is active now? 154 self.active_configs = None # which configs are active now? 155 156 ############################ 157 def start(self): 158 """Start the threads that make up the LoggerManager operation: 159 160 1. Configuration update loop 161 2. Loop to read logger stderr/status and either output it or 162 transmit it to a cached data server 163 164 Start threads as daemons so that they'll automatically terminate 165 if the main thread does. 166 """ 167 logging.info('Starting LoggerManager') 168 169 # Check logger status in a separate thread. If we've got the 170 # address of a data server websocket, send our updates to it. 171 self.check_logger_status_loop_thread = threading.Thread( 172 name='check_logger_status_loop', 173 target=self._check_logger_status_loop, daemon=True) 174 self.check_logger_status_loop_thread.start() 175 176 # Update configs in a separate thread. 177 self.update_configs_thread = threading.Thread( 178 name='update_configs_loop', 179 target=self._update_configs_loop, daemon=True) 180 self.update_configs_thread.start() 181 182 # Check cruise definition in a separate thread. If we've got the 183 # address of a data server websocket, send our updates to it. 184 self.send_cruise_definition_loop_thread = threading.Thread( 185 name='send_cruise_definition_loop', 186 target=self._send_cruise_definition_loop, daemon=True) 187 self.send_cruise_definition_loop_thread.start() 188 189 ############################ 190 def quit(self): 191 """Exit the loop and shut down all loggers.""" 192 self.quit_flag = True 193 194 ############################ 195 def _load_new_definition_from_api(self): 196 """Fetch a new cruise definition from API and build local maps. Then 197 send an updated cruise definition to the console. 198 """ 199 logging.info('Fetching new cruise definitions from API') 200 with self.config_lock: 201 try: 202 self.loggers = self.api.get_loggers() 203 self.config_to_logger = {} 204 for logger, logger_configs in self.loggers.items(): 205 # Map config_name->logger 206 for config in self.loggers[logger].get('configs', []): 207 self.config_to_logger[config] = logger 208 209 # This is a redundant grab of data when we're called from 210 # _send_cruise_definition_loop(), but we may also be called 211 # from a callback when the API alerts us that something has 212 # changed. So we need to re-grab self.cruise 213 self.cruise = self.api.get_configuration() # a Cruise object 214 self.cruise_filename = self.cruise.get('config_filename') 215 loaded_time = self.cruise.get('loaded_time') 216 self.cruise_loaded_time = datetime.datetime.timestamp(loaded_time) 217 self.active_mode = self.api.get_active_mode() 218 219 # Send updated cruise definition to CDS for console to read. 220 cruise_dict = { 221 'cruise_id': self.cruise.get('id', ''), 222 'filename': self.cruise_filename, 223 'config_timestamp': self.cruise_loaded_time, 224 'loggers': self.loggers, 225 'modes': self.cruise.get('modes', {}), 226 'active_mode': self.active_mode, 227 } 228 logging.info('Sending updated cruise definitions to CDS.') 229 self._write_record_to_data_server( 230 'status:cruise_definition', cruise_dict) 231 except (AttributeError, ValueError, TypeError) as e: 232 logging.info('Failed to update cruise definition: %s', e) 233 234 ############################ 235 def _check_logger_status_loop(self): 236 """Grab logger status message from supervisor and send to cached data 237 server via websocket. Also send cruise mode as separate message. 238 """ 239 while not self.quit_flag: 240 now = time.time() 241 try: 242 config_status = self.supervisor.get_status() 243 with self.config_lock: 244 # Stash status, note time and send update 245 self.config_status = config_status 246 self.status_time = now 247 self._write_record_to_data_server('status:logger_status', config_status) 248 249 # Now get and send cruise mode 250 mode_map = {'active_mode': self.api.get_active_mode()} 251 self._write_record_to_data_server('status:cruise_mode', mode_map) 252 except ValueError as e: 253 logging.warning('Error while trying to send logger status: %s', e) 254 time.sleep(self.interval) 255 256 ############################ 257 def _update_configs_loop(self): 258 """Iteratively check the API for updated configs and send them to the 259 appropriate LoggerRunners. 260 """ 261 while not self.quit_flag: 262 self._update_configs() 263 time.sleep(self.interval) 264 265 ############################ 266 def _update_configs(self): 267 """Get list of new (latest) configs. Send to logger supervisor to make 268 any necessary changes. 269 270 Note: we can't fold this into _update_configs_loop() because we may 271 need to ask the api to call it independently as a callback when it 272 notices that the config has changed. Search for the line: 273 274 api.on_update(callback=logger_manager._update_configs) 275 276 in this file to see where. 277 """ 278 with self.config_lock: 279 # Get new configs in dict {logger:{'configs':[config_name,...]}} 280 logger_configs = self.api.get_logger_configs() 281 if logger_configs: 282 supervisor.update_configs(logger_configs) 283 self.active_configs = logger_configs 284 285 ############################ 286 def _send_cruise_definition_loop(self): 287 """Iteratively assemble information from DB about what loggers should 288 exist and what states they *should* be in. We'll send this to the 289 cached data server whenever it changes (or if it's been a while 290 since we have). 291 292 Also, if the logger or config names have changed, signal that we 293 need to create a new config file for the supervisord process to 294 use. 295 296 Looks like: 297 {'active_mode': 'log', 298 'cruise_id': 'NBP1406', 299 'loggers': {'PCOD': {'active': 'PCOD->file/net', 300 'configs': ['PCOD->off', 301 'PCOD->net', 302 'PCOD->file/net', 303 'PCOD->file/net/db']}, 304 next_logger: next_configs, 305 ... 306 }, 307 'modes': ['off', 'monitor', 'log', 'log+db'] 308 } 309 310 """ 311 last_loaded_timestamp = 0 312 313 while not self.quit_flag: 314 try: 315 self.cruise = self.api.get_configuration() # a Cruise object 316 if not self.cruise: 317 logging.info('No cruise definition found in API') 318 time.sleep(self.interval * 2) 319 continue 320 self.cruise_filename = self.cruise.get('config_filename') 321 loaded_time = self.cruise.get('loaded_time') 322 self.cruise_loaded_time = datetime.datetime.timestamp(loaded_time) 323 324 # Has cruise definition file changed since we loaded it? If so, 325 # send a notification to console so it can ask if user wants to 326 # reload. 327 if self.cruise_filename: 328 try: 329 mtime = os.path.getmtime(self.cruise_filename) 330 if mtime > self.cruise_loaded_time: 331 logging.debug('Cruise file timestamp changed!') 332 self._write_record_to_data_server('status:file_update', mtime) 333 except FileNotFoundError: 334 logging.debug('Cruise file "%s" has disappeared?', self.cruise_filename) 335 336 # Does database have a cruise definition with a newer timestamp? 337 # Means user loaded/reloaded definition. Update our maps to 338 # reflect the new values and send an updated cruise_definition 339 # to the console. 340 if self.cruise_loaded_time > last_loaded_timestamp: 341 last_loaded_timestamp = self.cruise_loaded_time 342 logging.info('New cruise definition detected - rebuilding maps.') 343 self._load_new_definition_from_api() 344 345 except KeyboardInterrupt: # (AttributeError, ValueError, TypeError): 346 logging.warning('No cruise definition found in API') 347 348 # Whether or not we've sent an update, sleep 349 time.sleep(self.interval * 2) 350 351 ############################ 352 def _write_record_to_data_server(self, field_name, record): 353 """Format and label a record and send it to the cached data server. 354 """ 355 if self.data_server_writer: 356 das_record = DASRecord(fields={field_name: record}) 357 logging.debug('DASRecord: %s' % das_record) 358 self.data_server_writer.write(das_record) 359 else: 360 logging.info('Update: %s: %s', field_name, record) 361 362################################################################################ 363 364 365def run_data_server(data_server_websocket, 366 data_server_back_seconds, data_server_cleanup_interval, 367 data_server_interval): 368 """Run a CachedDataServer (to be called as a separate process), 369 accepting websocket connections to receive data to be cached and 370 served. 371 """ 372 # First get the port that we're going to run the data server on. Because 373 # we're running it locally, it should only have a port, not a hostname. 374 # We should try to handle it if they prefix with a ':', though. 375 data_server_websocket = data_server_websocket or DEFAULT_DATA_SERVER_WEBSOCKET 376 websocket_port = int(data_server_websocket.split(':')[-1]) 377 server = CachedDataServer(port=websocket_port, interval=data_server_interval) 378 379 # The server will start serving in its own thread after 380 # initialization, but we need to manually fire up the cleanup loop 381 # if we want it. Maybe we should have this also run automatically in 382 # its own thread after initialization? 383 server.cleanup_loop() 384 385 386################################################################################ 387if __name__ == '__main__': # noqa: C901 388 import argparse 389 import atexit 390 import readline 391 392 from server.server_api_command_line import ServerAPICommandLine 393 394 parser = argparse.ArgumentParser() 395 parser.add_argument('--config', dest='config', action='store', 396 help='Name of configuration file to load.') 397 parser.add_argument('--mode', dest='mode', action='store', default=None, 398 help='Optional name of mode to start system in.') 399 400 database_choices = ['memory', 'django'] 401 if SQLITE_API_DEFINED: 402 database_choices.append('sqlite') 403 parser.add_argument('--database', dest='database', action='store', 404 choices=database_choices, 405 default='memory', help='What backing store database ' 406 'to use.') 407 408 parser.add_argument('--stderr_file_pattern', dest='stderr_file_pattern', 409 default='/var/log/openrvdas/{logger}.stderr', 410 help='Pattern into which logger name will be ' 411 'interpolated to create the file path/name to which ' 412 'the logger\'s stderr will be written. E.g. ' 413 '\'/var/log/openrvdas/{logger}.stderr\'') 414 415 # Arguments for cached data server 416 parser.add_argument('--data_server_websocket', dest='data_server_websocket', 417 action='store', default=None, 418 help='Address at which to connect to cached data server ' 419 'to send status updates.') 420 parser.add_argument('--start_data_server', dest='start_data_server', 421 action='store_true', default=False, 422 help='Whether to start our own cached data server.') 423 parser.add_argument('--data_server_back_seconds', 424 dest='data_server_back_seconds', action='store', 425 type=float, default=480, 426 help='Maximum number of seconds of old data to keep ' 427 'for serving to new clients.') 428 parser.add_argument('--data_server_cleanup_interval', 429 dest='data_server_cleanup_interval', 430 action='store', type=float, default=60, 431 help='How often to clean old data out of the cache.') 432 parser.add_argument('--data_server_interval', dest='data_server_interval', 433 action='store', type=float, default=1, 434 help='How many seconds to sleep between successive ' 435 'sends of data to clients.') 436 437 parser.add_argument('--interval', dest='interval', action='store', 438 type=float, default=0.5, 439 help='How many seconds to sleep between logger checks.') 440 parser.add_argument('--max_tries', dest='max_tries', action='store', type=int, 441 default=DEFAULT_MAX_TRIES, 442 help='Number of times to retry failed loggers.') 443 444 parser.add_argument('--no-console', dest='no_console', default=False, 445 action='store_true', help='Run without a console ' 446 'that reads commands from stdin.') 447 448 parser.add_argument('-v', '--verbosity', dest='verbosity', 449 default=0, action='count', 450 help='Increase output verbosity') 451 parser.add_argument('-V', '--logger_verbosity', dest='logger_verbosity', 452 default=0, action='count', 453 help='Increase output verbosity of component loggers') 454 args = parser.parse_args() 455 456 # Set up logging first of all 457 LOG_LEVELS = {0: logging.WARNING, 1: logging.INFO, 2: logging.DEBUG} 458 459 log_level = LOG_LEVELS[min(args.verbosity, max(LOG_LEVELS))] 460 logging.basicConfig(format=DEFAULT_LOGGING_FORMAT, level=log_level) 461 462 # What level do we want our component loggers to write? 463 logger_log_level = LOG_LEVELS[min(args.logger_verbosity, max(LOG_LEVELS))] 464 465 ############################ 466 # First off, start any servers we're supposed to be running 467 logging.info('Preparing to start LoggerManager.') 468 469 # If we're supposed to be running our own CachedDataServer, start it 470 # here in its own daemon process (daemon so that it dies when we exit). 471 if args.start_data_server: 472 data_server_proc = multiprocessing.Process( 473 name='openrvdas_data_server', 474 target=run_data_server, 475 args=(args.data_server_websocket, 476 args.data_server_back_seconds, args.data_server_cleanup_interval, 477 args.data_server_interval), 478 daemon=True) 479 data_server_proc.start() 480 481 ############################ 482 # If we do have a data server, add a handler that will echo all 483 # logger_manager stderr output to it 484 if args.data_server_websocket: 485 stderr_writer = ComposedWriter( 486 transforms=ToDASRecordTransform(field_name='stderr:logger_manager'), 487 writers=[CachedDataWriter(data_server=args.data_server_websocket)]) 488 logging.getLogger().addHandler(StdErrLoggingHandler(stderr_writer, 489 parse_to_json=True)) 490 491 ############################ 492 # Instantiate API - a Are we using an in-memory store or Django 493 # database as our backing store? Do our imports conditionally, so 494 # they don't actually have to have Django if they're not using it. 495 if args.database == 'django': 496 from django_gui.django_server_api import DjangoServerAPI 497 api = DjangoServerAPI() 498 elif args.database == 'memory': 499 from server.in_memory_server_api import InMemoryServerAPI 500 api = InMemoryServerAPI() 501 elif args.database == 'sqlite': 502 from server.sqlite_server_api import SQLiteServerAPI # noqa F811 503 api = SQLiteServerAPI() 504 else: 505 raise ValueError('Illegal arg for --database: "%s"' % args.database) 506 507 # Now that API is defined, tack on one more logging handler: one 508 # that passes messages to API. 509 # TODO: decide if we even need this. Disabled for now 510 # logging.getLogger().addHandler(WriteToAPILoggingHandler(api)) 511 512 ############################ 513 # Create our logger supervisor. 514 supervisor = LoggerSupervisor( 515 configs=None, 516 stderr_file_pattern=args.stderr_file_pattern, 517 stderr_data_server=args.data_server_websocket, 518 max_tries=args.max_tries, 519 interval=args.interval, 520 logger_log_level=logger_log_level) 521 522 ############################ 523 # Create our LoggerManager 524 logger_manager = LoggerManager( 525 api=api, supervisor=supervisor, 526 data_server_websocket=args.data_server_websocket, 527 stderr_file_pattern=args.stderr_file_pattern, 528 interval=args.interval, 529 log_level=log_level, 530 logger_log_level=logger_log_level) 531 532 # When told to quit, shut down gracefully 533 api.on_quit(callback=logger_manager.quit) 534 api.on_quit(callback=supervisor.quit) 535 536 # When an active config changes in the database, update our configs here 537 api.on_update(callback=logger_manager._update_configs) 538 539 # When new configs are loaded, update our file of config processes 540 api.on_load(callback=logger_manager._load_new_definition_from_api) 541 542 ############################ 543 # If they've given us an initial configuration, get and load it. 544 if args.config: 545 config = read_config(args.config) 546 config = expand_cruise_definition(config) 547 548 # Hacky bit: need to stash the config filename for posterity 549 if 'cruise' not in config or config['cruise'] is None: 550 config['cruise'] = {} 551 552 config['cruise']['config_filename'] = args.config 553 api.load_configuration(config) 554 555 active_mode = args.mode or api.get_default_mode() 556 api.set_active_mode(active_mode) 557 api.message_log(source=SOURCE_NAME, user='(%s@%s)' % (USER, HOSTNAME), 558 log_level=api.INFO, 559 message='started with: %s, mode %s' % 560 (args.config, active_mode)) 561 562 ############################ 563 # Start all the various LoggerManager threads running 564 logger_manager.start() 565 566 try: 567 # If no console, just wait for the configuration update thread to 568 # end as a signal that we're done. 569 if args.no_console: 570 logging.warning('--no-console specified; waiting for LoggerManager ' 571 'to exit.') 572 if logger_manager.update_configs_thread: 573 logger_manager.update_configs_thread.join() 574 else: 575 logging.warning('LoggerManager has no update_configs_thread? ' 576 'Exiting...') 577 else: 578 # Create reader to read/process commands from stdin. Note: this 579 # needs to be in main thread for Ctl-C termination to be properly 580 # caught and processed, otherwise interrupts go to the wrong places. 581 582 # Set up command line interface to get commands. Start by 583 # reading history file, if one exists, to get past commands. 584 hist_filename = '.openrvdas_logger_manager_history' 585 hist_path = os.path.join(os.path.expanduser('~'), hist_filename) 586 try: 587 readline.read_history_file(hist_path) 588 # default history len is -1 (infinite), which may grow unruly 589 readline.set_history_length(1000) 590 except (FileNotFoundError, PermissionError, OSError): 591 pass 592 atexit.register(readline.write_history_file, hist_path) 593 594 command_line_reader = ServerAPICommandLine(api=api) 595 command_line_reader.run() 596 597 except KeyboardInterrupt: 598 pass 599 logging.debug('Done with logger_manager.py - exiting') 600 601 # Ask our SupervisorConnector to shutdown. 602 if supervisor: 603 supervisor.quit()
DEFAULT_MAX_TRIES =
3
SOURCE_NAME =
'LoggerManager'
USER =
'runner'
HOSTNAME =
'runnervmeorf1'
DEFAULT_DATA_SERVER_WEBSOCKET =
'localhost:8766'
def
kill_handler(self, signum):
49def kill_handler(self, signum): 50 """Translate an external signal (such as we'd get from os.kill) into a 51 KeyboardInterrupt, which will signal the start() loop to exit nicely.""" 52 raise KeyboardInterrupt('Received external kill signal')
Translate an external signal (such as we'd get from os.kill) into a KeyboardInterrupt, which will signal the start() loop to exit nicely.
class
LoggerManager:
58class LoggerManager: 59 ############################ 60 def __init__(self, 61 api, supervisor, data_server_websocket=None, 62 stderr_file_pattern='/var/log/openrvdas/{logger}.stderr', 63 interval=0.25, log_level=logging.info, logger_log_level=logging.WARNING): 64 """Read desired/current logger configs from Django DB and try to run the 65 loggers specified in those configs. 66 ``` 67 api - ServerAPI (or subclass) instance by which LoggerManager will get 68 its data store updates 69 70 supervisor - a LoggerSupervisor object to use to manage logger 71 processes. 72 73 data_server_websocket - cached data server host:port to which we are 74 going to send our status updates. 75 76 stderr_file_pattern - Pattern into which logger name will be 77 interpolated to create the file path/name to which the 78 logger's stderr will be written. E.g. 79 '/var/log/openrvdas/{logger}.stderr' If 80 data_server_websocket is defined, will write logger 81 stderr to it. 82 83 interval - number of seconds to sleep between checking/updating loggers 84 85 log_level - LoggerManager's log level 86 87 logger_log_level - At what logging level our component loggers 88 should operate. 89 ``` 90 91 """ 92 # Set signal to catch SIGTERM and convert it into a 93 # KeyboardInterrupt so we can shut things down gracefully. 94 try: 95 signal.signal(signal.SIGTERM, kill_handler) 96 except ValueError: 97 logging.warning('LoggerManager not running in main thread; ' 98 'shutting down with Ctl-C may not work.') 99 100 # api class must be subclass of ServerAPI 101 if not issubclass(type(api), ServerAPI): 102 raise ValueError('Passed api "%s" must be subclass of ServerAPI' % api) 103 self.api = api 104 self.supervisor = supervisor 105 106 # Data server to which we're going to send status updates 107 if data_server_websocket: 108 self.data_server_writer = CachedDataWriter(data_server_websocket) 109 else: 110 self.data_server_writer = None 111 112 self.stderr_file_pattern = stderr_file_pattern 113 self.interval = interval 114 self.logger_log_level = logger_log_level 115 116 # Try to set up logging, right off the bat: reset logging to its 117 # freshly-imported state and add handler that also sends logged 118 # messages to the cached data server. 119 reload(logging) 120 logging.basicConfig(format=DEFAULT_LOGGING_FORMAT, level=log_level) 121 122 if self.data_server_writer: 123 cds_writer = ComposedWriter( 124 transforms=ToDASRecordTransform(data_id='stderr', 125 field_name='stderr:logger_manager'), 126 writers=self.data_server_writer) 127 logging.getLogger().addHandler(StdErrLoggingHandler(cds_writer)) 128 129 # How our various loops and threads will know it's time to quit 130 self.quit_flag = False 131 132 # Where we store the latest cruise definition and status reports. 133 self.cruise = None 134 self.cruise_filename = None 135 self.cruise_loaded_time = 0 136 137 self.loggers = {} 138 self.config_to_logger = {} 139 140 self.logger_status = None 141 self.status_time = 0 142 143 # We loop to check the logger status and pass it off to the cached 144 # data server. Do this in a separate thread. 145 self.check_logger_status_thread = None 146 147 # We'll loop to check the API for updates to our desired 148 # configs. Do this in a separate thread. Also keep track of 149 # currently active configs so that we know when an update is 150 # actually needed. 151 self.update_configs_thread = None 152 self.config_lock = threading.Lock() 153 154 self.active_mode = None # which mode is active now? 155 self.active_configs = None # which configs are active now? 156 157 ############################ 158 def start(self): 159 """Start the threads that make up the LoggerManager operation: 160 161 1. Configuration update loop 162 2. Loop to read logger stderr/status and either output it or 163 transmit it to a cached data server 164 165 Start threads as daemons so that they'll automatically terminate 166 if the main thread does. 167 """ 168 logging.info('Starting LoggerManager') 169 170 # Check logger status in a separate thread. If we've got the 171 # address of a data server websocket, send our updates to it. 172 self.check_logger_status_loop_thread = threading.Thread( 173 name='check_logger_status_loop', 174 target=self._check_logger_status_loop, daemon=True) 175 self.check_logger_status_loop_thread.start() 176 177 # Update configs in a separate thread. 178 self.update_configs_thread = threading.Thread( 179 name='update_configs_loop', 180 target=self._update_configs_loop, daemon=True) 181 self.update_configs_thread.start() 182 183 # Check cruise definition in a separate thread. If we've got the 184 # address of a data server websocket, send our updates to it. 185 self.send_cruise_definition_loop_thread = threading.Thread( 186 name='send_cruise_definition_loop', 187 target=self._send_cruise_definition_loop, daemon=True) 188 self.send_cruise_definition_loop_thread.start() 189 190 ############################ 191 def quit(self): 192 """Exit the loop and shut down all loggers.""" 193 self.quit_flag = True 194 195 ############################ 196 def _load_new_definition_from_api(self): 197 """Fetch a new cruise definition from API and build local maps. Then 198 send an updated cruise definition to the console. 199 """ 200 logging.info('Fetching new cruise definitions from API') 201 with self.config_lock: 202 try: 203 self.loggers = self.api.get_loggers() 204 self.config_to_logger = {} 205 for logger, logger_configs in self.loggers.items(): 206 # Map config_name->logger 207 for config in self.loggers[logger].get('configs', []): 208 self.config_to_logger[config] = logger 209 210 # This is a redundant grab of data when we're called from 211 # _send_cruise_definition_loop(), but we may also be called 212 # from a callback when the API alerts us that something has 213 # changed. So we need to re-grab self.cruise 214 self.cruise = self.api.get_configuration() # a Cruise object 215 self.cruise_filename = self.cruise.get('config_filename') 216 loaded_time = self.cruise.get('loaded_time') 217 self.cruise_loaded_time = datetime.datetime.timestamp(loaded_time) 218 self.active_mode = self.api.get_active_mode() 219 220 # Send updated cruise definition to CDS for console to read. 221 cruise_dict = { 222 'cruise_id': self.cruise.get('id', ''), 223 'filename': self.cruise_filename, 224 'config_timestamp': self.cruise_loaded_time, 225 'loggers': self.loggers, 226 'modes': self.cruise.get('modes', {}), 227 'active_mode': self.active_mode, 228 } 229 logging.info('Sending updated cruise definitions to CDS.') 230 self._write_record_to_data_server( 231 'status:cruise_definition', cruise_dict) 232 except (AttributeError, ValueError, TypeError) as e: 233 logging.info('Failed to update cruise definition: %s', e) 234 235 ############################ 236 def _check_logger_status_loop(self): 237 """Grab logger status message from supervisor and send to cached data 238 server via websocket. Also send cruise mode as separate message. 239 """ 240 while not self.quit_flag: 241 now = time.time() 242 try: 243 config_status = self.supervisor.get_status() 244 with self.config_lock: 245 # Stash status, note time and send update 246 self.config_status = config_status 247 self.status_time = now 248 self._write_record_to_data_server('status:logger_status', config_status) 249 250 # Now get and send cruise mode 251 mode_map = {'active_mode': self.api.get_active_mode()} 252 self._write_record_to_data_server('status:cruise_mode', mode_map) 253 except ValueError as e: 254 logging.warning('Error while trying to send logger status: %s', e) 255 time.sleep(self.interval) 256 257 ############################ 258 def _update_configs_loop(self): 259 """Iteratively check the API for updated configs and send them to the 260 appropriate LoggerRunners. 261 """ 262 while not self.quit_flag: 263 self._update_configs() 264 time.sleep(self.interval) 265 266 ############################ 267 def _update_configs(self): 268 """Get list of new (latest) configs. Send to logger supervisor to make 269 any necessary changes. 270 271 Note: we can't fold this into _update_configs_loop() because we may 272 need to ask the api to call it independently as a callback when it 273 notices that the config has changed. Search for the line: 274 275 api.on_update(callback=logger_manager._update_configs) 276 277 in this file to see where. 278 """ 279 with self.config_lock: 280 # Get new configs in dict {logger:{'configs':[config_name,...]}} 281 logger_configs = self.api.get_logger_configs() 282 if logger_configs: 283 supervisor.update_configs(logger_configs) 284 self.active_configs = logger_configs 285 286 ############################ 287 def _send_cruise_definition_loop(self): 288 """Iteratively assemble information from DB about what loggers should 289 exist and what states they *should* be in. We'll send this to the 290 cached data server whenever it changes (or if it's been a while 291 since we have). 292 293 Also, if the logger or config names have changed, signal that we 294 need to create a new config file for the supervisord process to 295 use. 296 297 Looks like: 298 {'active_mode': 'log', 299 'cruise_id': 'NBP1406', 300 'loggers': {'PCOD': {'active': 'PCOD->file/net', 301 'configs': ['PCOD->off', 302 'PCOD->net', 303 'PCOD->file/net', 304 'PCOD->file/net/db']}, 305 next_logger: next_configs, 306 ... 307 }, 308 'modes': ['off', 'monitor', 'log', 'log+db'] 309 } 310 311 """ 312 last_loaded_timestamp = 0 313 314 while not self.quit_flag: 315 try: 316 self.cruise = self.api.get_configuration() # a Cruise object 317 if not self.cruise: 318 logging.info('No cruise definition found in API') 319 time.sleep(self.interval * 2) 320 continue 321 self.cruise_filename = self.cruise.get('config_filename') 322 loaded_time = self.cruise.get('loaded_time') 323 self.cruise_loaded_time = datetime.datetime.timestamp(loaded_time) 324 325 # Has cruise definition file changed since we loaded it? If so, 326 # send a notification to console so it can ask if user wants to 327 # reload. 328 if self.cruise_filename: 329 try: 330 mtime = os.path.getmtime(self.cruise_filename) 331 if mtime > self.cruise_loaded_time: 332 logging.debug('Cruise file timestamp changed!') 333 self._write_record_to_data_server('status:file_update', mtime) 334 except FileNotFoundError: 335 logging.debug('Cruise file "%s" has disappeared?', self.cruise_filename) 336 337 # Does database have a cruise definition with a newer timestamp? 338 # Means user loaded/reloaded definition. Update our maps to 339 # reflect the new values and send an updated cruise_definition 340 # to the console. 341 if self.cruise_loaded_time > last_loaded_timestamp: 342 last_loaded_timestamp = self.cruise_loaded_time 343 logging.info('New cruise definition detected - rebuilding maps.') 344 self._load_new_definition_from_api() 345 346 except KeyboardInterrupt: # (AttributeError, ValueError, TypeError): 347 logging.warning('No cruise definition found in API') 348 349 # Whether or not we've sent an update, sleep 350 time.sleep(self.interval * 2) 351 352 ############################ 353 def _write_record_to_data_server(self, field_name, record): 354 """Format and label a record and send it to the cached data server. 355 """ 356 if self.data_server_writer: 357 das_record = DASRecord(fields={field_name: record}) 358 logging.debug('DASRecord: %s' % das_record) 359 self.data_server_writer.write(das_record) 360 else: 361 logging.info('Update: %s: %s', field_name, record)
LoggerManager( api, supervisor, data_server_websocket=None, stderr_file_pattern='/var/log/openrvdas/{logger}.stderr', interval=0.25, log_level=<function info>, logger_log_level=30)
60 def __init__(self, 61 api, supervisor, data_server_websocket=None, 62 stderr_file_pattern='/var/log/openrvdas/{logger}.stderr', 63 interval=0.25, log_level=logging.info, logger_log_level=logging.WARNING): 64 """Read desired/current logger configs from Django DB and try to run the 65 loggers specified in those configs. 66 ``` 67 api - ServerAPI (or subclass) instance by which LoggerManager will get 68 its data store updates 69 70 supervisor - a LoggerSupervisor object to use to manage logger 71 processes. 72 73 data_server_websocket - cached data server host:port to which we are 74 going to send our status updates. 75 76 stderr_file_pattern - Pattern into which logger name will be 77 interpolated to create the file path/name to which the 78 logger's stderr will be written. E.g. 79 '/var/log/openrvdas/{logger}.stderr' If 80 data_server_websocket is defined, will write logger 81 stderr to it. 82 83 interval - number of seconds to sleep between checking/updating loggers 84 85 log_level - LoggerManager's log level 86 87 logger_log_level - At what logging level our component loggers 88 should operate. 89 ``` 90 91 """ 92 # Set signal to catch SIGTERM and convert it into a 93 # KeyboardInterrupt so we can shut things down gracefully. 94 try: 95 signal.signal(signal.SIGTERM, kill_handler) 96 except ValueError: 97 logging.warning('LoggerManager not running in main thread; ' 98 'shutting down with Ctl-C may not work.') 99 100 # api class must be subclass of ServerAPI 101 if not issubclass(type(api), ServerAPI): 102 raise ValueError('Passed api "%s" must be subclass of ServerAPI' % api) 103 self.api = api 104 self.supervisor = supervisor 105 106 # Data server to which we're going to send status updates 107 if data_server_websocket: 108 self.data_server_writer = CachedDataWriter(data_server_websocket) 109 else: 110 self.data_server_writer = None 111 112 self.stderr_file_pattern = stderr_file_pattern 113 self.interval = interval 114 self.logger_log_level = logger_log_level 115 116 # Try to set up logging, right off the bat: reset logging to its 117 # freshly-imported state and add handler that also sends logged 118 # messages to the cached data server. 119 reload(logging) 120 logging.basicConfig(format=DEFAULT_LOGGING_FORMAT, level=log_level) 121 122 if self.data_server_writer: 123 cds_writer = ComposedWriter( 124 transforms=ToDASRecordTransform(data_id='stderr', 125 field_name='stderr:logger_manager'), 126 writers=self.data_server_writer) 127 logging.getLogger().addHandler(StdErrLoggingHandler(cds_writer)) 128 129 # How our various loops and threads will know it's time to quit 130 self.quit_flag = False 131 132 # Where we store the latest cruise definition and status reports. 133 self.cruise = None 134 self.cruise_filename = None 135 self.cruise_loaded_time = 0 136 137 self.loggers = {} 138 self.config_to_logger = {} 139 140 self.logger_status = None 141 self.status_time = 0 142 143 # We loop to check the logger status and pass it off to the cached 144 # data server. Do this in a separate thread. 145 self.check_logger_status_thread = None 146 147 # We'll loop to check the API for updates to our desired 148 # configs. Do this in a separate thread. Also keep track of 149 # currently active configs so that we know when an update is 150 # actually needed. 151 self.update_configs_thread = None 152 self.config_lock = threading.Lock() 153 154 self.active_mode = None # which mode is active now? 155 self.active_configs = None # which configs are active now?
Read desired/current logger configs from Django DB and try to run the loggers specified in those configs.
api - ServerAPI (or subclass) instance by which LoggerManager will get
its data store updates
supervisor - a LoggerSupervisor object to use to manage logger
processes.
data_server_websocket - cached data server host:port to which we are
going to send our status updates.
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' If
data_server_websocket is defined, will write logger
stderr to it.
interval - number of seconds to sleep between checking/updating loggers
log_level - LoggerManager's log level
logger_log_level - At what logging level our component loggers
should operate.
def
start(self):
158 def start(self): 159 """Start the threads that make up the LoggerManager operation: 160 161 1. Configuration update loop 162 2. Loop to read logger stderr/status and either output it or 163 transmit it to a cached data server 164 165 Start threads as daemons so that they'll automatically terminate 166 if the main thread does. 167 """ 168 logging.info('Starting LoggerManager') 169 170 # Check logger status in a separate thread. If we've got the 171 # address of a data server websocket, send our updates to it. 172 self.check_logger_status_loop_thread = threading.Thread( 173 name='check_logger_status_loop', 174 target=self._check_logger_status_loop, daemon=True) 175 self.check_logger_status_loop_thread.start() 176 177 # Update configs in a separate thread. 178 self.update_configs_thread = threading.Thread( 179 name='update_configs_loop', 180 target=self._update_configs_loop, daemon=True) 181 self.update_configs_thread.start() 182 183 # Check cruise definition in a separate thread. If we've got the 184 # address of a data server websocket, send our updates to it. 185 self.send_cruise_definition_loop_thread = threading.Thread( 186 name='send_cruise_definition_loop', 187 target=self._send_cruise_definition_loop, daemon=True) 188 self.send_cruise_definition_loop_thread.start()
Start the threads that make up the LoggerManager operation:
- Configuration update loop
- Loop to read logger stderr/status and either output it or transmit it to a cached data server
Start threads as daemons so that they'll automatically terminate if the main thread does.
def
run_data_server( data_server_websocket, data_server_back_seconds, data_server_cleanup_interval, data_server_interval):
366def run_data_server(data_server_websocket, 367 data_server_back_seconds, data_server_cleanup_interval, 368 data_server_interval): 369 """Run a CachedDataServer (to be called as a separate process), 370 accepting websocket connections to receive data to be cached and 371 served. 372 """ 373 # First get the port that we're going to run the data server on. Because 374 # we're running it locally, it should only have a port, not a hostname. 375 # We should try to handle it if they prefix with a ':', though. 376 data_server_websocket = data_server_websocket or DEFAULT_DATA_SERVER_WEBSOCKET 377 websocket_port = int(data_server_websocket.split(':')[-1]) 378 server = CachedDataServer(port=websocket_port, interval=data_server_interval) 379 380 # The server will start serving in its own thread after 381 # initialization, but we need to manually fire up the cleanup loop 382 # if we want it. Maybe we should have this also run automatically in 383 # its own thread after initialization? 384 server.cleanup_loop()
Run a CachedDataServer (to be called as a separate process), accepting websocket connections to receive data to be cached and served.