controller.py
1256 lines
| 36.2 KiB
| text/x-python
|
PythonLexer
/ schainpy / controller.py
|
r197 | ''' | ||
|
r1171 | Updated on January , 2018, for multiprocessing purposes | ||
Author: Sergio Cortez | ||||
|
r197 | Created on September , 2012 | ||
''' | ||||
|
r1171 | from platform import python_version | ||
|
r672 | import sys | ||
|
r514 | import ast | ||
|
r687 | import datetime | ||
|
r672 | import traceback | ||
|
r898 | import math | ||
|
r931 | import time | ||
|
r1171 | import zmq | ||
|
r1052 | from multiprocessing import Process, cpu_count | ||
|
r1177 | from threading import Thread | ||
|
r681 | from xml.etree.ElementTree import ElementTree, Element, SubElement, tostring | ||
from xml.dom import minidom | ||||
|
r1171 | |||
r1129 | from schainpy.admin import Alarm, SchainWarning | |||
|
r1191 | from schainpy.model import * | ||
|
r1052 | from schainpy.utils import log | ||
|
r1191 | |||
|
r1052 | |||
DTYPES = { | ||||
'Voltage': '.r', | ||||
'Spectra': '.pdata' | ||||
} | ||||
|
r1082 | |||
|
r1052 | def MPProject(project, n=cpu_count()): | ||
''' | ||||
Project wrapper to run schain in n processes | ||||
''' | ||||
rconf = project.getReadUnitObj() | ||||
op = rconf.getOperationObj('run') | ||||
dt1 = op.getParameterValue('startDate') | ||||
dt2 = op.getParameterValue('endDate') | ||||
|
r1112 | tm1 = op.getParameterValue('startTime') | ||
tm2 = op.getParameterValue('endTime') | ||||
r892 | days = (dt2 - dt1).days | |||
|
r1082 | |||
for day in range(days + 1): | ||||
r892 | skip = 0 | |||
cursor = 0 | ||||
processes = [] | ||||
|
r1052 | dt = dt1 + datetime.timedelta(day) | ||
dt_str = dt.strftime('%Y/%m/%d') | ||||
reader = JRODataReader() | ||||
paths, files = reader.searchFilesOffLine(path=rconf.path, | ||||
|
r1082 | startDate=dt, | ||
endDate=dt, | ||||
|
r1112 | startTime=tm1, | ||
endTime=tm2, | ||||
|
r1082 | ext=DTYPES[rconf.datatype]) | ||
|
r1052 | nFiles = len(files) | ||
if nFiles == 0: | ||||
|
r999 | continue | ||
|
r1112 | skip = int(math.ceil(nFiles / n)) | ||
|
r1082 | while nFiles > cursor * skip: | ||
rconf.update(startDate=dt_str, endDate=dt_str, cursor=cursor, | ||||
skip=skip) | ||||
p = project.clone() | ||||
|
r1052 | p.start() | ||
processes.append(p) | ||||
|
r924 | cursor += 1 | ||
def beforeExit(exctype, value, trace): | ||||
r892 | for process in processes: | |||
process.terminate() | ||||
process.join() | ||||
|
r1167 | print(traceback.print_tb(trace)) | ||
|
r957 | |||
|
r924 | sys.excepthook = beforeExit | ||
r892 | ||||
for process in processes: | ||||
process.join() | ||||
|
r924 | process.terminate() | ||
|
r958 | |||
|
r931 | time.sleep(3) | ||
r892 | ||||
|
r1177 | def wait(context): | ||
time.sleep(1) | ||||
c = zmq.Context() | ||||
receiver = c.socket(zmq.SUB) | ||||
receiver.connect('ipc:///tmp/schain_{}_pub'.format(self.id)) | ||||
receiver.setsockopt(zmq.SUBSCRIBE, self.id.encode()) | ||||
msg = receiver.recv_multipart()[1] | ||||
context.terminate() | ||||
|
r197 | class ParameterConf(): | ||
r889 | ||||
|
r197 | id = None | ||
name = None | ||||
value = None | ||||
|
r199 | format = None | ||
r889 | ||||
|
r529 | __formated_value = None | ||
r889 | ||||
|
r197 | ELEMENTNAME = 'Parameter' | ||
r889 | ||||
|
r197 | def __init__(self): | ||
r889 | ||||
|
r199 | self.format = 'str' | ||
r889 | ||||
|
r197 | def getElementName(self): | ||
r889 | ||||
|
r197 | return self.ELEMENTNAME | ||
r889 | ||||
|
r197 | def getValue(self): | ||
|
r600 | |||
value = self.value | ||||
format = self.format | ||||
r889 | ||||
|
r529 | if self.__formated_value != None: | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
r892 | if format == 'obj': | |||
return value | ||||
|
r600 | if format == 'str': | ||
|
r596 | self.__formated_value = str(value) | ||
return self.__formated_value | ||||
r889 | ||||
|
r596 | if value == '': | ||
|
r1167 | raise ValueError('%s: This parameter value is empty' % self.name) | ||
r889 | ||||
|
r600 | if format == 'list': | ||
|
r1193 | strList = [s.strip() for s in value.split(',')] | ||
|
r529 | self.__formated_value = strList | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
|
r600 | if format == 'intlist': | ||
|
r1052 | ''' | ||
|
r535 | Example: | ||
value = (0,1,2) | ||||
|
r1052 | ''' | ||
r889 | ||||
|
r735 | new_value = ast.literal_eval(value) | ||
r889 | ||||
|
r735 | if type(new_value) not in (tuple, list): | ||
new_value = [int(new_value)] | ||||
r889 | ||||
|
r735 | self.__formated_value = new_value | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
|
r600 | if format == 'floatlist': | ||
|
r1052 | ''' | ||
|
r535 | Example: | ||
value = (0.5, 1.4, 2.7) | ||||
|
r1052 | ''' | ||
r889 | ||||
|
r735 | new_value = ast.literal_eval(value) | ||
r889 | ||||
|
r735 | if type(new_value) not in (tuple, list): | ||
new_value = [float(new_value)] | ||||
r889 | ||||
|
r741 | self.__formated_value = new_value | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
|
r600 | if format == 'date': | ||
|
r224 | strList = value.split('/') | ||
|
r197 | intList = [int(x) for x in strList] | ||
date = datetime.date(intList[0], intList[1], intList[2]) | ||||
r889 | ||||
|
r529 | self.__formated_value = date | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
|
r600 | if format == 'time': | ||
|
r224 | strList = value.split(':') | ||
|
r197 | intList = [int(x) for x in strList] | ||
time = datetime.time(intList[0], intList[1], intList[2]) | ||||
r889 | ||||
|
r529 | self.__formated_value = time | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
|
r600 | if format == 'pairslist': | ||
|
r1052 | ''' | ||
|
r226 | Example: | ||
value = (0,1),(1,2) | ||||
|
r1052 | ''' | ||
|
r226 | |||
|
r735 | new_value = ast.literal_eval(value) | ||
r889 | ||||
|
r735 | if type(new_value) not in (tuple, list): | ||
|
r1167 | raise ValueError('%s has to be a tuple or list of pairs' % value) | ||
r889 | ||||
|
r735 | if type(new_value[0]) not in (tuple, list): | ||
if len(new_value) != 2: | ||||
|
r1167 | raise ValueError('%s has to be a tuple or list of pairs' % value) | ||
|
r735 | new_value = [new_value] | ||
r889 | ||||
|
r735 | for thisPair in new_value: | ||
if len(thisPair) != 2: | ||||
|
r1167 | raise ValueError('%s has to be a tuple or list of pairs' % value) | ||
r889 | ||||
|
r735 | self.__formated_value = new_value | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
|
r600 | if format == 'multilist': | ||
|
r1052 | ''' | ||
|
r514 | Example: | ||
value = (0,1,2),(3,4,5) | ||||
|
r1052 | ''' | ||
|
r514 | multiList = ast.literal_eval(value) | ||
r889 | ||||
|
r600 | if type(multiList[0]) == int: | ||
|
r1052 | multiList = ast.literal_eval('(' + value + ')') | ||
r889 | ||||
|
r529 | self.__formated_value = multiList | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
r889 | ||||
|
r677 | if format == 'bool': | ||
value = int(value) | ||||
r889 | ||||
|
r677 | if format == 'int': | ||
value = float(value) | ||||
r889 | ||||
|
r600 | format_func = eval(format) | ||
r889 | ||||
|
r529 | self.__formated_value = format_func(value) | ||
r889 | ||||
|
r529 | return self.__formated_value | ||
|
r596 | |||
def updateId(self, new_id): | ||||
r889 | ||||
|
r596 | self.id = str(new_id) | ||
r889 | ||||
|
r199 | def setup(self, id, name, value, format='str'): | ||
|
r596 | self.id = str(id) | ||
|
r197 | self.name = name | ||
r892 | if format == 'obj': | |||
self.value = value | ||||
else: | ||||
self.value = str(value) | ||||
|
r535 | self.format = str.lower(format) | ||
r889 | ||||
|
r735 | self.getValue() | ||
r889 | ||||
|
r643 | return 1 | ||
r889 | ||||
|
r577 | def update(self, name, value, format='str'): | ||
r889 | ||||
|
r577 | self.name = name | ||
self.value = str(value) | ||||
self.format = format | ||||
r889 | ||||
|
r197 | def makeXml(self, opElement): | ||
|
r898 | if self.name not in ('queue',): | ||
parmElement = SubElement(opElement, self.ELEMENTNAME) | ||||
parmElement.set('id', str(self.id)) | ||||
parmElement.set('name', self.name) | ||||
parmElement.set('value', self.value) | ||||
parmElement.set('format', self.format) | ||||
|
r1184 | |||
|
r197 | def readXml(self, parmElement): | ||
r889 | ||||
|
r197 | self.id = parmElement.get('id') | ||
self.name = parmElement.get('name') | ||||
self.value = parmElement.get('value') | ||||
|
r568 | self.format = str.lower(parmElement.get('format')) | ||
r889 | ||||
|
r1082 | # Compatible with old signal chain version | ||
|
r568 | if self.format == 'int' and self.name == 'idfigure': | ||
self.name = 'id' | ||||
r889 | ||||
|
r197 | def printattr(self): | ||
r889 | ||||
|
r1184 | print('Parameter[%s]: name = %s, value = %s, format = %s, project_id = %s' % (self.id, self.name, self.value, self.format, self.project_id)) | ||
|
r1082 | |||
|
r1052 | class OperationConf(): | ||
r889 | ||||
|
r197 | ELEMENTNAME = 'Operation' | ||
r889 | ||||
|
r197 | def __init__(self): | ||
r889 | ||||
|
r589 | self.id = '0' | ||
|
r568 | self.name = None | ||
self.priority = None | ||||
|
r1171 | self.topic = None | ||
r889 | ||||
|
r197 | def __getNewId(self): | ||
r889 | ||||
|
r1082 | return int(self.id) * 10 + len(self.parmConfObjList) + 1 | ||
|
r596 | |||
|
r1171 | def getId(self): | ||
return self.id | ||||
|
r596 | def updateId(self, new_id): | ||
r889 | ||||
|
r596 | self.id = str(new_id) | ||
r889 | ||||
|
r596 | n = 1 | ||
for parmObj in self.parmConfObjList: | ||||
r889 | ||||
|
r1082 | idParm = str(int(new_id) * 10 + n) | ||
|
r596 | parmObj.updateId(idParm) | ||
r889 | ||||
|
r596 | n += 1 | ||
r889 | ||||
|
r197 | def getElementName(self): | ||
r889 | ||||
|
r197 | return self.ELEMENTNAME | ||
r889 | ||||
|
r197 | def getParameterObjList(self): | ||
r889 | ||||
|
r197 | return self.parmConfObjList | ||
r889 | ||||
|
r577 | def getParameterObj(self, parameterName): | ||
r889 | ||||
|
r577 | for parmConfObj in self.parmConfObjList: | ||
r889 | ||||
|
r577 | if parmConfObj.name != parameterName: | ||
continue | ||||
r889 | ||||
|
r577 | return parmConfObj | ||
r889 | ||||
|
r577 | return None | ||
|
r606 | def getParameterObjfromValue(self, parameterValue): | ||
r889 | ||||
|
r577 | for parmConfObj in self.parmConfObjList: | ||
r889 | ||||
|
r577 | if parmConfObj.getValue() != parameterValue: | ||
continue | ||||
r889 | ||||
|
r577 | return parmConfObj.getValue() | ||
r889 | ||||
|
r577 | return None | ||
r889 | ||||
|
r577 | def getParameterValue(self, parameterName): | ||
r889 | ||||
|
r577 | parameterObj = self.getParameterObj(parameterName) | ||
|
r1082 | |||
|
r1004 | # if not parameterObj: | ||
|
r1011 | # return None | ||
|
r1082 | |||
|
r577 | value = parameterObj.getValue() | ||
r889 | ||||
|
r577 | return value | ||
r889 | ||||
def getKwargs(self): | ||||
kwargs = {} | ||||
for parmConfObj in self.parmConfObjList: | ||||
if self.name == 'run' and parmConfObj.name == 'datatype': | ||||
continue | ||||
kwargs[parmConfObj.name] = parmConfObj.getValue() | ||||
return kwargs | ||||
|
r1177 | def setup(self, id, name, priority, type, project_id): | ||
r889 | ||||
|
r596 | self.id = str(id) | ||
|
r1177 | self.project_id = project_id | ||
|
r197 | self.name = name | ||
self.type = type | ||||
self.priority = priority | ||||
self.parmConfObjList = [] | ||||
r889 | ||||
|
r577 | def removeParameters(self): | ||
r889 | ||||
|
r577 | for obj in self.parmConfObjList: | ||
del obj | ||||
r889 | ||||
|
r577 | self.parmConfObjList = [] | ||
r889 | ||||
|
r199 | def addParameter(self, name, value, format='str'): | ||
|
r1082 | |||
|
r1052 | if value is None: | ||
return None | ||||
|
r197 | id = self.__getNewId() | ||
r889 | ||||
|
r197 | parmConfObj = ParameterConf() | ||
|
r643 | if not parmConfObj.setup(id, name, value, format): | ||
return None | ||||
r889 | ||||
|
r197 | self.parmConfObjList.append(parmConfObj) | ||
r889 | ||||
|
r197 | return parmConfObj | ||
r889 | ||||
|
r577 | def changeParameter(self, name, value, format='str'): | ||
r889 | ||||
|
r577 | parmConfObj = self.getParameterObj(name) | ||
parmConfObj.update(name, value, format) | ||||
r889 | ||||
|
r577 | return parmConfObj | ||
r889 | ||||
|
r681 | def makeXml(self, procUnitElement): | ||
r889 | ||||
|
r681 | opElement = SubElement(procUnitElement, self.ELEMENTNAME) | ||
|
r197 | opElement.set('id', str(self.id)) | ||
opElement.set('name', self.name) | ||||
opElement.set('type', self.type) | ||||
opElement.set('priority', str(self.priority)) | ||||
r889 | ||||
|
r197 | for parmConfObj in self.parmConfObjList: | ||
parmConfObj.makeXml(opElement) | ||||
r889 | ||||
|
r1184 | def readXml(self, opElement, project_id): | ||
r889 | ||||
|
r197 | self.id = opElement.get('id') | ||
self.name = opElement.get('name') | ||||
self.type = opElement.get('type') | ||||
self.priority = opElement.get('priority') | ||||
|
r1192 | self.project_id = str(project_id) | ||
r889 | ||||
|
r1082 | # Compatible with old signal chain version | ||
# Use of 'run' method instead 'init' | ||||
|
r568 | if self.type == 'self' and self.name == 'init': | ||
self.name = 'run' | ||||
r889 | ||||
|
r197 | self.parmConfObjList = [] | ||
r889 | ||||
|
r860 | parmElementList = opElement.iter(ParameterConf().getElementName()) | ||
r889 | ||||
|
r197 | for parmElement in parmElementList: | ||
parmConfObj = ParameterConf() | ||||
parmConfObj.readXml(parmElement) | ||||
r889 | ||||
|
r1082 | # Compatible with old signal chain version | ||
# If an 'plot' OPERATION is found, changes name operation by the value of its type PARAMETER | ||||
|
r568 | if self.type != 'self' and self.name == 'Plot': | ||
if parmConfObj.format == 'str' and parmConfObj.name == 'type': | ||||
self.name = parmConfObj.value | ||||
continue | ||||
r889 | ||||
|
r568 | self.parmConfObjList.append(parmConfObj) | ||
r889 | ||||
|
r197 | def printattr(self): | ||
r889 | ||||
|
r1184 | print('%s[%s]: name = %s, type = %s, priority = %s, project_id = %s' % (self.ELEMENTNAME, | ||
|
r1082 | self.id, | ||
self.name, | ||||
self.type, | ||||
|
r1184 | self.priority, | ||
self.project_id)) | ||||
r889 | ||||
|
r197 | for parmConfObj in self.parmConfObjList: | ||
parmConfObj.printattr() | ||||
r889 | ||||
|
r1171 | def createObject(self): | ||
r889 | ||||
|
r1171 | className = eval(self.name) | ||
|
r1184 | |||
|
r1177 | if self.type == 'other': | ||
opObj = className() | ||||
elif self.type == 'external': | ||||
kwargs = self.getKwargs() | ||||
opObj = className(self.id, self.project_id, **kwargs) | ||||
opObj.start() | ||||
r889 | ||||
|
r197 | return opObj | ||
r889 | ||||
|
r197 | class ProcUnitConf(): | ||
r889 | ||||
|
r197 | ELEMENTNAME = 'ProcUnit' | ||
r889 | ||||
|
r197 | def __init__(self): | ||
r889 | ||||
|
r197 | self.id = None | ||
|
r199 | self.datatype = None | ||
|
r197 | self.name = None | ||
|
r1171 | self.inputId = None | ||
|
r197 | self.opConfObjList = [] | ||
self.procUnitObj = None | ||||
self.opObjDict = {} | ||||
r889 | ||||
|
r197 | def __getPriority(self): | ||
r889 | ||||
|
r1082 | return len(self.opConfObjList) + 1 | ||
r889 | ||||
|
r197 | def __getNewId(self): | ||
r889 | ||||
|
r1082 | return int(self.id) * 10 + len(self.opConfObjList) + 1 | ||
r889 | ||||
|
r197 | def getElementName(self): | ||
r889 | ||||
|
r197 | return self.ELEMENTNAME | ||
r889 | ||||
|
r197 | def getId(self): | ||
r889 | ||||
|
r589 | return self.id | ||
|
r596 | |||
|
r1177 | def updateId(self, new_id): | ||
|
r1171 | ''' | ||
|
r1082 | new_id = int(parentId) * 10 + (int(self.id) % 10) | ||
new_inputId = int(parentId) * 10 + (int(self.inputId) % 10) | ||||
r889 | ||||
|
r1082 | # If this proc unit has not inputs | ||
|
r1171 | #if self.inputId == '0': | ||
#new_inputId = 0 | ||||
r889 | ||||
|
r596 | n = 1 | ||
for opConfObj in self.opConfObjList: | ||||
r889 | ||||
|
r1082 | idOp = str(int(new_id) * 10 + n) | ||
|
r596 | opConfObj.updateId(idOp) | ||
r889 | ||||
|
r596 | n += 1 | ||
r889 | ||||
|
r596 | self.parentId = str(parentId) | ||
self.id = str(new_id) | ||||
|
r1171 | #self.inputId = str(new_inputId) | ||
''' | ||||
n = 1 | ||||
|
r1177 | |||
|
r197 | def getInputId(self): | ||
r889 | ||||
|
r589 | return self.inputId | ||
r889 | ||||
|
r197 | def getOperationObjList(self): | ||
r889 | ||||
|
r197 | return self.opConfObjList | ||
r889 | ||||
|
r577 | def getOperationObj(self, name=None): | ||
r889 | ||||
|
r577 | for opConfObj in self.opConfObjList: | ||
r889 | ||||
|
r577 | if opConfObj.name != name: | ||
continue | ||||
r889 | ||||
|
r577 | return opConfObj | ||
r889 | ||||
|
r577 | return None | ||
r889 | ||||
|
r606 | def getOpObjfromParamValue(self, value=None): | ||
r889 | ||||
|
r577 | for opConfObj in self.opConfObjList: | ||
if opConfObj.getParameterObjfromValue(parameterValue=value) != value: | ||||
continue | ||||
return opConfObj | ||||
return None | ||||
r889 | ||||
|
r197 | def getProcUnitObj(self): | ||
r889 | ||||
|
r197 | return self.procUnitObj | ||
r889 | ||||
|
r1177 | def setup(self, project_id, id, name, datatype, inputId): | ||
|
r1171 | ''' | ||
id sera el topico a publicar | ||||
inputId sera el topico a subscribirse | ||||
''' | ||||
|
r1082 | # Compatible with old signal chain version | ||
if datatype == None and name == None: | ||||
|
r1167 | raise ValueError('datatype or name should be defined') | ||
r889 | ||||
|
r1171 | #Definir una condicion para inputId cuando sea 0 | ||
|
r1082 | if name == None: | ||
|
r596 | if 'Proc' in datatype: | ||
name = datatype | ||||
else: | ||||
|
r1082 | name = '%sProc' % (datatype) | ||
r889 | ||||
|
r1082 | if datatype == None: | ||
datatype = name.replace('Proc', '') | ||||
r889 | ||||
|
r596 | self.id = str(id) | ||
|
r1177 | self.project_id = project_id | ||
|
r197 | self.name = name | ||
|
r199 | self.datatype = datatype | ||
|
r1171 | self.inputId = inputId | ||
|
r197 | self.opConfObjList = [] | ||
r889 | ||||
|
r1171 | self.addOperation(name='run', optype='self') | ||
r889 | ||||
|
r577 | def removeOperations(self): | ||
r889 | ||||
|
r577 | for obj in self.opConfObjList: | ||
del obj | ||||
r889 | ||||
|
r577 | self.opConfObjList = [] | ||
self.addOperation(name='run') | ||||
r889 | ||||
|
r219 | def addParameter(self, **kwargs): | ||
|
r580 | ''' | ||
|
r1052 | Add parameters to 'run' operation | ||
|
r580 | ''' | ||
|
r219 | opObj = self.opConfObjList[0] | ||
r889 | ||||
|
r219 | opObj.addParameter(**kwargs) | ||
r889 | ||||
|
r219 | return opObj | ||
r889 | ||||
|
r1177 | def addOperation(self, name, optype='self'): | ||
|
r1171 | ''' | ||
Actualizacion - > proceso comunicacion | ||||
En el caso de optype='self', elminar. DEfinir comuncacion IPC -> Topic | ||||
definir el tipoc de socket o comunicacion ipc++ | ||||
''' | ||||
r889 | ||||
|
r197 | id = self.__getNewId() | ||
|
r1171 | priority = self.__getPriority() # Sin mucho sentido, pero puede usarse | ||
|
r197 | opConfObj = OperationConf() | ||
|
r1177 | opConfObj.setup(id, name=name, priority=priority, type=optype, project_id=self.project_id) | ||
|
r197 | self.opConfObjList.append(opConfObj) | ||
r889 | ||||
|
r197 | return opConfObj | ||
r889 | ||||
|
r681 | def makeXml(self, projectElement): | ||
r889 | ||||
|
r681 | procUnitElement = SubElement(projectElement, self.ELEMENTNAME) | ||
procUnitElement.set('id', str(self.id)) | ||||
procUnitElement.set('name', self.name) | ||||
procUnitElement.set('datatype', self.datatype) | ||||
procUnitElement.set('inputId', str(self.inputId)) | ||||
r889 | ||||
|
r197 | for opConfObj in self.opConfObjList: | ||
|
r681 | opConfObj.makeXml(procUnitElement) | ||
r889 | ||||
|
r1184 | def readXml(self, upElement, project_id): | ||
r889 | ||||
|
r197 | self.id = upElement.get('id') | ||
self.name = upElement.get('name') | ||||
|
r199 | self.datatype = upElement.get('datatype') | ||
|
r197 | self.inputId = upElement.get('inputId') | ||
|
r1184 | self.project_id = str(project_id) | ||
r889 | ||||
|
r1052 | if self.ELEMENTNAME == 'ReadUnit': | ||
self.datatype = self.datatype.replace('Reader', '') | ||||
r889 | ||||
|
r1052 | if self.ELEMENTNAME == 'ProcUnit': | ||
self.datatype = self.datatype.replace('Proc', '') | ||||
r889 | ||||
|
r589 | if self.inputId == 'None': | ||
self.inputId = '0' | ||||
r889 | ||||
|
r197 | self.opConfObjList = [] | ||
r889 | ||||
|
r860 | opElementList = upElement.iter(OperationConf().getElementName()) | ||
r889 | ||||
|
r197 | for opElement in opElementList: | ||
opConfObj = OperationConf() | ||||
|
r1184 | opConfObj.readXml(opElement, project_id) | ||
|
r197 | self.opConfObjList.append(opConfObj) | ||
r889 | ||||
|
r197 | def printattr(self): | ||
r889 | ||||
|
r1184 | print('%s[%s]: name = %s, datatype = %s, inputId = %s, project_id = %s' % (self.ELEMENTNAME, | ||
|
r1082 | self.id, | ||
self.name, | ||||
self.datatype, | ||||
|
r1184 | self.inputId, | ||
self.project_id)) | ||||
|
r1082 | |||
|
r197 | for opConfObj in self.opConfObjList: | ||
opConfObj.printattr() | ||||
r889 | ||||
def getKwargs(self): | ||||
opObj = self.opConfObjList[0] | ||||
kwargs = opObj.getKwargs() | ||||
return kwargs | ||||
|
r1177 | def createObjects(self): | ||
|
r1171 | ''' | ||
|
r1177 | Instancia de unidades de procesamiento. | ||
|
r1171 | ''' | ||
|
r1192 | |||
|
r197 | className = eval(self.name) | ||
r889 | kwargs = self.getKwargs() | |||
|
r1177 | procUnitObj = className(self.id, self.inputId, self.project_id, **kwargs) # necesitan saber su id y su entrada por fines de ipc | ||
log.success('creating process...', self.name) | ||||
r889 | ||||
|
r924 | for opConfObj in self.opConfObjList: | ||
|
r1177 | |||
if opConfObj.type == 'self' and opConfObj.name == 'run': | ||||
|
r906 | continue | ||
|
r1082 | elif opConfObj.type == 'self': | ||
|
r1177 | opObj = getattr(procUnitObj, opConfObj.name) | ||
else: | ||||
opObj = opConfObj.createObject() | ||||
|
r1171 | |||
|
r1177 | log.success('creating operation: {}, type:{}'.format( | ||
opConfObj.name, | ||||
opConfObj.type), self.name) | ||||
|
r1171 | |||
|
r1177 | procUnitObj.addOperation(opConfObj, opObj) | ||
|
r1171 | |||
procUnitObj.start() | ||||
|
r197 | self.procUnitObj = procUnitObj | ||
|
r1171 | |||
|
r573 | def close(self): | ||
r889 | ||||
|
r577 | for opConfObj in self.opConfObjList: | ||
if opConfObj.type == 'self': | ||||
continue | ||||
r889 | ||||
|
r577 | opObj = self.procUnitObj.getOperationObj(opConfObj.id) | ||
opObj.close() | ||||
r889 | ||||
|
r577 | self.procUnitObj.close() | ||
r889 | ||||
|
r573 | return | ||
r889 | ||||
|
r1082 | |||
|
r197 | class ReadUnitConf(ProcUnitConf): | ||
r889 | ||||
|
r197 | ELEMENTNAME = 'ReadUnit' | ||
r889 | ||||
|
r197 | def __init__(self): | ||
r889 | ||||
|
r197 | self.id = None | ||
|
r199 | self.datatype = None | ||
|
r197 | self.name = None | ||
|
r589 | self.inputId = None | ||
|
r197 | self.opConfObjList = [] | ||
r889 | ||||
|
r197 | def getElementName(self): | ||
r889 | ||||
|
r1171 | return self.ELEMENTNAME | ||
|
r1177 | def setup(self, project_id, id, name, datatype, path='', startDate='', endDate='', | ||
startTime='', endTime='', server=None, **kwargs): | ||||
|
r1052 | |||
|
r1171 | |||
''' | ||||
*****el id del proceso sera el Topico | ||||
Adicion de {topic}, si no esta presente -> error | ||||
kwargs deben ser trasmitidos en la instanciacion | ||||
''' | ||||
|
r1082 | # Compatible with old signal chain version | ||
if datatype == None and name == None: | ||||
|
r1167 | raise ValueError('datatype or name should be defined') | ||
|
r1088 | if name == None: | ||
|
r596 | if 'Reader' in datatype: | ||
name = datatype | ||||
|
r1088 | datatype = name.replace('Reader','') | ||
|
r596 | else: | ||
|
r1088 | name = '{}Reader'.format(datatype) | ||
if datatype == None: | ||||
if 'Reader' in name: | ||||
datatype = name.replace('Reader','') | ||||
else: | ||||
datatype = name | ||||
name = '{}Reader'.format(name) | ||||
r889 | ||||
|
r197 | self.id = id | ||
|
r1177 | self.project_id = project_id | ||
|
r197 | self.name = name | ||
|
r199 | self.datatype = datatype | ||
|
r963 | if path != '': | ||
self.path = os.path.abspath(path) | ||||
|
r197 | self.startDate = startDate | ||
self.endDate = endDate | ||||
self.startTime = startTime | ||||
self.endTime = endTime | ||||
|
r963 | self.server = server | ||
|
r224 | self.addRunOperation(**kwargs) | ||
r889 | ||||
|
r1052 | def update(self, **kwargs): | ||
|
r589 | |||
|
r1052 | if 'datatype' in kwargs: | ||
datatype = kwargs.pop('datatype') | ||||
|
r589 | if 'Reader' in datatype: | ||
|
r1052 | self.name = datatype | ||
|
r589 | else: | ||
|
r1082 | self.name = '%sReader' % (datatype) | ||
|
r1052 | self.datatype = self.name.replace('Reader', '') | ||
r889 | ||||
|
r1082 | attrs = ('path', 'startDate', 'endDate', | ||
|
r1177 | 'startTime', 'endTime') | ||
|
r1082 | |||
|
r1052 | for attr in attrs: | ||
if attr in kwargs: | ||||
setattr(self, attr, kwargs.pop(attr)) | ||||
|
r1082 | |||
|
r577 | self.updateRunOperation(**kwargs) | ||
r889 | ||||
|
r687 | def removeOperations(self): | ||
r889 | ||||
|
r687 | for obj in self.opConfObjList: | ||
del obj | ||||
r889 | ||||
|
r687 | self.opConfObjList = [] | ||
r889 | ||||
|
r1171 | def addRunOperation(self, **kwargs): | ||
r889 | ||||
|
r1171 | opObj = self.addOperation(name='run', optype='self') | ||
r889 | ||||
|
r963 | if self.server is None: | ||
|
r1082 | opObj.addParameter( | ||
name='datatype', value=self.datatype, format='str') | ||||
|
r1052 | opObj.addParameter(name='path', value=self.path, format='str') | ||
|
r1082 | opObj.addParameter( | ||
name='startDate', value=self.startDate, format='date') | ||||
opObj.addParameter( | ||||
name='endDate', value=self.endDate, format='date') | ||||
opObj.addParameter( | ||||
name='startTime', value=self.startTime, format='time') | ||||
opObj.addParameter( | ||||
name='endTime', value=self.endTime, format='time') | ||||
|
r1167 | for key, value in list(kwargs.items()): | ||
|
r1082 | opObj.addParameter(name=key, value=value, | ||
format=type(value).__name__) | ||||
|
r963 | else: | ||
|
r1082 | opObj.addParameter(name='server', value=self.server, format='str') | ||
r889 | ||||
|
r197 | return opObj | ||
r889 | ||||
|
r577 | def updateRunOperation(self, **kwargs): | ||
r889 | ||||
|
r1052 | opObj = self.getOperationObj(name='run') | ||
|
r577 | opObj.removeParameters() | ||
r889 | ||||
|
r1052 | opObj.addParameter(name='datatype', value=self.datatype, format='str') | ||
opObj.addParameter(name='path', value=self.path, format='str') | ||||
|
r1082 | opObj.addParameter( | ||
name='startDate', value=self.startDate, format='date') | ||||
|
r1052 | opObj.addParameter(name='endDate', value=self.endDate, format='date') | ||
|
r1082 | opObj.addParameter( | ||
name='startTime', value=self.startTime, format='time') | ||||
|
r1052 | opObj.addParameter(name='endTime', value=self.endTime, format='time') | ||
|
r1082 | |||
|
r1167 | for key, value in list(kwargs.items()): | ||
|
r1082 | opObj.addParameter(name=key, value=value, | ||
format=type(value).__name__) | ||||
r889 | ||||
|
r577 | return opObj | ||
|
r1082 | |||
|
r1184 | def readXml(self, upElement, project_id): | ||
r889 | ||||
|
r681 | self.id = upElement.get('id') | ||
self.name = upElement.get('name') | ||||
self.datatype = upElement.get('datatype') | ||||
|
r1184 | self.project_id = str(project_id) #yong | ||
r889 | ||||
|
r1052 | if self.ELEMENTNAME == 'ReadUnit': | ||
self.datatype = self.datatype.replace('Reader', '') | ||||
r889 | ||||
|
r681 | self.opConfObjList = [] | ||
r889 | ||||
|
r860 | opElementList = upElement.iter(OperationConf().getElementName()) | ||
r889 | ||||
|
r681 | for opElement in opElementList: | ||
opConfObj = OperationConf() | ||||
|
r1184 | opConfObj.readXml(opElement, project_id) | ||
|
r681 | self.opConfObjList.append(opConfObj) | ||
r889 | ||||
|
r681 | if opConfObj.name == 'run': | ||
self.path = opConfObj.getParameterValue('path') | ||||
self.startDate = opConfObj.getParameterValue('startDate') | ||||
self.endDate = opConfObj.getParameterValue('endDate') | ||||
self.startTime = opConfObj.getParameterValue('startTime') | ||||
self.endTime = opConfObj.getParameterValue('endTime') | ||||
r889 | ||||
|
r1082 | |||
|
r1040 | class Project(Process): | ||
|
r1052 | |||
|
r224 | ELEMENTNAME = 'Project' | ||
r889 | ||||
|
r1171 | def __init__(self): | ||
|
r1052 | |||
|
r1040 | Process.__init__(self) | ||
|
r1171 | self.id = None | ||
|
r1177 | self.filename = None | ||
|
r197 | self.description = None | ||
|
r1126 | self.email = None | ||
r1130 | self.alarm = None | |||
|
r197 | self.procUnitConfObjDict = {} | ||
|
r1052 | |||
|
r197 | def __getNewId(self): | ||
r889 | ||||
|
r1167 | idList = list(self.procUnitConfObjDict.keys()) | ||
|
r1082 | id = int(self.id) * 10 | ||
r889 | ||||
|
r764 | while True: | ||
id += 1 | ||||
r889 | ||||
|
r764 | if str(id) in idList: | ||
continue | ||||
r889 | ||||
|
r764 | break | ||
r889 | ||||
|
r197 | return str(id) | ||
r889 | ||||
|
r197 | def getElementName(self): | ||
r889 | ||||
|
r197 | return self.ELEMENTNAME | ||
|
r577 | |||
def getId(self): | ||||
r889 | ||||
|
r577 | return self.id | ||
r889 | ||||
|
r596 | def updateId(self, new_id): | ||
r889 | ||||
|
r596 | self.id = str(new_id) | ||
r889 | ||||
|
r1167 | keyList = list(self.procUnitConfObjDict.keys()) | ||
|
r596 | keyList.sort() | ||
r889 | ||||
|
r596 | n = 1 | ||
newProcUnitConfObjDict = {} | ||||
r889 | ||||
|
r596 | for procKey in keyList: | ||
r889 | ||||
|
r596 | procUnitConfObj = self.procUnitConfObjDict[procKey] | ||
|
r1082 | idProcUnit = str(int(self.id) * 10 + n) | ||
|
r1177 | procUnitConfObj.updateId(idProcUnit) | ||
|
r596 | newProcUnitConfObjDict[idProcUnit] = procUnitConfObj | ||
n += 1 | ||||
r889 | ||||
|
r596 | self.procUnitConfObjDict = newProcUnitConfObjDict | ||
r889 | ||||
|
r1177 | def setup(self, id=1, name='', description='', email=None, alarm=[]): | ||
r889 | ||||
|
r1171 | print(' ') | ||
|
r1167 | print('*' * 60) | ||
|
r1171 | print('* Starting SIGNAL CHAIN PROCESSING (Multiprocessing) v%s *' % schainpy.__version__) | ||
|
r1167 | print('*' * 60) | ||
|
r1171 | print("* Python " + python_version() + " *") | ||
print('*' * 19) | ||||
print(' ') | ||||
|
r596 | self.id = str(id) | ||
|
r1171 | self.description = description | ||
|
r1126 | self.email = email | ||
self.alarm = alarm | ||||
|
r577 | |||
|
r1126 | def update(self, **kwargs): | ||
r889 | ||||
|
r1167 | for key, value in list(kwargs.items()): | ||
|
r1126 | setattr(self, key, value) | ||
r889 | ||||
|
r1052 | def clone(self): | ||
p = Project() | ||||
p.procUnitConfObjDict = self.procUnitConfObjDict | ||||
return p | ||||
|
r687 | def addReadUnit(self, id=None, datatype=None, name=None, **kwargs): | ||
|
r1052 | |||
|
r1171 | ''' | ||
Actualizacion: | ||||
Se agrego un nuevo argumento: topic -relativo a la forma de comunicar los procesos simultaneos | ||||
* El id del proceso sera el topico al que se deben subscribir los procUnits para recibir la informacion(data) | ||||
''' | ||||
|
r687 | if id is None: | ||
idReadUnit = self.__getNewId() | ||||
else: | ||||
idReadUnit = str(id) | ||||
r889 | ||||
|
r197 | readUnitConfObj = ReadUnitConf() | ||
|
r1177 | readUnitConfObj.setup(self.id, idReadUnit, name, datatype, **kwargs) | ||
|
r197 | self.procUnitConfObjDict[readUnitConfObj.getId()] = readUnitConfObj | ||
|
r1171 | |||
|
r197 | return readUnitConfObj | ||
r889 | ||||
|
r589 | def addProcUnit(self, inputId='0', datatype=None, name=None): | ||
r889 | ||||
|
r1171 | ''' | ||
Actualizacion: | ||||
Se agrego dos nuevos argumentos: topic_read (lee data de otro procUnit) y topic_write(escribe o envia data a otro procUnit) | ||||
Deberia reemplazar a "inputId" | ||||
** A fin de mantener el inputID, este sera la representaacion del topicoal que deben subscribirse. El ID propio de la intancia | ||||
(proceso) sera el topico de la publicacion, todo sera asignado de manera dinamica. | ||||
''' | ||||
idProcUnit = self.__getNewId() #Topico para subscripcion | ||||
|
r197 | procUnitConfObj = ProcUnitConf() | ||
|
r1177 | procUnitConfObj.setup(self.id, idProcUnit, name, datatype, inputId) #topic_read, topic_write, | ||
|
r197 | self.procUnitConfObjDict[procUnitConfObj.getId()] = procUnitConfObj | ||
r889 | ||||
|
r197 | return procUnitConfObj | ||
r889 | ||||
|
r587 | def removeProcUnit(self, id): | ||
r889 | ||||
|
r1167 | if id in list(self.procUnitConfObjDict.keys()): | ||
|
r587 | self.procUnitConfObjDict.pop(id) | ||
r889 | ||||
|
r577 | def getReadUnitId(self): | ||
r889 | ||||
|
r577 | readUnitConfObj = self.getReadUnitObj() | ||
r889 | ||||
|
r577 | return readUnitConfObj.id | ||
r889 | ||||
|
r577 | def getReadUnitObj(self): | ||
r889 | ||||
|
r1167 | for obj in list(self.procUnitConfObjDict.values()): | ||
|
r1052 | if obj.getElementName() == 'ReadUnit': | ||
|
r577 | return obj | ||
r889 | ||||
|
r577 | return None | ||
r889 | ||||
|
r606 | def getProcUnitObj(self, id=None, name=None): | ||
r889 | ||||
|
r606 | if id != None: | ||
return self.procUnitConfObjDict[id] | ||||
r889 | ||||
|
r606 | if name != None: | ||
return self.getProcUnitObjByName(name) | ||||
r889 | ||||
|
r606 | return None | ||
r889 | ||||
|
r581 | def getProcUnitObjByName(self, name): | ||
r889 | ||||
|
r1167 | for obj in list(self.procUnitConfObjDict.values()): | ||
|
r581 | if obj.name == name: | ||
return obj | ||||
r889 | ||||
|
r581 | return None | ||
r889 | ||||
|
r606 | def procUnitItems(self): | ||
r889 | ||||
|
r1167 | return list(self.procUnitConfObjDict.items()) | ||
r889 | ||||
def makeXml(self): | ||||
|
r224 | projectElement = Element('Project') | ||
|
r197 | projectElement.set('id', str(self.id)) | ||
projectElement.set('name', self.name) | ||||
projectElement.set('description', self.description) | ||||
r889 | ||||
|
r1167 | for procUnitConfObj in list(self.procUnitConfObjDict.values()): | ||
|
r197 | procUnitConfObj.makeXml(projectElement) | ||
r889 | ||||
|
r197 | self.projectElement = projectElement | ||
r889 | ||||
|
r702 | def writeXml(self, filename=None): | ||
r889 | ||||
|
r702 | if filename == None: | ||
|
r708 | if self.filename: | ||
filename = self.filename | ||||
else: | ||||
|
r1052 | filename = 'schain.xml' | ||
r889 | ||||
|
r702 | if not filename: | ||
|
r1167 | print('filename has not been defined. Use setFilename(filename) for do it.') | ||
|
r702 | return 0 | ||
r889 | ||||
|
r687 | abs_file = os.path.abspath(filename) | ||
r889 | ||||
|
r687 | if not os.access(os.path.dirname(abs_file), os.W_OK): | ||
|
r1167 | print('No write permission on %s' % os.path.dirname(abs_file)) | ||
|
r681 | return 0 | ||
r889 | ||||
|
r687 | if os.path.isfile(abs_file) and not(os.access(abs_file, os.W_OK)): | ||
|
r1167 | print('File %s already exists and it could not be overwriten' % abs_file) | ||
|
r681 | return 0 | ||
r889 | ||||
|
r681 | self.makeXml() | ||
r889 | ||||
|
r687 | ElementTree(self.projectElement).write(abs_file, method='xml') | ||
r889 | ||||
|
r708 | self.filename = abs_file | ||
r889 | ||||
|
r681 | return 1 | ||
|
r197 | |||
|
r1082 | def readXml(self, filename=None): | ||
r889 | ||||
|
r830 | if not filename: | ||
|
r1167 | print('filename is not defined') | ||
|
r830 | return 0 | ||
r889 | ||||
|
r687 | abs_file = os.path.abspath(filename) | ||
r889 | ||||
|
r687 | if not os.path.isfile(abs_file): | ||
|
r1167 | print('%s file does not exist' % abs_file) | ||
|
r681 | return 0 | ||
r889 | ||||
|
r197 | self.projectElement = None | ||
self.procUnitConfObjDict = {} | ||||
r889 | ||||
|
r735 | try: | ||
self.projectElement = ElementTree().parse(abs_file) | ||||
except: | ||||
|
r1167 | print('Error reading %s, verify file format' % filename) | ||
|
r735 | return 0 | ||
r889 | ||||
|
r197 | self.project = self.projectElement.tag | ||
r889 | ||||
|
r197 | self.id = self.projectElement.get('id') | ||
self.name = self.projectElement.get('name') | ||||
|
r1082 | self.description = self.projectElement.get('description') | ||
readUnitElementList = self.projectElement.iter( | ||||
ReadUnitConf().getElementName()) | ||||
r889 | ||||
|
r197 | for readUnitElement in readUnitElementList: | ||
readUnitConfObj = ReadUnitConf() | ||||
|
r1184 | readUnitConfObj.readXml(readUnitElement, self.id) | ||
|
r197 | self.procUnitConfObjDict[readUnitConfObj.getId()] = readUnitConfObj | ||
r889 | ||||
|
r1082 | procUnitElementList = self.projectElement.iter( | ||
ProcUnitConf().getElementName()) | ||||
r889 | ||||
|
r197 | for procUnitElement in procUnitElementList: | ||
procUnitConfObj = ProcUnitConf() | ||||
|
r1184 | procUnitConfObj.readXml(procUnitElement, self.id) | ||
|
r197 | self.procUnitConfObjDict[procUnitConfObj.getId()] = procUnitConfObj | ||
r889 | ||||
|
r708 | self.filename = abs_file | ||
r889 | ||||
|
r681 | return 1 | ||
r889 | ||||
|
r1177 | def __str__(self): | ||
r889 | ||||
|
r1184 | print('Project[%s]: name = %s, description = %s, project_id = %s' % (self.id, | ||
|
r1082 | self.name, | ||
|
r1184 | self.description, | ||
self.project_id)) | ||||
|
r1082 | |||
|
r1177 | for procUnitConfObj in self.procUnitConfObjDict.values(): | ||
print(procUnitConfObj) | ||||
r889 | ||||
|
r197 | def createObjects(self): | ||
r889 | ||||
|
r1192 | |||
keys = list(self.procUnitConfObjDict.keys()) | ||||
keys.sort() | ||||
for key in keys: | ||||
self.procUnitConfObjDict[key].createObjects() | ||||
r889 | ||||
r1130 | def __handleError(self, procUnitConfObj, modes=None, stdout=True): | |||
r889 | ||||
|
r687 | import socket | ||
r889 | ||||
r1129 | if modes is None: | |||
modes = self.alarm | ||||
|
r1163 | |||
if not self.alarm: | ||||
modes = [] | ||||
r1129 | ||||
|
r687 | err = traceback.format_exception(sys.exc_info()[0], | ||
|
r1082 | sys.exc_info()[1], | ||
sys.exc_info()[2]) | ||||
|
r1171 | |||
r1128 | log.error('{}'.format(err[-1]), procUnitConfObj.name) | |||
|
r735 | |||
|
r1052 | message = ''.join(err) | ||
r889 | ||||
r1130 | if stdout: | |||
sys.stderr.write(message) | ||||
r889 | ||||
|
r1082 | subject = 'SChain v%s: Error running %s\n' % ( | ||
schainpy.__version__, procUnitConfObj.name) | ||||
r889 | ||||
|
r1082 | subtitle = '%s: %s\n' % ( | ||
procUnitConfObj.getElementName(), procUnitConfObj.name) | ||||
subtitle += 'Hostname: %s\n' % socket.gethostbyname( | ||||
socket.gethostname()) | ||||
subtitle += 'Working directory: %s\n' % os.path.abspath('./') | ||||
subtitle += 'Configuration file: %s\n' % self.filename | ||||
subtitle += 'Time: %s\n' % str(datetime.datetime.now()) | ||||
r889 | ||||
|
r687 | readUnitConfObj = self.getReadUnitObj() | ||
if readUnitConfObj: | ||||
|
r1052 | subtitle += '\nInput parameters:\n' | ||
|
r1082 | subtitle += '[Data path = %s]\n' % readUnitConfObj.path | ||
subtitle += '[Data type = %s]\n' % readUnitConfObj.datatype | ||||
subtitle += '[Start date = %s]\n' % readUnitConfObj.startDate | ||||
subtitle += '[End date = %s]\n' % readUnitConfObj.endDate | ||||
subtitle += '[Start time = %s]\n' % readUnitConfObj.startTime | ||||
subtitle += '[End time = %s]\n' % readUnitConfObj.endTime | ||||
r889 | ||||
r1129 | a = Alarm( | |||
modes=modes, | ||||
|
r1126 | email=self.email, | ||
message=message, | ||||
subject=subject, | ||||
subtitle=subtitle, | ||||
filename=self.filename | ||||
) | ||||
|
r1082 | |||
r1130 | return a | |||
r1129 | ||||
|
r672 | def isPaused(self): | ||
return 0 | ||||
r889 | ||||
|
r672 | def isStopped(self): | ||
return 0 | ||||
r889 | ||||
|
r672 | def runController(self): | ||
|
r1052 | ''' | ||
|
r672 | returns 0 when this process has been stopped, 1 otherwise | ||
|
r1052 | ''' | ||
r889 | ||||
|
r672 | if self.isPaused(): | ||
|
r1167 | print('Process suspended') | ||
r889 | ||||
|
r672 | while True: | ||
|
r1052 | time.sleep(0.1) | ||
r889 | ||||
|
r672 | if not self.isPaused(): | ||
break | ||||
r889 | ||||
|
r672 | if self.isStopped(): | ||
break | ||||
r889 | ||||
|
r1167 | print('Process reinitialized') | ||
r889 | ||||
|
r672 | if self.isStopped(): | ||
|
r1167 | print('Process stopped') | ||
|
r672 | return 0 | ||
r889 | ||||
|
r672 | return 1 | ||
|
r691 | |||
def setFilename(self, filename): | ||||
r889 | ||||
|
r691 | self.filename = filename | ||
r889 | ||||
|
r1171 | def setProxyCom(self): | ||
|
r1177 | |||
if not os.path.exists('/tmp/schain'): | ||||
os.mkdir('/tmp/schain') | ||||
|
r1171 | |||
|
r1177 | self.ctx = zmq.Context() | ||
xpub = self.ctx.socket(zmq.XPUB) | ||||
xpub.bind('ipc:///tmp/schain/{}_pub'.format(self.id)) | ||||
xsub = self.ctx.socket(zmq.XSUB) | ||||
xsub.bind('ipc:///tmp/schain/{}_sub'.format(self.id)) | ||||
|
r1171 | |||
|
r1177 | try: | ||
zmq.proxy(xpub, xsub) | ||||
|
r1187 | except: # zmq.ContextTerminated: | ||
|
r1177 | xpub.close() | ||
xsub.close() | ||||
r889 | ||||
|
r1052 | def run(self): | ||
|
r1004 | |||
|
r1177 | log.success('Starting {}: {}'.format(self.name, self.id), tag='') | ||
|
r1112 | self.start_time = time.time() | ||
|
r1177 | self.createObjects() | ||
# t = Thread(target=wait, args=(self.ctx, )) | ||||
# t.start() | ||||
|
r1171 | self.setProxyCom() | ||
|
r1177 | |||
|
r1171 | # Iniciar todos los procesos .start(), monitoreo de procesos. ELiminar lo de abajo | ||
r889 | ||||
|
r1187 | log.success('{} Done (time: {}s)'.format( | ||
|
r1112 | self.name, | ||
|
r1177 | time.time()-self.start_time)) | ||