File indexing completed on 2026-08-12 09:36:11
0001
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
0032
0033
0034
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
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
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
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
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()