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))