R"""A producer/consumer example

   Author          : Anders Andersen
\\ Created On      : Tue Apr 21 10:16:51 1999
\\ Last Modified By: 
\\ Last Modified On: Mon Jul 05 23:23:43 1999
\\ Status          : Unknown, Use with caution!

Copyright {\copyright} 1999 Lancaster University, UK and NORUT
Information Technology Ltd., Norway.  See COPYING for details.

An management example with a producer, a buffer and a consumer.

"""


from lbind import *
from sigbind import *
from component import *
from fc2 import FC2
from amcomp import *
from encaps import *

import thread
import time
import random


class PC:
    R"""A generic class for producers and consumers

    A generic class implementing the common features of producers and
    consumers.  It provides a thread safe solution to start and stop a
    thread that runs a loop executing its task (the {\aacodefont
    \_do\_it} method implements the task).  The {\aacodefont
    \_resched} method is used to schedule next iteration of the loop
    (default sleeping a random period).

    """
    
    def __init__(self):
        self.ctrl = IRef(self, ["start", "stop", "running"], [])
        self.lock = thread.allocate_lock()
        self._running = 0

    def _run(self):
        while 1:
            self.lock.acquire()
            if self._running:
                self.lock.release()
                self._do_it()
            else:
                self.lock.release()
                break
            self._resched()

    def _resched(self):
        time.sleep(random.random())

    def start(self):
        self.lock.acquire()
        if not self._running:
            self._running = 1
            self.lock.release()
            thread.start_new_thread(self._run, ())
        else:
            self.lock.release()

    def stop(self):
        self.lock.acquire()
        if self._running:
            self._running = 0
            self.lock.release()
        else:
            self.lock.release()

    def running(self):
        return self._running

    def _do_it(self):
        pass


class Producer(PC):
    R"""A producer

    A producer with a control ({\aacodefont ctrl}), a request
    ({\aacodefont req}) and a reply ({\aacodefont rep}) interface.
    When the producer is running it produces sequence numbered request
    that are sent through the {\aacodefont req} interface.  Replays
    are received through the {\aacodefont req} interface and managed
    by the {\aacodefont result} method.  The producer will resend
    request where not replay are received.

    """
    
    def __init__(self):
        self.val = 0
        self.requests = []
        self.resend = []
        self.req = IRef(self, [], ["add_request", "add_request_first"])
        self.rep = IRef(self, ["result"], [])
        PC.__init__(self)

    def _do_it(self):
        self.val = self.val + 1
        self.requests.append(self.val)
        print "P:", self.val
        self.req.add_request(self.val)

    def result(self, val):
        for i in range(len(self.requests)):
            if self.requests[i] < val:
                if self.requests[i] in self.resend:
                    for j in range(len(self.resend)):
                        if self.requests[i] == self.resend[j]:
                            del self.resend[j]
                            break
                else:
                    self.resend.append(self.requests[i])
                    print "P (resend):", self.requests[i]
                    self.req.add_request_first(self.requests[i])
            elif self.requests[i] == val:
                del self.requests[i]
                break
            else:
                break


class Consumer(PC):
    R"""A consumer

    A consumer with a consumer and a request interface.

    """

    def __init__(self):
        self.req = IRef(self, [], ["get_request"])
        self.rep = IRef(self, [], ["result"])
        self.val = 0
        PC.__init__(self)

    def _do_it(self):
        val = self.req.get_request()
        print "C:", val
        self.rep.result(val)
        


class Buffer:
    R"""A buffer

    A buffer between a producer and a consumer.

    """

    def __init__(self, size):
        self.inn = IRef(self, ["add_request", "add_request_first"], [])
        self.out = IRef(self, ["get_request"], [])
        self.overflow = SigSrcIRef()
        self.empty = SigSrcIRef()
        self.size = size + 1
        self.buf = [""] * self.size
        self.bgn = 0
        self.end = 0
        self.lock = thread.allocate_lock()

    def add_request(self, item):
        self.lock.acquire()
        self.end = (self.end + 1) % self.size
        if self.end == self.bgn:
            self.end = (self.end - 1) % self.size
            self.overflow.event()
        else:
            self.buf[self.end] = item
        print "B:", item
        self.lock.release()

    def add_request_first(self, item):
        self.lock.acquire()
        if ((self.bgn - 1) % self.size) == self.end:
            self.overflow.event()
        else:
            self.buf[self.bgn] = item
            self.bgn = (self.bgn - 1) % self.size
        print "B:", item
        self.lock.release()

    def get_request(self):
        self.lock.acquire()
        if self.end == self.bgn:
            self.empty.event()
            item = None
        else:
            self.bgn = (self.bgn + 1) % self.size
            item = self.buf[self.bgn]
        self.lock.release()
        return item


def ipresched(self):
    time.sleep(random.random()/2)

class SAIP:
    def __init__(self, obj):
        self.obj = obj
        self.ip = IRef(self, ["event"], [])
        self.mp = IRef(self, [], ["event"])
        self.encaps = None
    def event(self):
        if not self.encaps:
            self.encaps = encapsulation(self.obj)
        self.encaps.addMethod("_resched", ipresched, 1)
        print "SAIP done ip"
        thread.start_new_thread(self._sendmp, ())
    def _sendmp(self):
        time.sleep(0.5)
        try:
            self.mp.event()
        except AttributeError:
            pass


class LRRResched:
    def __init__(self, fac):
        self.fac = fac
    def __call__(self, *arg):
        time.sleep(random.random()*self.fac)

class SALRR:
    def __init__(self, obj):
        self.fac = 1
        self.obj = obj
        self.lrr = IRef(self, ["event"], [])
        self.encaps = None
    def event(self):
        self.fac = self.fac * 2
        if not self.encaps:
            self.encaps = encapsulation(self.obj)
        self.encaps.addPreMethod("_resched", LRRResched(self.fac))
        print "SALRR done lrr"


# Create the producer, the consumer and the buffer
prod = componentFactory(["ctrl", "req", "rep"], Producer)
cons = componentFactory(["ctrl", "req", "rep"], Consumer)
buff = componentFactory(["inn", "out", "overflow", "empty"], Buffer, (5,))

# Create bindings between producer-buffer-consumer
localBind(prod.interfaces["req"], buff.interfaces["inn"])
localBind(cons.interfaces["req"], buff.interfaces["out"])
localBind(cons.interfaces["rep"], prod.interfaces["rep"])

# Create control interfaces for the producer and the consumer
cprod = IRef(None, [], ["start", "stop", "running"])
ccons = IRef(None, [], ["start", "stop", "running"])

# Create bindgs from the control interfaces to the producer and the consumer
localBind(cprod, prod.interfaces["ctrl"])
localBind(ccons, cons.interfaces["ctrl"])

# Monitor
moni = AmComp(FC2(open("moni.fc2")).fc2py)
localBind(buff.interfaces["overflow"], moni.interfaces["a_in_overflow"])

# Strategy selector
ssel = AmComp(FC2(open("ssel.fc2")).fc2py)
localBind(moni.interfaces["a_out_oto"], ssel.interfaces["a_in_oto"])

# Strategy activator IP
saip = componentFactory(["ip", "mp"], SAIP, (cons.object,))
localBind(ssel.interfaces["a_out_ip"], saip.interfaces["ip"])
localBind(saip.interfaces["mp"], ssel.interfaces["a_in_mp"])

# Strategy activator LRR
salrr = componentFactory(["lrr"], SALRR, (prod.object,))
localBind(ssel.interfaces["a_out_lrr"], salrr.interfaces["lrr"])

# Catch event
#evnt = Event("LRR")

# Start the automata
moni.run()
ssel.run()


#
# Run a test
#

# Test sequence spec (prod [on/off], cons [on/off], sleep [sec])
seq = [(1,0,4),(1,1,4),(1,0,2),(1,1,4),(0,0,0)]

# The test loop
for (p, c, s) in seq:

    # Start or stop producer
    if p:
        if not cprod.running():
            cprod.start()
    else:
        if cprod.running():
            cprod.stop()

    # Start or stop consumer
    if c:
        if not ccons.running():
            ccons.start()
    else:
        if ccons.running():
            ccons.stop()

    # Wait for next action
    time.sleep(s)
