#!/usr/bin/python3

#Needs an interface to delete old data associated with a config-spec, so clients
#can update their formats and purge incompatible data during upgrades
#May also need a way to export data

import argparse
import hashlib
import logging
import logging.handlers
import os
import signal
import sys
import time
import traceback

import libstorward.config
import libstorward.service
import libstorward.system

_logger = logging.getLogger('')

def _termhandler(signum, frame):
    if not libstorward.system.ALIVE:
        _logger.critical("A second TERM signal was received; exiting immediately")
        sys.exit(3)
    libstorward.system.ALIVE = False
    _logger.warning("A TERM signal was received; beginning graceful shutdown")

def _setup_logging(config):
    logging.root.setLevel(logging.DEBUG)
    
    if config.logging_level == libstorward.config.LOGGING_DEBUG.upper():
        formatter = logging.Formatter(
         "%(asctime)s : %(levelname)s : %(name)s:%(lineno)d[%(threadName)s] : %(message)s"
        )
    else:
        formatter = logging.Formatter(
         "%(asctime)s : %(levelname)s : %(name)s : %(message)s"
        )
        
    if config.logging_level: #Disk-based logging is desired
        if logging.root.handlers:
            _logger.info("Configuring file-based logging for {file}...".format(
                file=config.logging_file,
            ))
        file_logger = logging.handlers.RotatingFileHandler(
            config.logging_file,
            maxBytes=config.logging_file_max_size,
            backupCount=config.logging_file_history_count,
        )
        if logging.root.handlers:
            _logger.info("Configured file-based logging with rotation every {count} bytes and history of {filecount}".format(
                count=config.logging_file_max_size,
                filecount=config.logging_file_history_count,
            ))
        file_logger.setLevel(getattr(logging, config.logging_level))
        file_logger.setFormatter(formatter)
        logging.root.addHandler(file_logger)
        _logger.info("File-based logging online")
        
if __name__ == '__main__':
    parser = argparse.ArgumentParser(description='Store-and-forward daemon for Intelligent Endpoints')
    parser.add_argument('--config', type=str,
        default='/etc/3d-p/iep-storward/iep-storward.json',
        help='The path from which to read configuration settings',
    )
    args = parser.parse_args()
    
    config = libstorward.config.load_config(args.config)
    del args
    
    _setup_logging(config)
    _logger.debug("Configuration loaded")
    
    _logger.debug("Installing SIGTERM handler...")
    signal.signal(signal.SIGTERM, _termhandler)
    
    _logger.debug("Preparing bandwidth-splitter...")
    libstorward.system.initialise_bandwidth_splitter(config.bandwidth_details)
    
    _logger.debug("Preparing internal input channel...")
    libstorward.system.initialise_internal()
    
    _logger.debug("Preparing TCP input channel...")
    libstorward.system.initialise_tcp(config.communication_details_tcp)
    
    _logger.debug("Connecting to MQTT...")
    libstorward.system.initialise_mqtt(config.communication_details_mqtt)
    
    _logger.debug("Preparing geofence monitor...")
    libstorward.system.initialise_geofence()
    
    _logger.debug("Preparing service engines...")
    input_work_list = libstorward.service.InputWorkList()
    service_engines = []
    for service_spec in config.services: #Reads from the service-path specified in the config file and yields an object for each
        try:
            unique_id = '{name}-{path_hash}'.format(
                name=service_spec.name,
                path_hash=hashlib.md5(service_spec.path.encode('utf-8')).hexdigest(),
            )
            service_engine = libstorward.service.ServiceEngine(
                service_spec=service_spec,
                service_id=unique_id,
                workpath=os.path.join(config.buffering_filesystem_path, unique_id),
                input_list_callback=input_work_list.callback,
            )
        except Exception as e:
            _logger.error("Unable to initialise service defined at {path}:\n{error}".format(
                path=service_spec.path,
                error=traceback.format_exc(),
            ))
            continue
        else:
            _logger.info("Prepared '{name}' with ID '{unique_id}'".format(
                name=service_spec.name,
                unique_id=unique_id,
            ))
            service_engines.append(service_engine)
    _logger.info("Prepared {count} service engines".format(
        count=len(service_engines),
    ))
    
    _logger.debug("Preparing worker threads...")
    thread_count = config.performance_details.get('worker_threads')
    if thread_count is None:
        thread_count = 3 #Having fewer than two general-purpose threads could let one expensive service block everything; one will be slaved to inputs
        thread_count += int(len(service_engines) / 10) #Arbitrary density
    thread_count = min(int(len(service_engines) * 2), max(3, thread_count)) #No reason to have more than one thread per input/output
    if thread_count < (config.performance_details.get('worker_threads') or 0):
        _logger.info("Requested {requested} worker-threads, but some would always be idle".format(
            requested=config.performance_details.get('worker_threads'),
        ))
    _logger.info("Preparing {count} worker threads...".format(
        count=thread_count,
    ))
    service_engine_output_work_pool = libstorward.service.ServiceEngineOutputWorkPool(service_engines)
    service_threads = []
    for i in range(thread_count - 1):
        service_thread = libstorward.service.WorkThread(service_engine_output_work_pool, input_work_list)
        service_threads.append(service_thread)
        service_thread.start()
    #Slave one thread to only handle inputs so starvation doesn't occur due to network anomalies
    service_thread = libstorward.service.WorkThread(service_engine_output_work_pool, input_work_list, input_only=True)
    service_threads.append(service_thread)
    service_thread.start()
    del thread_count
    del service_engines
    del service_thread
    
    del config
    _logger.info("Beginning normal operation...")
    try:
        while libstorward.system.ALIVE:
            #TODO: if dynamic worker-thread spawning/culling is needed, implement it here
            time.sleep(1)
    except KeyboardInterrupt:
        _logger.warning("System shutdown requested by operator")
    except Exception as e:
        _logger.critical("System shutdown triggered by unknown failure:\n{error}".format(
            error=traceback.format_exc(),
        ))
    finally:
        libstorward.system.ALIVE = False
        
        _logger.debug("Ensuring all service engines have stopped...")
        for service_thread in service_threads:
            service_thread.join()
        _logger.debug("All service engines have stopped")
        
