#!/usr/bin/python
#
# arch-tag: interface to node power controller
# Time-stamp: <2005-10-19 14:11:29 mike>
#
import threading
import struct
import time
try:
    import timeoutsocket
    timeoutsocket.setDefaultSocketTimeout(5)
    commError = (timeoutsocket.Timeout, timeoutsocket.Error)
except:
    import socket
    commError = socket.error
import xmlrpclib
from NPCData import adc_offset, dio_offset

__version__ = '1.5'

class NPC(object):
    def __init__(self, host):
        """Initialize the XMLRPC connection to the node power control
        system and determine the format of the raw data records.

        arguments:
            host      -  hostname of the node power controller (npc)
        """
        self.server = xmlrpclib.ServerProxy("http://%s:8080/RPC2" % host)
        sample = self.server.getSample().data
        self.record = list(struct.unpack('>lll', sample[0:12]))
        self._fmt = '>' + ('4B32h' * self.record[2])
        if len(sample) != (self.record[2]*68 + 12):
            sample = self.server.getSample().data
        self.record.extend(list(struct.unpack(self._fmt, sample[12:])))
        self.interval = 1
        self.running = 0
        self.monitor = None
        self.host = host
        self.listeners = []
        self.logger = None
        self.hostup = 0
        self.methods = self.server.system.listMethods()
        
    def __getstate__(self):
        d = {}
        d.update(self.__dict__)
        d['server'] = None
        d['logger'] = None
        return d

    def __setstate__(self, d):
        self.__dict__.update(d)
        self.server = xmlrpclib.ServerProxy("http://%s:8080/RPC2" % self.host)
        sample = self.server.getSample().data
        self.record = list(struct.unpack('>lll', sample[0:12]))
        self.record.extend(list(struct.unpack(self._fmt, sample[12:])))
        self.methods = self.server.system.listMethods()
        
    def _action(self):
        record = None
        try:
            t = time.time()
            sample = self.server.getSample().data
            record = list(struct.unpack('>lll', sample[0:12]))
            record.extend(struct.unpack(self._fmt, sample[12:]))
            record[0] = int(t)
            record[1] = int((t - record[0])*1e6)
            self.hostup = 1
        except commError, e:
            if self.logger:
                self.logger("Communication error, trying to reconnect to NPC")
            self.hostup = 0
            self.server = xmlrpclib.ServerProxy("http://%s:8080/RPC2" % self.host)
        if record:
            self._lock.acquire()
            self.record = record
            self._lock.release()
        
    def __call__(self):
        """Thread function.  Executes the _action method on a periodic
        schedule using the specified interval.
        """
        listeners = self.listeners
        action = self._action
        inc = self.interval
        self.running = 1
        t0 = time.time()
        while self.running:
            t = time.time()
            action()
            for l in listeners:
                l.recvSample(self.record, self.hostup)
            t0 += inc
            if t0 < t:
                t0 = t
            time.sleep(t0-t)

    def register(self, obj):
        """Add an object to the listeners list.  Each listener must implement
        a method named recvSample.  On each poll, this method will be called
        with the sample record as an argument.
        """
        try:
            if callable(getattr(obj, 'recvSample')):
                self.listeners.append(obj)
        except AttributeError:
            pass

    def set_logger(self, logger=None):
        self.logger = logger
        
    def getSample(self):
        """Return the most recent sample record"""
        if not self.running:
            return self.record
        self._lock.acquire()
        rval = self.record
        self._lock.release()
        return rval

    def data_format(self):
        """Return the struct.unpack format string for the data
        portion of the sample record.
        """
        return self._fmt
    
    def start_poll(self, interval=1):
        """Start a thread to poll the NPC at the specified interval
        (in seconds).
        """
        self.interval = interval
        self._lock = threading.Lock()
        self.monitor = threading.Thread(target=self)
        self.monitor.setDaemon(1)
        self.monitor.start()

    def stop_poll(self):
        self.running = 0
        self.monitor.join(timeout=self.interval*2)
        self._lock = None
        self.monitor = None
        
    def is_running(self):
        return self.running

    def sample_time(self):
        """Return the timestamp of the most recent sample as a two-element
        list; [seconds, microseconds].
        """
        return self.record[0:2]

    def _rpc_dispatch(self, name, *args):
        """Intercept calls to the XML-RPC server so we can catch
        timeout errors.
        """
        ## TODO: intercept calls to setADmax and cache the values
        ## so they can be re-sent to the NPC if communication is
        ## lost.
        f = getattr(self.server, name)
        try:
            return f(*args)
        except commError:
            if self.logger:
                self.logger("Comm error in %s" % name)
            return None
        
    def __getattr__(self, name):
        """Dispatch remote function calls to the Server Proxy"""
        return lambda *args: self._rpc_dispatch(name, *args)

    
