2017-03-09 17:39:44 +00:00
|
|
|
import logging
|
2017-09-15 13:44:20 +01:00
|
|
|
from datetime import datetime
|
2017-03-09 17:39:44 +00:00
|
|
|
|
2017-03-17 15:57:05 +00:00
|
|
|
from wa.framework import pluginloader, signal
|
2017-03-20 16:24:22 +00:00
|
|
|
from wa.framework.configuration.core import Status
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
|
|
|
|
class Job(object):
|
|
|
|
|
2017-07-04 15:44:05 +01:00
|
|
|
_workload_cache = {}
|
|
|
|
|
2017-03-09 17:39:44 +00:00
|
|
|
@property
|
|
|
|
def id(self):
|
|
|
|
return self.spec.id
|
|
|
|
|
|
|
|
@property
|
2017-03-16 17:54:48 +00:00
|
|
|
def label(self):
|
|
|
|
return self.spec.label
|
|
|
|
|
|
|
|
@property
|
|
|
|
def classifiers(self):
|
|
|
|
return self.spec.classifiers
|
2017-03-09 17:39:44 +00:00
|
|
|
|
2017-03-20 16:24:22 +00:00
|
|
|
@property
|
|
|
|
def status(self):
|
|
|
|
return self._status
|
|
|
|
|
|
|
|
@status.setter
|
|
|
|
def status(self, value):
|
|
|
|
self._status = value
|
|
|
|
if self.output:
|
|
|
|
self.output.status = value
|
|
|
|
|
2017-03-09 17:39:44 +00:00
|
|
|
def __init__(self, spec, iteration, context):
|
|
|
|
self.logger = logging.getLogger('job')
|
|
|
|
self.spec = spec
|
|
|
|
self.iteration = iteration
|
|
|
|
self.context = context
|
|
|
|
self.workload = None
|
|
|
|
self.output = None
|
2017-09-15 13:44:20 +01:00
|
|
|
self.run_time = None
|
2017-03-09 17:39:44 +00:00
|
|
|
self.retries = 0
|
2017-03-20 16:24:22 +00:00
|
|
|
self._status = Status.NEW
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def load(self, target, loader=pluginloader):
|
2017-03-17 16:28:21 +00:00
|
|
|
self.logger.info('Loading job {}'.format(self.id))
|
2017-07-04 15:44:05 +01:00
|
|
|
if self.iteration == 1:
|
|
|
|
self.workload = loader.get_workload(self.spec.workload_name,
|
|
|
|
target,
|
|
|
|
**self.spec.workload_parameters)
|
|
|
|
self.workload.init_resources(self.context)
|
|
|
|
self.workload.validate()
|
|
|
|
self._workload_cache[self.id] = self.workload
|
|
|
|
else:
|
|
|
|
self.workload = self._workload_cache[self.id]
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def initialize(self, context):
|
|
|
|
self.logger.info('Initializing job {}'.format(self.id))
|
2017-03-17 15:57:05 +00:00
|
|
|
with signal.wrap('WORKLOAD_INITIALIZED', self, context):
|
|
|
|
self.workload.initialize(context)
|
2017-09-18 16:12:03 +01:00
|
|
|
self.set_status(Status.PENDING)
|
2017-03-16 17:54:48 +00:00
|
|
|
context.update_job_state(self)
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def configure_target(self, context):
|
|
|
|
self.logger.info('Configuring target for job {}'.format(self.id))
|
2017-03-06 17:29:15 +00:00
|
|
|
context.tm.commit_runtime_parameters(self.spec.runtime_parameters)
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def setup(self, context):
|
|
|
|
self.logger.info('Setting up job {}'.format(self.id))
|
2017-03-17 15:57:05 +00:00
|
|
|
with signal.wrap('WORKLOAD_SETUP', self, context):
|
|
|
|
self.workload.setup(context)
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def run(self, context):
|
|
|
|
self.logger.info('Running job {}'.format(self.id))
|
2017-03-17 15:57:05 +00:00
|
|
|
with signal.wrap('WORKLOAD_EXECUTION', self, context):
|
2017-09-15 13:44:20 +01:00
|
|
|
start_time = datetime.utcnow()
|
|
|
|
try:
|
|
|
|
self.workload.run(context)
|
|
|
|
finally:
|
|
|
|
self.run_time = datetime.utcnow() - start_time
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def process_output(self, context):
|
|
|
|
self.logger.info('Processing output for job {}'.format(self.id))
|
2017-03-27 17:31:44 +01:00
|
|
|
with signal.wrap('WORKLOAD_RESULT_EXTRACTION', self, context):
|
|
|
|
self.workload.extract_results(context)
|
|
|
|
context.extract_results()
|
|
|
|
with signal.wrap('WORKLOAD_OUTPUT_UPDATE', self, context):
|
|
|
|
self.workload.update_output(context)
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def teardown(self, context):
|
|
|
|
self.logger.info('Tearing down job {}'.format(self.id))
|
2017-03-17 15:57:05 +00:00
|
|
|
with signal.wrap('WORKLOAD_TEARDOWN', self, context):
|
|
|
|
self.workload.teardown(context)
|
2017-03-09 17:39:44 +00:00
|
|
|
|
|
|
|
def finalize(self, context):
|
|
|
|
self.logger.info('Finalizing job {}'.format(self.id))
|
2017-03-17 15:57:05 +00:00
|
|
|
with signal.wrap('WORKLOAD_FINALIZED', self, context):
|
|
|
|
self.workload.finalize(context)
|
2017-09-18 16:12:03 +01:00
|
|
|
|
|
|
|
def set_status(self, status, force=False):
|
|
|
|
status = Status(status)
|
|
|
|
if force or self.status < status:
|
|
|
|
self.status = status
|