X-Git-Url: http://mmka.chem.univ.gda.pl/gitweb/?a=blobdiff_plain;f=qcg%2Futils.py;h=5555e16b2972a9afaf10a9a536e523b1e9a26a06;hb=dd8e6b0785c78c42c3866f161dd853a54e47275a;hp=bc9135f1e414ab94c912d80f4114c025d9e2731d;hpb=b9bcd422c66e0cd9b0f20a0c037dbc2d811bb59f;p=qcg-portal.git diff --git a/qcg/utils.py b/qcg/utils.py index bc9135f..5555e16 100644 --- a/qcg/utils.py +++ b/qcg/utils.py @@ -1,10 +1,15 @@ +import os +import string +import random + from django.core.paginator import Paginator -from django.db import transaction from django.utils.formats import date_format -from django.utils.timezone import now, localtime -from pyqcg.service import Registry -from pyqcg.utils import Credential, TimePeriod +from django.utils.timezone import localtime +from pyqcg import QCG +from pyqcg.utils import Credential +from pyqcg.description import JobDescription +from filex.ftp import FTPOperation from qcg import constants @@ -18,51 +23,6 @@ def username_from_dn(dn): return username -@transaction.atomic -def update_user_data(user, proxy): - from qcg.models import User, Job, Task, Allocation, NodeInfo - - credential = Credential(proxy) - registry = Registry(credential) - - # put lock on user record (hopefully..?) - user = User.objects.select_for_update().get(pk=user.pk) - - changed_filter = {'changed': TimePeriod(after=user.last_update)} - - ################################### - # Jobs - ################################### - for qcg_job in registry.jobs(**changed_filter): - params = Job.qcg_map(qcg_job, user) - job_id = params.pop('job_id') - - Job.objects.update_or_create(job_id=job_id, defaults=params) - - ################################### - # Tasks - ################################### - jobs_cache = {j.job_id: j for j in Job.objects.filter(owner=user)} - for qcg_task in registry.tasks(**changed_filter): - params = Task.qcg_map(qcg_task, jobs_cache) - task_id = params.pop('task_id') - - task, created = Task.objects.update_or_create(job__job_id=qcg_task.job_id, task_id=task_id, defaults=params) - - if not created: - task.allocations.all().delete() - - for qcg_alloc in qcg_task.allocations: - alloc = task.allocations.create(**Allocation.qcg_map(qcg_alloc)) - - for qcg_node in qcg_alloc.nodes: - alloc.nodes.create(**NodeInfo.qcg_map(qcg_node)) - - # release user lock - user.last_update = now() - user.save() - - def try_parse_int(s, default=None, base=10): try: return int(s, base) @@ -85,3 +45,102 @@ def paginator_context(request, objects, per_page=constants.PER_PAGE): def localtime_str(datetime): return date_format(localtime(datetime), 'DATETIME_FORMAT') + + +def random_id(size=8, chars=string.ascii_uppercase + string.digits): + return ''.join(random.choice(chars) for _ in range(size)) + + +def chunks(seq, size): + return (seq[pos:pos + size] for pos in xrange(0, len(seq), size)) + + +def to_job_desc(params, proxy): + QCG.start() + desc = JobDescription(Credential(proxy)) + + direct_map = ('env_variables', 'executable', 'arguments', 'note', 'grant', 'hosts', 'properties', 'queue', 'procs', + 'wall_time', 'memory', 'memory_per_slot', 'modules', 'input', 'stage_in', 'native', 'notify', + 'preprocess', 'postprocess', 'persistent') + + for name in direct_map: + if params[name]: + setattr(desc, name, params[name]) + + if params['application']: + desc.set_application(*params['application']) + desc.stage_in += [params['master_file']] + desc.arguments.insert(0, os.path.basename(params['master_file'])) + if params['script']: + ftp = FTPOperation(proxy) + + ftp.mkdir(constants.QCG_DATA_URL, parents=True) + url = os.path.join(constants.QCG_DATA_URL, 'script.{}.sh'.format(random_id())) + ftp.put(url) + + for chunk in chunks(params['script'], 4096): + ftp.stream.put(chunk) + ftp.stream.put(None) + ftp.wait() + + desc.executable = url + if params['nodes']: + desc.set_nodes(*params['nodes']) + if params['reservation']: + desc.set_reservation(params['reservation']) + if params['watch_output']: + desc.set_watch_output(params['watch_output'], params['watch_output_pattern']) + # TODO monitoring + + return desc + + +def to_form_data(xml): + # prevent circular import errors + from forms import JobDescriptionForm + + QCG.start() + desc = JobDescription() + desc.xml_description = xml + + direct_map = ('env_variables', 'executable', 'arguments', 'note', 'grant', 'hosts', 'properties', 'queue', 'procs', + 'wall_time', 'memory', 'memory_per_slot', 'modules', 'input', 'stage_in', 'native', 'persistent') + + params = {} + for name in direct_map: + attr = getattr(desc, name) + if isinstance(attr, bool) or attr: + params[name] = attr + + if desc.application is not None: + app_name, app_ver = desc.application + params['application'] = app_name if app_ver is None else app_name + '/' + app_name + stage_in = params['stage_in'] + params['stage_in'], params['master_file'] = stage_in[:-1], stage_in[-1] + params['arguments'] = params['arguments'][1:] + if desc.nodes is not None: + params['nodes'] = ':'.join(map(str, desc.nodes)) + if desc.reservation is not None: + res_id, res_type = desc.reservation + params['reservation'] = res_id + if desc.notify is not None: + params['notify_type'], params['notify_address'] = desc.notify.split(':') + if desc.watch_output is not None: + watch_output, params['watch_output_pattern'] = desc.watch_output + params['watch_output_type'], params['watch_output_address'] = watch_output.split(':') + if desc.preprocess is not None: + if desc.preprocess.startswith('gsiftp://'): + params['preprocess_type'] = JobDescriptionForm.Process.SCRIPT + params['preprocess_script'] = desc.preprocess + else: + params['preprocess_type'] = JobDescriptionForm.Process.CMD + params['preprocess_cmd'] = desc.preprocess + if desc.postprocess is not None: + if desc.postprocess.startswith('gsiftp://'): + params['postprocess_type'] = JobDescriptionForm.Process.SCRIPT + params['postprocess_script'] = desc.postprocess + else: + params['postprocess_type'] = JobDescriptionForm.Process.CMD + params['postprocess_cmd'] = desc.postprocess + + return params