Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-07-22 09:42:03

0001 import json
0002 import os
0003 import subprocess
0004 import sys
0005 from pathlib import Path
0006 
0007 from django.conf import settings
0008 from django.core.management.base import BaseCommand, CommandError
0009 from django.db import OperationalError, ProgrammingError
0010 
0011 from monitor_app.epicprod_inventory import (
0012     build_expected_files_for_task,
0013     sync_expected_files_for_task,
0014     sync_job_from_study_data,
0015 )
0016 from monitor_app.panda.queries import study_job
0017 
0018 
0019 CACHE_PAYLOAD_LOG = (
0020     Path(__file__).resolve().parents[4] / 'scripts' / 'cache-payload-log.py'
0021 )
0022 STUDY_FETCH_TIMEOUT = int(os.environ.get('EPICPROD_STUDY_FETCH_TIMEOUT', '150'))
0023 
0024 
0025 class Command(BaseCommand):
0026     help = "Build or refresh ePIC production job/file inventory."
0027 
0028     def add_arguments(self, parser):
0029         group = parser.add_mutually_exclusive_group(required=True)
0030         group.add_argument('--pandaid', type=int, help='PanDA job id to refresh')
0031         group.add_argument('--jeditaskid', type=int, help='JEDI task id to build expected files for')
0032         group.add_argument('--prod-task', help='PCS task name/composed name to build expected files for')
0033         parser.add_argument('--spec-file',
0034                             help='EVGEN spec JSON to use as the expected-file source')
0035         parser.add_argument('--dry-run', action='store_true',
0036                             help='Print derived rows without writing database changes')
0037 
0038     def handle(self, *args, **options):
0039         try:
0040             if options['pandaid']:
0041                 self._sync_pandaid(options['pandaid'], options['dry_run'])
0042             else:
0043                 task = self._resolve_task(
0044                     jeditaskid=options.get('jeditaskid'),
0045                     name=options.get('prod_task'),
0046                 )
0047                 spec = None
0048                 if options.get('spec_file'):
0049                     with open(options['spec_file']) as f:
0050                         spec = json.load(f)
0051                 self._sync_task(task, options['dry_run'], spec=spec)
0052         except (OperationalError, ProgrammingError) as exc:
0053             raise CommandError(
0054                 f'epicprod inventory tables are not available; run migrations first: {exc}'
0055             ) from exc
0056 
0057     def _resolve_task(self, *, jeditaskid=None, name=None):
0058         from pcs.models import ProdTask
0059         qs = ProdTask.objects.select_related('dataset', 'prod_config')
0060         if jeditaskid:
0061             from pcs.models import PandaTasks
0062             assoc = (
0063                 PandaTasks.objects
0064                 .select_related('prod_task', 'prod_task__dataset', 'prod_task__prod_config')
0065                 .filter(jedi_task_id=jeditaskid)
0066                 .first()
0067             )
0068             task = assoc.prod_task if assoc else qs.filter(panda_task_id=jeditaskid).first()
0069             if not task:
0070                 raise CommandError(f'No PCS task association records jediTaskID={jeditaskid}')
0071             return task
0072         task = qs.filter(name=name).first()
0073         if not task:
0074             for t in qs.all():
0075                 if t.composed_name == name:
0076                     return t
0077             raise CommandError(f'No PCS task found for {name!r}')
0078         return task
0079 
0080     def _sync_task(self, task, dry_run, spec=None):
0081         if dry_run:
0082             rows = build_expected_files_for_task(task, spec=spec)
0083             self.stdout.write(json.dumps(self._jsonable(rows), indent=2))
0084             return
0085         rows = sync_expected_files_for_task(task, spec=spec)
0086         self.stdout.write(
0087             self.style.SUCCESS(
0088                 f'synced {len(rows)} expected file row(s) for {task.composed_name}'
0089             )
0090         )
0091 
0092     def _sync_pandaid(self, pandaid, dry_run):
0093         data = study_job(pandaid)
0094         if 'error' in data:
0095             raise CommandError(data['error'])
0096         if not dry_run:
0097             self._cache_payload_log_before_study(pandaid, data)
0098         if dry_run:
0099             self.stdout.write(json.dumps(self._jsonable(data), indent=2))
0100             return
0101         job = sync_job_from_study_data(data)
0102         self.stdout.write(
0103             self.style.SUCCESS(
0104                 f'synced epicprod job {job.pandaid}: phase={job.phase or "(none)"}'
0105             )
0106         )
0107 
0108     def _cache_payload_log_before_study(self, pandaid, data):
0109         job = data.get('job') or {}
0110         log_file = data.get('log_file') or {}
0111         jeditaskid = job.get('jeditaskid')
0112         scope = log_file.get('scope')
0113         lfn = log_file.get('lfn')
0114         if not (jeditaskid and scope and lfn):
0115             return
0116         cache_root = getattr(settings, 'SWF_TMP_DIR', '/data/swf-tmp')
0117         done = os.path.join(cache_root, 'panda-logs', str(jeditaskid), str(pandaid), '.done')
0118         if os.path.isfile(done):
0119             return
0120         cmd = [
0121             sys.executable, str(CACHE_PAYLOAD_LOG),
0122             '--scope', str(scope),
0123             '--lfn', str(lfn),
0124             '--jeditaskid', str(jeditaskid),
0125             '--pandaid', str(pandaid),
0126         ]
0127         self.stdout.write(f'caching payload log before study: pandaid={pandaid}')
0128         try:
0129             p = subprocess.run(
0130                 cmd, capture_output=True, text=True,
0131                 timeout=STUDY_FETCH_TIMEOUT,
0132             )
0133         except subprocess.TimeoutExpired as exc:
0134             raise CommandError(
0135                 f'payload log fetch timed out before study for pandaid={pandaid} '
0136                 f'after {STUDY_FETCH_TIMEOUT}s'
0137             ) from exc
0138         if p.returncode != 0:
0139             stderr = (p.stderr or '').strip()
0140             reason = stderr.splitlines()[-1] if stderr else f'rc={p.returncode}'
0141             raise CommandError(
0142                 f'payload log fetch failed before study for pandaid={pandaid}: {reason}'
0143             )
0144         for line in (p.stderr or '').splitlines():
0145             self.stdout.write(f'  cache-payload-log: {line}')
0146 
0147     @staticmethod
0148     def _jsonable(value):
0149         return json.loads(json.dumps(value, default=str))