#!/usr/bin/python
#
# Run a powstream test
#

import atexit
import datetime
import fcntl
import json
import pscheduler
import sys
import os
import signal
import time
import pytz
from powstream_defaults import *
from powstream_utils import get_config, parse_raw_owamp_output, cleanup_dir, cleanup_file, handle_run_error, sleep_or_end, graceful_exit
from subprocess import Popen, PIPE, call

#track when this run starts - make sure it is aware that it is UTC
start_time = datetime.datetime.utcnow().replace(tzinfo=pytz.utc)
#create tag to include in data directory based on current time
#prevents old processes from stomping on new during restart
time_tag = start_time.strftime("%Y%b%dT%H%M%S%f")

#Init logging
log = pscheduler.Log(prefix="tool-powstream", quiet=True)

#DEBUGGING: Set static values below
# task_uuid = 'ABC123'
# participant = 0
# participant_data = [{}, {'ctrl-port': 861}]
# test_spec = {'source': '10.0.1.28', 'dest': '10.0.1.25'}
# duration = pscheduler.iso8601_as_timedelta('PT2M')
# input = { 'schedule': {'until': "2016-09-15T15:22:40Z" }}

#parse JSON input
input = pscheduler.json_load(exit_on_error=True)
try:
    task_uuid = input['task-uuid']
    participant = input['participant']
    participant_data = input['participant-data']
    test_spec = input['test']['spec']
    duration = pscheduler.iso8601_as_timedelta(input['schedule']['duration'])
except KeyError as e:
    pscheduler.fail("Missing required key in run input: %s" % e)
except:
    pscheduler.fail("Error parsing run input: %s" % sys.exc_info()[0])

#check if given start time
if ('schedule' in input) and ('start' in input['schedule']):
    start_time = pscheduler.iso8601_as_datetime(input['schedule']['start'])
    
#set end time - use lesser of now + duration and until time
#NOTE: I don't think we are given until time anymore, leaving this here anyways in case it gets added back
#      Assuming it doesn't we are always given a start and duration, so should be sufficient
end_time = start_time + duration
if ('schedule' in input) and ('until' in input['schedule']):
    until_time = pscheduler.iso8601_as_datetime(input['schedule']['until'])
    if until_time < end_time:
        end_time = until_time
log.debug("Powstream run ends at %s" % end_time)
#determine whether this is a reverse test
flip = test_spec.get('flip', False)

#setup infrastructure for killing processes on normal exit
proc = None
def terminate_proc():
    try:
        if proc:
            proc.terminate()
            log.debug("Terminated powstream process %s" % proc.pid)
    except:
        pass
atexit.register(terminate_proc)

#constants
ADDR_FORMAT = "[%s]:%d"
DEFAULT_POWSTREAM_CMD = '/usr/bin/powstream'
DEFAULT_OWSTATS_CMD = '/usr/bin/owstats'
DEFAULT_PKILL_CMD = '/usr/bin/pkill'
DEFAULT_DATA_DIR = '/var/lib/pscheduler/tool/powstream'
DEFAULT_BUCKET_WIDTH = TIME_SCALE #convert to ms
DEFAULT_RAW_OUTPUT = False #don't display raw packets by default
POWSTREAM_RANGE_ARGS = [
    ('data-ports', '-P'),
]
POWSTREAM_VAL_ARGS = [
    ('ip-tos', '-D'),    
    ('packet-padding', '-s')
]

#read config file
config = get_config()

powstream_cmd = DEFAULT_POWSTREAM_CMD
if config and config.has_option(CONFIG_SECTION, CONFIG_OPT_POWSTREAM_CMD):
    powstream_cmd = config.get(CONFIG_SECTION, CONFIG_OPT_POWSTREAM_CMD)
    
owstats_cmd = DEFAULT_OWSTATS_CMD
if config and config.has_option(CONFIG_SECTION, CONFIG_OPT_OWSTATS_CMD):
    owstats_cmd = config.get(CONFIG_SECTION, CONFIG_OPT_OWSTATS_CMD)

pkill_cmd = DEFAULT_PKILL_CMD
if config and config.has_option(CONFIG_SECTION, CONFIG_OPT_PKILL_CMD):
    pkill_cmd = config.get(CONFIG_SECTION, CONFIG_OPT_PKILL_CMD)
    
keep_data_files = False
if config and config.has_option(CONFIG_SECTION, CONFIG_OPT_KEEP_DATA_FILES):
    keep_data_files = config.getboolean(CONFIG_SECTION, CONFIG_OPT_KEEP_DATA_FILES)

#determine data directory
parent_data_dir = DEFAULT_DATA_DIR
data_dir=None
if config and config.has_option(CONFIG_SECTION, CONFIG_OPT_DATA_DIR):
    parent_data_dir = config.get(CONFIG_SECTION, CONFIG_OPT_DATA_DIR)
if not parent_data_dir.endswith("/"):
    parent_data_dir += "/" 
data_dir = parent_data_dir + task_uuid + '-' + time_tag

#Always print files (-p)
powstream_args = [powstream_cmd, '-p', '-d', data_dir]

#register various handlers that make sure data dir is removed on exit
atexit.register(cleanup_dir, data_dir, keep_data_files=keep_data_files)
cleanup_handler = lambda signum, frame: graceful_exit(data_dir, keep_data_files=keep_data_files, pkill_cmd=pkill_cmd)
signal.signal(signal.SIGTERM, cleanup_handler)

# set log level if needed
if config and config.has_option(CONFIG_SECTION, CONFIG_OPT_LOG_LEVEL):
    log_level = config.get(CONFIG_SECTION, CONFIG_OPT_LOG_LEVEL)
    powstream_args.append('-g')
    powstream_args.append(log_level)
    
#build basic arguments
for arg in POWSTREAM_VAL_ARGS:
    if arg[0] in test_spec:
        powstream_args.append(arg[1])
        powstream_args.append(str(test_spec[arg[0]]))
for rarg in POWSTREAM_RANGE_ARGS:
    if rarg[0] in test_spec:
        powstream_args.append(rarg[1])
        powstream_args.append("%d-%d" % (test_spec[rarg[0]]['lower'], test_spec[rarg[0]]['upper']))
        
#set interval,count and timeout to ensure consistent with duration
count = test_spec.get('packet-count', DEFAULT_PACKET_COUNT)
interval = test_spec.get('packet-interval', DEFAULT_PACKET_INTERVAL)
packet_timeout = test_spec.get('packet-timeout', 0)
powstream_args.append('-c')
powstream_args.append(str(count))
powstream_args.append('-i')
powstream_args.append(str(interval))
if packet_timeout > 0:
    powstream_args.append('-L')
    powstream_args.append(str(packet_timeout))
#calculate min time between results
result_sleep = count * interval + packet_timeout

#set if ipv4 only or ipv6 only
ip_version = str(test_spec.get('ip-version', ''))
if ip_version == '4':
    powstream_args.append('-4')
elif ip_version == '6':
    powstream_args.append('-6')

#bucket width is used for rounding delay values used as buckets for histogram
bucket_width = test_spec.get('bucket-width', DEFAULT_BUCKET_WIDTH)

#determine whether we will return raw packets
raw_output = test_spec.get('output-raw', DEFAULT_RAW_OUTPUT)

#determine control port
control_port = int(test_spec.get('ctrl-port', DEFAULT_OWAMPD_PORT))
    
#finally, set the addresses and packet flow direction
if flip:
    #reverse test
    if 'dest' in test_spec:
        powstream_args.append('-S')
        powstream_args.append(test_spec['dest'])
    powstream_args.append(ADDR_FORMAT % (test_spec['source'], control_port))
else:
    #forward test
    powstream_args.append('-t')
    if 'source' in test_spec:
        powstream_args.append('-S')
        powstream_args.append(test_spec['source'])
    powstream_args.append(ADDR_FORMAT % (test_spec['dest'], control_port))

#Run the process
time.sleep(DEFAULT_CLIENT_SLEEP) #wait for server to boot
powstream_start_attempts = 0
got_result = False
parent_pid = os.getppid()
log.debug("Parent pid is %s" % parent_pid)
while True:
    try:
        #make sure we are not past end time
        if sleep_or_end(0, end_time, parent_pid): break
        
        #cleanup an previous runs - i really don't trust powstream
        try:
            #kill anything referencing our data_dir
            call([pkill_cmd, '-9', '-f', parent_data_dir + task_uuid ], shell=False)
        except:
            #best effort clean-up
            pass
        if os.path.exists(data_dir):
            cleanup_dir(data_dir, keep_data_files=keep_data_files, ignore_errors=True)
        if not os.path.exists(data_dir):
            os.makedirs(data_dir)
        
        #spawn process
        log.debug("Running command: %s" % " ".join(powstream_args))
        proc = Popen(powstream_args, stdout=PIPE, stderr=PIPE, shell=False)
        #update signal handler to make sure process is killed when terminated
        cleanup_handler = lambda signum, frame: graceful_exit(data_dir, keep_data_files=keep_data_files, proc=proc, pkill_cmd=pkill_cmd)
        signal.signal(signal.SIGTERM, cleanup_handler)
        #give it a minute to fire-up
        if sleep_or_end(DEFAULT_CLIENT_SLEEP, end_time, parent_pid):
            break
            
        #check that it is still running
        if proc.poll() is not None:
            powstream_start_attempts += 1
            pow_error = ''
            for line in proc.stderr:
                pow_error += line.strip() + " "
            if powstream_start_attempts > MAX_POWSTREAM_START_ATTEMPTS:
                handle_run_error("powstream returned an error on startup for %d consectutive attempts. Giving up. This is the error it reported: %s" % (powstream_start_attempts, pow_error))
            elif sleep_or_end(DEFAULT_RESTART_SLEEP, end_time, parent_pid):
                handle_run_error("powstream returned an error on startup, end time reached, so giving up" % pow_error)
                break
            else:
                handle_run_error("powstream returned an error on startup, attempting restart after waiting %d seconds: %s" % (DEFAULT_RESTART_SLEEP, pow_error))
                continue
        else:
            powstream_start_attempts = 0
        
        #make stdout non blocking so our program doesn't hang forever if powstream stop reporting
        fd = proc.stdout.fileno()
        fl = fcntl.fcntl(fd, fcntl.F_GETFL)
        fcntl.fcntl(fd, fcntl.F_SETFL, fl | os.O_NONBLOCK)
        
        #sleep until expect result. don't exit if stop earlier so we can check for results
        sleep_or_end(result_sleep, end_time, parent_pid)
        attempts_without_result = 0
        while True:
            #make sure we don't get trapped in this loop forever
            if sleep_or_end(0, end_time, parent_pid): break
            
            #can't use proc.communicate because waits for EOF
            try:
                line = proc.stdout.readline()
            except IOError:
                #throws exception in non-blocking mode if no data, so catch that here
                line = ''
            line = line.strip()
            if line == '' and proc.poll() is not None:
                break
            elif line.endswith('.sum'):
                cleanup_file(line, keep_data_files=keep_data_files)
            elif line.endswith('.owp'):
                got_result = True
                attempts_without_result = 0
                #run owstats to get output
                owstats_args = [owstats_cmd, '-R', line]
                log.debug("Running owstats command: %s" % " ".join(owstats_args))
                try:
                    stats_returncode, stats_stdout, stats_stderr = pscheduler.run_program(owstats_args, timeout=30)
                except OSError as e:
                    handle_run_error("owstats encountered an OS error: %s" % e)
                    cleanup_file(line, keep_data_files=keep_data_files)
                    continue
                except Exception:
                    handle_run_error("owstats failed to complete execution: %s" % sys.exc_info()[0])
                    cleanup_file(line, keep_data_files=keep_data_files)
                    continue
            
                #see if command completed successfully
                log.debug("owstats returned status %s" % stats_returncode)
                if stats_returncode:
                    handle_run_error("owstats completed but returned error: %s" % stats_stderr)
                    cleanup_file(line, keep_data_files=keep_data_files)
                    continue
            
                #no longer need file
                cleanup_file(line, keep_data_files=keep_data_files)
            
                #parse output
                results = parse_raw_owamp_output(stats_stdout, raw_output=raw_output, bucket_width=bucket_width)
            
                #print
                print pscheduler.json_dump(results)
                print pscheduler.api_result_delimiter()
                sys.stdout.flush()
            elif line == '':
                #make sure we don't go too long without a result. 
                # powstream is notorious for this
                attempts_without_result += 1
                if attempts_without_result > DEFAULT_MAX_RETRIES:
                    proc.terminate()
                    break
                
                #sleep until next result
                if got_result and attempts_without_result == 1:
                    #if first blank line since we got result, sleep until next result
                    if sleep_or_end(result_sleep, end_time, parent_pid):
                        break
                else:
                    #...otherwise try again after a shorter sleep
                    if sleep_or_end(DEFAULT_RETRY_SLEEP, end_time, parent_pid):
                        break
        
        #command completed, this should not happen. log and restart
        log.debug("powstream returned status %s" % proc.returncode)
        if proc.returncode:
            owp_error = ''
            for line in proc.stderr:
                owp_error += line.rstrip().lstrip() + " "
            if owp_error == '':
                owp_error = "Nothing on stderr, may have been killed by external command (e.g. something ran 'kill')"
            handle_run_error("powstream exited with error, will attempt restart: %s" % owp_error)
        elif not sleep_or_end(0, end_time, parent_pid):
            #if powstream exited before end time, report that fact
            log.debug("powstream exited with no errors but should have kept running, will attempt restart")
    except OSError as e:
        log.error("powstream encountered an OS error: %s" % e)
        try:
           if proc: proc.terminate()
        except:
            pass
        handle_run_error("The powstream command failed during execution. See server logs for more details.", do_log=False)
    except Exception:
        log.error("powstream failed to complete execution: %s" % sys.exc_info()[0])
        try:
           if proc: proc.terminate()
        except:
            pass
        handle_run_error("The powstream command failed during execution. See server logs for more details.", do_log=False)

pscheduler.succeed()
