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.
api
supervisor
stderr_file_pattern
interval
logger_log_level
quit_flag
cruise
cruise_filename
cruise_loaded_time
loggers
config_to_logger
logger_status
status_time
check_logger_status_thread
update_configs_thread
config_lock
active_mode
active_configs
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:

  1. Configuration update loop
  2. 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 quit(self):
191    def quit(self):
192        """Exit the loop and shut down all loggers."""
193        self.quit_flag = True

Exit the loop and shut down all loggers.

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.