Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-08-12 09:36:11

0001 #!/usr/bin/env python
0002 """
0003 find_stale_jobs.py — report production jobs that look stuck.
0004 
0005 A job is considered a candidate when:
0006   - status matches --status  (default: 'running')
0007   - memoryprovisioned > --min-memory  (default: 11000 MB)
0008   - the 'running' timestamp is older than --running-hours  (default: 72 h)
0009 
0010 For each candidate the 'out' log file mtime is checked:
0011   - if the file is older than --out-hours  (default: 24 h) it is flagged STALE
0012   - if the file is newer it is flagged ACTIVE  (job may still be writing)
0013   - if the file cannot be stat'd it is flagged MISSING
0014 
0015 Does NOT depend on any other module in this project.
0016 """
0017 
0018 import argparse
0019 import logging
0020 import os
0021 import random
0022 import sys
0023 import time
0024 from datetime import datetime, timezone
0025 from pathlib import Path
0026 
0027 
0028 try:
0029     import pyodbc
0030 except ImportError:
0031     pyodbc = None  # allow --dryrun without pyodbc installed
0032 
0033 # ============================================================================
0034 # Inline logger
0035 # ============================================================================
0036 CHATTY_LEVEL_NUM = 5
0037 logging.addLevelName(CHATTY_LEVEL_NUM, "CHATTY")
0038 
0039 def _chatty(self, message, *args, **kws):
0040     if self.isEnabledFor(CHATTY_LEVEL_NUM):
0041         self._log(CHATTY_LEVEL_NUM, message, args, stacklevel=2, **kws)
0042 logging.Logger.chatty = _chatty
0043 
0044 
0045 class _Fmt(logging.Formatter):
0046     show_datetime = True
0047     grey     = "\x1b[38;20m"
0048     yellow   = "\x1b[33;20m"
0049     green    = "\x1b[32;20m"
0050     blue     = "\x1b[36;20m"
0051     red      = "\x1b[31;20m"
0052     bold_red = "\x1b[31;1m"
0053     reset    = "\x1b[0m"
0054     _datetime_fmt = "%(asctime)s [%(levelname)s] - %(message)s"
0055     _plain_fmt    = "[%(levelname)s] - %(message)s"
0056 
0057     def _base_format(self):
0058         return self._datetime_fmt if self.show_datetime else self._plain_fmt
0059 
0060     def format(self, record):
0061         base = self._base_format()
0062         formats = {
0063             CHATTY_LEVEL_NUM: self.yellow   + base + " (%(filename)s:%(lineno)d) " + self.reset,
0064             logging.DEBUG:    self.grey     + base + " (%(filename)s:%(lineno)d) " + self.reset,
0065             logging.INFO:     self.green    + base + self.reset,
0066             logging.WARNING:  self.blue     + base + " (%(filename)s:%(lineno)d) " + self.reset,
0067             logging.ERROR:    self.red      + base + " (%(filename)s:%(lineno)d) " + self.reset,
0068             logging.CRITICAL: self.bold_red + base + " (%(filename)s:%(lineno)d) " + self.reset,
0069         }
0070         formatter = logging.Formatter(formats.get(record.levelno, base))
0071         return formatter.format(record)
0072 
0073 
0074 _log = logging.getLogger('find_stale_jobs')
0075 if not _log.hasHandlers():
0076     _ch = logging.StreamHandler()
0077     _ch.setFormatter(_Fmt())
0078     _log.addHandler(_ch)
0079 
0080 CHATTY = _log.chatty
0081 DEBUG  = _log.debug
0082 INFO   = _log.info
0083 WARN   = _log.warning
0084 ERROR  = _log.error
0085 
0086 # ============================================================================
0087 # DB connection strings
0088 # ============================================================================
0089 if os.uname().sysname == 'Darwin':
0090     _STATR = 'DRIVER=PostgreSQL Unicode;SERVER=localhost;DATABASE=productiondb;READONLY=True;UID=eickolja'
0091 else:
0092     _STATR = 'DSN=Production_read;READONLY=True;UID=argouser'
0093 
0094 _RETRYABLE = {'40001', '53300', '57P03', '08006', '08001'}
0095 
0096 def _db_query(query: str, ntries: int = 5):
0097     CHATTY(f'[sql]\n{query}')
0098     if pyodbc is None:
0099         ERROR("pyodbc not available; cannot query DB.")
0100         sys.exit(1)
0101     for itry in range(ntries):
0102         try:
0103             conn = pyodbc.connect(_STATR)
0104             curs = conn.cursor()
0105             curs.execute(query)
0106             return curs
0107         except pyodbc.Error as exc:
0108             state = exc.args[0]
0109             ERROR(f"Attempt {itry + 1}/{ntries}: {exc}")
0110             if state in _RETRYABLE:
0111                 delay = min(60, (2 ** itry) * (0.5 + random.random()))
0112                 WARN(f"Retrying in {delay:.1f}s …")
0113                 time.sleep(delay)
0114             else:
0115                 ERROR("Non-retryable DB error. Stop.")
0116                 sys.exit(41)
0117         except Exception as exc:
0118             ERROR(f"Unexpected error: {exc}")
0119             sys.exit(41)
0120     ERROR("Exhausted all DB attempts. Stop.")
0121     sys.exit(41)
0122 
0123 # ============================================================================
0124 # Helpers
0125 # ============================================================================
0126 def _age_hours(ts) -> float:
0127     """Return how many hours ago `ts` (datetime or epoch float) was."""
0128     if isinstance(ts, (int, float)):
0129         dt = datetime.fromtimestamp(ts, tz=timezone.utc)
0130     elif ts.tzinfo is None:
0131         dt = ts.replace(tzinfo=timezone.utc)
0132     else:
0133         dt = ts
0134     return (datetime.now(tz=timezone.utc) - dt).total_seconds() / 3600.0
0135 
0136 
0137 def _out_status(out_path: str, out_hours: float) -> tuple[str, str]:
0138     """
0139     Return (label, detail) for the given out file path.
0140     label: 'STALE' | 'ACTIVE' | 'MISSING'
0141     """
0142     try:
0143         mtime = Path(out_path).stat().st_mtime
0144     except OSError:
0145         return 'MISSING', 'cannot stat'
0146     age_h = _age_hours(mtime)
0147     age_str = f"{age_h:.1f} h ago"
0148     if age_h > out_hours:
0149         return 'STALE', age_str
0150     return 'ACTIVE', age_str
0151 
0152 # ============================================================================
0153 # Main logic
0154 # ============================================================================
0155 def _build_query(args) -> str:
0156     conditions = [
0157         f"status = '{args.status}'",
0158         f"memoryprovisioned > {args.min_memory}",
0159         f"started < NOW() - INTERVAL '{args.running_hours} hours'",
0160     ]
0161     if args.dataset:
0162         conditions.append(f"dataset = '{args.dataset}'")
0163     if args.dsttype:
0164         op = 'LIKE' if '%' in args.dsttype else '='
0165         conditions.append(f"dsttype {op} '{args.dsttype}'")
0166     if args.tag:
0167         conditions.append(f"tag = '{args.tag}'")
0168 
0169     where = '\n  AND '.join(conditions)
0170     return (
0171         "SELECT ClusterId, ProcId, submission_host, started, out, MemoryProvisioned, dataset, dsttype"
0172         " FROM production_jobs"
0173         f"\nWHERE {where}"
0174         "\nORDER BY started ASC;"
0175     )
0176 
0177 
0178 def main():
0179     args = _parse_args()
0180     if args.running_hours is None:
0181         args.running_hours = args.out_hours
0182     _set_loglevel(args)
0183     _Fmt.show_datetime = False
0184 
0185     query = _build_query(args)
0186     INFO(f"Query:\n{query}")
0187 
0188     if args.dryrun:
0189         INFO("[dryrun] not executing query.")
0190         return
0191 
0192     curs = _db_query(query)
0193     rows = curs.fetchall()
0194     conn = getattr(curs, 'connection', None)
0195     curs.close()
0196     if conn:
0197         conn.close()
0198 
0199     if not rows:
0200         INFO("No matching jobs found.")
0201         return
0202 
0203     INFO(f"{len(rows)} candidate job(s) found. Checking out files …\n")
0204 
0205     counts = {'STALE': 0, 'ACTIVE': 0, 'MISSING': 0}
0206     stale_lines = []
0207     for cluster_id, proc_id, submission_host, started, out, mem_prov, dataset, dsttype in rows:
0208         condor_id = f"{cluster_id}.{proc_id}" if cluster_id is not None else "?.?"
0209         submit_host = submission_host or '?'
0210         running_age = _age_hours(started)
0211         label, detail = _out_status(out, args.out_hours)
0212         counts[label] += 1
0213         line = f"[{label:7s}]  running {running_age:6.1f} h ago  {detail:30s}  {out}"
0214         if label == 'STALE':
0215             ERROR(line)
0216             stale_lines.append((condor_id, submit_host, running_age, detail, mem_prov, dataset, dsttype, out))
0217         elif label == 'MISSING':
0218             WARN(line)
0219         else:
0220             INFO(line)
0221 
0222     INFO(f"\nSummary: {counts['STALE']} STALE  {counts['ACTIVE']} ACTIVE  {counts['MISSING']} MISSING  (of {len(rows)} total)")
0223 
0224     if stale_lines:
0225         cid_w  = max(len(r[0]) for r in stale_lines)
0226         host_w = max(len(r[1]) for r in stale_lines)
0227         ds_w   = max(len(r[5]) for r in stale_lines)
0228         dst_w  = max(len(r[6]) for r in stale_lines)
0229         mtime_w = max(len(r[3]) for r in stale_lines)
0230         header = (f"{'CONDOR_ID':<{cid_w}}  {'SUBMIT_HOST':<{host_w}}"
0231                   f"  {'STARTED_AGO':>11}  {'OUT_MTIME':<{mtime_w}}"
0232                   f"  {'MEM_PROV':>8}"
0233                   f"  {'DATASET':<{ds_w}}  {'DSTTYPE':<{dst_w}}")
0234         if args.show_out:
0235             header += "  OUT"
0236         print(f"\nStale jobs:\n{header}")
0237         for condor_id, submit_host, running_age, detail, mem_prov, dataset, dsttype, out in stale_lines:
0238             prov_s = f"{mem_prov:>8}" if mem_prov is not None else f"{'?':>8}"
0239             row = (f"{condor_id:<{cid_w}}  {submit_host:<{host_w}}"
0240                    f"  {running_age:>10.1f}h  {detail:<{mtime_w}}"
0241                    f"  {prov_s}"
0242                    f"  {dataset:<{ds_w}}  {dsttype:<{dst_w}}")
0243             if args.show_out:
0244                 row += f"  {out}"
0245             print(row)
0246 
0247 # ============================================================================
0248 # Argument parsing
0249 # ============================================================================
0250 def _add_verbosity(parser):
0251     vgroup = parser.add_mutually_exclusive_group()
0252     vgroup.add_argument('-v', '--verbose', action='count', default=0,
0253                         help='Increase verbosity (-v INFO, -vv DEBUG, -vvv CHATTY).')
0254     vgroup.add_argument('-d', '--debug',  action='store_true', help='Alias for -vv (DEBUG).')
0255     vgroup.add_argument('--chatty',       action='store_true', help='Alias for -vvv (CHATTY).')
0256 
0257 
0258 def _set_loglevel(args):
0259     if args.chatty or args.verbose >= 3:
0260         _log.setLevel(CHATTY_LEVEL_NUM)
0261     elif args.debug or args.verbose == 2:
0262         _log.setLevel(logging.DEBUG)
0263     else:
0264         _log.setLevel(logging.INFO)
0265 
0266 
0267 def _parse_args():
0268     parser = argparse.ArgumentParser(
0269         description='Find production jobs that appear stuck.',
0270         formatter_class=argparse.RawDescriptionHelpFormatter,
0271         epilog="""
0272 Examples:
0273   # Default: status=running, memoryprovisioned>11000, started >72 h ago
0274   find_stale_jobs.py
0275 
0276   # Narrow to a dataset and check out files older than 12 h
0277   find_stale_jobs.py --dataset run3oo --out-hours 12
0278 
0279   # Only show query, do not hit DB
0280   find_stale_jobs.py --dryrun
0281 """,
0282     )
0283 
0284     parser.add_argument('--status',        default='running',
0285                         help='Job status to match (default: running).')
0286     parser.add_argument('--min-memory',    dest='min_memory', type=int, default=11000,
0287                         help='Minimum memoryprovisioned in MB (default: 11000).')
0288     parser.add_argument('--running-hours', dest='running_hours', type=float, default=None,
0289                         help='Select jobs running for at least this many hours (default: same as --out-hours).')
0290     parser.add_argument('--out-hours',     dest='out_hours',     type=float, default=24.0,
0291                         help='Flag out file as STALE when its mtime is older than this many hours (default: 24).')
0292     parser.add_argument('--dataset',  default=None, help='Filter by dataset name.')
0293     parser.add_argument('--dsttype',  default=None, help='Filter by dsttype (%% triggers LIKE).')
0294     parser.add_argument('--tag',      default=None, help='Filter by production tag.')
0295     parser.add_argument('--show-out', dest='show_out', action='store_true', default=False,
0296                         help='Include the OUT file path in the final stale-jobs table.')
0297     parser.add_argument('-n', '--dryrun', action='store_true', default=False,
0298                         help='Print the SQL query without executing it.')
0299     _add_verbosity(parser)
0300 
0301     return parser.parse_args()
0302 
0303 
0304 if __name__ == '__main__':
0305     main()