From fe4eabbfecd7d544f1ddce0882c075474987a228 Mon Sep 17 00:00:00 2001 From: zaihuaji Date: Fri, 2 Oct 2026 11:07:57 -0500 Subject: [PATCH] guard every batch wait and catch SIGTERM so a killed run still reports; bump version to 3.0.20 A PBS run could pass its walltime and be killed without ever reporting. The slot-acquire wait check_child(..., -1) blocks for as long as all child slots stay busy, with the deadline check stranded outside that loop, so a run whose children each tar ~100GB walked past the cutoff while waiting. New wait_child_slot() polls instead and checks the deadline between polls, the same way wait_all_children() already guards the terminal wait; both of the remaining blocking waits now use them. Four more ways a run could die unheard: - report_dscheck_email() marked the report delivered before trying to park it, so a failed cache both parked nothing and disabled every later attempt - the quit path waited for children unguarded and was killed before its email - the deadline guard excluded -A 16/32/64/128, which are submitted with the same 24 hour walltime - ACTMSG has no wording for the hidden actions, so building the report raised KeyError instead PBS sends SIGTERM before the SIGKILL and nothing trapped it. A batch run now catches it and parks a report in the seconds it has left, as a backstop for the cases the walltime cutoff cannot cover. Co-Authored-By: Claude Opus 4.6 --- pyproject.toml | 2 +- src/rda_python_dsquasar/dsquasar.py | 155 ++++++++++++++++++++++------ 2 files changed, 124 insertions(+), 33 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 59a9745..dd3d3c6 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "rda_python_dsquasar" -version = "3.0.19" +version = "3.0.20" authors = [ { name="Zaihua Ji", email="zji@ucar.edu" }, ] diff --git a/src/rda_python_dsquasar/dsquasar.py b/src/rda_python_dsquasar/dsquasar.py index 259ce28..75cc322 100644 --- a/src/rda_python_dsquasar/dsquasar.py +++ b/src/rda_python_dsquasar/dsquasar.py @@ -23,6 +23,7 @@ import re import sys import time +import signal from os import path as op from time import time as tm from rda_python_common.pg_cmd import PgCMD @@ -96,6 +97,7 @@ def __init__(self): self.MPBLIMIT = 5000 # tar file count per batch process for uploading self.MPBMAX = 4 # maximum number of batch processes for uploading self.ONEHOUR = 3600 # seconds; headroom kept to report before the walltime + self.ERETRY = 600 # seconds to wait before retrying a failed progress report self.WALLTIME = 24*3600 # PBS walltime assumed, refreshed from PBS at start of a run self.MAXRUNTIME = self.WALLTIME - self.ONEHOUR # cutoff to cache a progress report # a repeat submit normally blocks as a duplicate while the batch job is running. if @@ -124,6 +126,7 @@ def __init__(self): 'doemail' : 0, 'starttime' : 0, # wall-clock start of the run, for the PBS walltime guard 'einfo' : 0, # set once a progress report is cached into dscheck.einfo + 'eretry' : 0, # earliest retry time after a progress report failed to cache 'maxworkers' : self.MAXWORKERS, # -W, maximum concurrent workers per command 'worker' : 1, # -w, this run's worker slot; >1 for an added extra worker 'cmd' : None @@ -200,6 +203,7 @@ def start_actions(self): self.cmdlog(self.PGBACK['cmd']) self.PGBACK['starttime'] = tm() self.set_walltime_deadline() + self.catch_batch_termination() if self.sopts['u']: if self.sopts['a']: self.pglog("-u: Dataset IDs must be provided to Unlock datasets", self.LOGWRN) self.unlock_datasets() @@ -238,7 +242,7 @@ def start_actions(self): self.PGLOG['ERRCNT'] or self.PGBACK['einfo']): self.PGBACK['doemail'] = 0 if self.PGBACK['doemail']: - amsg = self.ACTMSG[self.PGBACK['action']] + amsg = self.action_message() bmsg = self.BACKMSG[self.PGBACK['backflag']] if self.PGBACK['backflag'] else 'backup' dcnt = len(self.dsids) dmsg = self.dsids[0] if dcnt == 1 else "{} datasets".format(dcnt if dcnt > 1 else 'All') @@ -791,7 +795,7 @@ def quit_dsquasar(self, qinfo, msg = None): # a second report, drop the progress report parked in dscheck.einfo by the parent, and # mark the shared check record failed while the parent is still working if self.PGSIG['PPID'] > 1: sys.exit(1) - if self.PGBACK['mproc'] > 1: self.check_child(None, 0, self.LOGWRN, 1) + self.wait_all_children() if qinfo: if 'dslocks' in qinfo and qinfo['dslocks']: for dsid in qinfo['dslocks']: self.lock_dataset(dsid, 0, self.LGEREX) @@ -799,7 +803,7 @@ def quit_dsquasar(self, qinfo, msg = None): dcnt = len(qinfo['qdsids']) fcnt = qinfo['qfcnt'] ssize = self.format_float_value(qinfo['qsize']) - amsg = self.ACTMSG[self.PGBACK['action']] + amsg = self.action_message() bmsg = self.BACKMSG[qinfo['backflag']] dmsg = qinfo['qdsids'][0] if dcnt == 1 else "{} datasets".format(dcnt) msg = "Quit {}: {} {} files for {}({}) files of {}".format(amsg, qcnt, bmsg, fcnt, ssize, dmsg) @@ -815,6 +819,13 @@ def quit_dsquasar(self, qinfo, msg = None): if self.PGBACK['pstep']: self.record_dscheck_status("F") self.pgexit(0) + # ACTMSG carries no wording for the hidden actions (-A 32/64), so name those by number + # rather than raise a KeyError: a report that cannot be worded is still a report that + # has to go out, and crashing while building it loses the run as silently as a kill does + def action_message(self, act = None): + if act is None: act = self.PGBACK['action'] + return self.ACTMSG[act] if act in self.ACTMSG else "Action {}".format(act) + # a batch run normally gets the full 24 hours asked for at submit time, but the dscheck # daemon caps that to the maximum of the PBS queue it lands in - only 6 hours for the # default queue - and to any planned system down. read the granted walltime back from @@ -830,13 +841,54 @@ def set_walltime_deadline(self): self.PGLOG['CURBID'], self.seconds_to_string_time(self.WALLTIME)), self.LOGWRN) return wtime = 3600*int(ms.group(1)) + 60*int(ms.group(2)) + (int(ms.group(3)) if ms.group(3) else 0) - if wtime <= self.ONEHOUR: return # too short to keep the reporting headroom + if wtime <= self.ONEHOUR: # too short to keep the reporting headroom + self.pglog("PBS Job {} walltime {} is too short to report progress ahead of; assume {}".format( + self.PGLOG['CURBID'], self.seconds_to_string_time(wtime), + self.seconds_to_string_time(self.WALLTIME)), self.LOGWRN) + return self.WALLTIME = wtime self.MAXRUNTIME = wtime - self.ONEHOUR self.pglog("PBS Job {} walltime {}: report progress after {}".format( self.PGLOG['CURBID'], self.seconds_to_string_time(self.WALLTIME), self.seconds_to_string_time(self.MAXRUNTIME)), self.LOGWRN) + # PBS kills a job that runs out of walltime with SIGTERM first and SIGKILL a few seconds + # later (the MoM's kill_delay, 10 seconds by default), and the same pair is what qdel + # sends. SIGTERM is the only warning there is, and nothing in the common library traps + # it, so the run used to die on the spot with whatever it had to say unsaid. catching it + # turns the last seconds into a report. only a batch run arms this: on the command line + # SIGTERM must keep killing the process the way the user expects. + def catch_batch_termination(self): + if self.PGLOG['CURBID'] < 1: return + signal.signal(signal.SIGTERM, self.batch_term_handler) + + # park a progress report and then die of the signal that was sent. this runs on borrowed + # time, so it does only what it must: no tar queue counts (two database queries we may + # not get to finish - the done count comes from memory), no dataset unlocking (the + # dscheck daemon already cleans up after a dead pid). the default handler is restored + # first so that a second signal, or the SIGKILL that follows, ends the run outright + # instead of re-entering here should the report hang. + # the report may have to be written from inside an interrupted database call; that is + # why it goes through report_dscheck_email, whose cache_customized_email falls back to + # sending the mail directly when the UPDATE fails. + def batch_term_handler(self, signum, frame): + signal.signal(signum, signal.SIG_DFL) + # a forked child shares the parent's dscheck record: it must not report for the run + if self.PGSIG['PPID'] > 1: os._exit(1) + if not self.PGBACK['einfo'] and self.PGBACK['doemail'] and self.PGLOG['DSCHECK']: + etime = self.seconds_to_string_time(int(tm() - self.PGBACK['starttime'])) + amsg = self.action_message() + dmsg = self.batch_done_count() + wmsg = self.seconds_to_string_time(self.WALLTIME) + self.pglog(self.INDENT + "{}: Terminated by signal {} after {} of the {} PBS walltime, {}".format( + amsg, signum, etime, wmsg, dmsg), self.LOGACT) + self.set_email("{}: {} terminated after {} of the {} PBS walltime, {}!".format( + self.PGBACK['cmd'], amsg, etime, wmsg, dmsg), self.EMLTOP) + title = "dsquasar: {} terminated ({})".format(amsg, dmsg) + if self.PGBACK['errcnt']: title += " Error({})".format(self.PGBACK['errcnt']) + self.report_dscheck_email(title, 1) + os.kill(os.getpid(), signum) + # email a report for a run under dscheck control. the progress report is parked in # dscheck.einfo and the final report is sent right away. the run holds the check lock # while it works, and the dscheck daemon only mails unlocked records, so a parked @@ -857,8 +909,15 @@ def report_dscheck_email(self, title, cache = 0): ebuf = "From: {}\nTo: {}\n".format(sender, receiver) if self.PGLOG['CCDADDR']: ebuf += "Cc: {}\n".format(self.PGLOG['CCDADDR']) ebuf += "Subject: {}!\n\n{}\n".format(title, msg) - self.PGBACK['einfo'] = 1 # report the progress only once per run - return self.cache_customized_email("dscheck", "einfo", cnd, ebuf, self.LOGWRN) + # only a report that is actually parked (or sent directly by the fallback) counts + # as reported. marking it reported before the attempt turned a failed cache into a + # silent loss: nothing was parked, and the flag stopped every later attempt too + estat = self.cache_customized_email("dscheck", "einfo", cnd, ebuf, self.LOGWRN) + if estat: + self.PGBACK['einfo'] = 1 # report the progress only once per run + else: + self.PGBACK['eretry'] = tm() + self.ERETRY + return estat estat = self.build_customized_email("dscheck", "einfo", cnd, title, self.LOGWRN) # a failed send has already replaced the parked report via cache_customized_email() if estat == self.SUCCESS and self.PGBACK['einfo']: @@ -868,40 +927,48 @@ def report_dscheck_email(self, title, cache = 0): self.pglog("Cannot clean the progress report of {}; it is emailed again".format(cnd), self.LOGWRN) return estat - # guard the long PBS batch jobs (-A 2/3/4/6/7) against the walltime: once past MAXRUNTIME - # the run keeps going, but a progress report is parked in dscheck.einfo so that - # something is reported even if PBS kills the job at the walltime. the final report - # follows and drops the parked one if the run does finish in time. the live email - # buffers are saved and put back, so the final report still carries everything logged - # before the cutoff. - # both queue depths are reported rather than the current phase's, because -A 3 tars each - # input file right after creating it, so its status 'N' count stays near zero while that - # is the busy phase. they are worded as STATES ('left to build', 'tarred') and not as - # work this run will do: only -A 4 and -A 6 transfer, so calling the status 'T' count - # 'to transfer' in an -A 2 or -A 3 report claims work that run never performs. - # a no-op for command-line runs, other actions, before the cutoff, or once a progress - # report is cached. + # guard a long PBS batch job against the walltime: once past MAXRUNTIME the run keeps + # going, but a progress report is parked in dscheck.einfo so that something is reported + # even if PBS kills the job at the walltime. the final report follows and drops the + # parked one if the run does finish in time. the live email buffers are saved and put + # back, so the final report still carries everything logged before the cutoff. + # PBS sends no catchable warning before the kill - SIGTERM is not trapped and SIGKILL + # cannot be - so this parked report is the only thing standing between a killed job and + # a run that is never heard from. every action submitted to PBS is guarded, not just the + # tar and transfer ones: -A 16 and the hidden -A 128 are submitted with the same 24 hour + # walltime and used to be excluded, so they died silently. + # the tar queue depths are only reported for the actions they describe. both depths are + # reported rather than the current phase's, because -A 3 tars each input file right after + # creating it, so its status 'N' count stays near zero while that is the busy phase. they + # are worded as STATES ('left to build', 'tarred') and not as work this run will do: only + # -A 4 and -A 6 transfer, so calling the status 'T' count 'to transfer' in an -A 2 or + # -A 3 report claims work that run never performs. + # a no-op for command-line runs, before the cutoff, or once a progress report is cached. # called from the top of every per-item loop, not only where a tar is dispatched: a long # stretch of files that are skipped (already backed up, or not enough accumulated size to # tar yet) would otherwise walk past the cutoff without ever reaching the check, which is # how a 23h run reported nothing. def check_batch_deadline(self): - act = self.PGBACK['action'] if self.PGBACK['einfo'] or not (self.PGBACK['doemail'] and self.PGLOG['DSCHECK']): return - if self.PGLOG['CURBID'] < 1 or act not in (self.TARACT, self.CTACTS, self.BCKACT, self.TBACTS, self.CBACTS): return + if self.PGLOG['CURBID'] < 1: return # only a PBS batch run has a walltime to beat elapsed = tm() - self.PGBACK['starttime'] if elapsed < self.MAXRUNTIME: return - tcnt = self.batch_tar_count('N') - bcnt = self.batch_tar_count('T') + # a report that could not be parked is retried, but not once per file: the counts + # below are database queries and the cutoff leaves a whole hour of them otherwise + if self.PGBACK['eretry'] and tm() < self.PGBACK['eretry']: return + act = self.PGBACK['action'] etime = self.seconds_to_string_time(int(elapsed)) - amsg = self.ACTMSG[act] - bmsg = self.BACKMSG[self.PGBACK['backflag']] if self.PGBACK['backflag'] else 'backup' + amsg = self.action_message(act) dmsg = self.batch_done_count() - rmsg = "{} {} tar file(s) left to build and {} tarred".format(tcnt, bmsg, bcnt) - msg = "{}: Still running after {} of the {} PBS walltime, {}, with {}".format(amsg, etime, self.seconds_to_string_time(self.WALLTIME), dmsg, rmsg) + rmsg = '' + if act&self.CBACTS: # tar queue depths say nothing about the other actions + bmsg = self.BACKMSG[self.PGBACK['backflag']] if self.PGBACK['backflag'] else 'backup' + rmsg = ", with {} {} tar file(s) left to build and {} tarred".format( + self.batch_tar_count('N'), bmsg, self.batch_tar_count('T')) + msg = "{}: Still running after {} of the {} PBS walltime, {}{}".format(amsg, etime, self.seconds_to_string_time(self.WALLTIME), dmsg, rmsg) self.pglog(self.INDENT + msg, self.LOGACT) saved = {key : self.PGLOG[key] for key in ('EMLMSG', 'ERRMSG', 'ERRCNT', 'SUMMSG', 'PRGMSG')} - self.set_email("{}: {} still in progress after {}, {}, with {}!".format(self.PGBACK['cmd'], amsg, etime, dmsg, rmsg), self.EMLTOP) + self.set_email("{}: {} still in progress after {}, {}{}!".format(self.PGBACK['cmd'], amsg, etime, dmsg, rmsg), self.EMLTOP) title = "dsquasar: {} in progress ({})".format(amsg, dmsg) if self.PGBACK['errcnt']: title += " Error({})".format(self.PGBACK['errcnt']) self.report_dscheck_email(title, 1) @@ -921,6 +988,26 @@ def wait_all_children(self): self.show_wait_message(i, "{}: wait child processes".format(self.PGSIG['DSTR']), self.LOGWRN, 1) i += 1 + # wait for a free child process slot, checking the walltime deadline between polls. + # check_child(..., -1) does the waiting inside its own loop too, breaking only once fewer + # than MPROC children are left running, so a parent whose children each tar many GB sits + # there for as long as every slot stays busy - well past the cutoff - while the deadline + # check right below the call is never reached. That is the other half of the problem + # wait_all_children() fixed: that one guarded the terminal wait, this one guards the far + # more frequent wait for a slot. Polling with dowait 0 keeps check_child's own cadence and + # wait message, and returns the free slot count exactly as check_child(..., -1) does. + def wait_child_slot(self): + if self.PGBACK['mproc'] < 2: return 0 + i = 0 + while True: + self.check_batch_deadline() + pcnt = self.check_child(None, 0, self.LOGWRN, 0) + ccnt = self.PGSIG['MPROC'] - pcnt + if ccnt > 0: return ccnt + self.show_wait_message(i, "{}: wait {}/{} child processes".format( + self.PGSIG['DSTR'], pcnt, self.PGSIG['MPROC']), self.LOGWRN, 1) + i += 1 + # recompute the confirmed backup counts from RDADB after all child processes # finished, so a multi-process summary reflects succeeded (not just started) work def confirm_quasar_counts(self, qinfo, bids, dcnd): @@ -955,7 +1042,7 @@ def duplicate_backup_file(self, dsid, qfile): # backup one Quasar Backup or Backup&Drdata from one or multiple inputs and, # reset the quasar backup dict def process_one_backup_file(self, qinfo, addback, keepid = False): - ccnt = self.check_child(None, 0, self.LOGWRN, -1) if self.PGBACK['mproc'] > 1 else 0 + ccnt = self.wait_child_slot() if self.PGSIG['QUIT']: self.quit_dsquasar(qinfo) self.check_batch_deadline() dsids = qinfo['dsids'] @@ -1035,7 +1122,7 @@ def process_one_backup_file(self, qinfo, addback, keepid = False): # Transfer multiple tarfiles to Quasar Backup or Backup&Drdata, and # reset the quasar backup dict def transfer_quasar_tarfiles(self, qinfo): - ccnt = self.check_child(None, 0, self.LOGWRN, -1) if self.PGBACK['mproc'] > 1 else 0 + ccnt = self.wait_child_slot() if self.PGSIG['QUIT']: self.quit_dsquasar(qinfo) self.check_batch_deadline() # prepare for backup one tar file @@ -1623,6 +1710,7 @@ def process_dataset_chkfiles(self, dsfiles, backflag): qinfo = {'backflag' : backflag, 'bid' : 0, 'dsid' : None, 'size' : 0, 'bfile' : None, 'bqfiles' : {}, 'dqfiles' : {}, 'qdsids' : [], 'qcnt' : 0, 'ncnt' : 0} for bid in bfiles: + self.check_batch_deadline() qinfo['bid'] = bid binfo = bfiles[bid] if qinfo['dsid'] and binfo['dsid'] != qinfo['dsid']: @@ -1741,12 +1829,13 @@ def process_dataset_pathfiles(self, dsfiles, backflag): qinfo = {'backflag' : backflag, 'bid' : 0, 'dsids' : [], 'fcnt' : 0, 'size' : 0, 'bfile' : None, 'pfile' : None, 'qdsids' : [], 'qfcnt' : 0, 'qsize' : 0, 'qcnt' : 0} for bid in bfiles: + self.check_batch_deadline() qinfo['bid'] = bid binfo = bfiles[bid] for bkey in binfo: qinfo[bkey] = binfo[bkey] self.process_one_quasar_pathfile(qinfo) if self.PGBACK['mproc'] > 1: - self.check_child(None, 0, self.LOGWRN, 1) # wait all child processes done + self.wait_all_children() # wait all child processes done self.confirm_quasar_counts(qinfo, list(bfiles), "bfile LIKE 'G%/%.tar'") # recount confirmed renames from RDADB qcnt = qinfo['qcnt'] if qcnt > 0: @@ -1761,7 +1850,7 @@ def process_dataset_pathfiles(self, dsfiles, backflag): # backup one Quasar Backup or Backup&Drdata from one or multiple inputs and, # reset the quasar backup dict def process_one_quasar_pathfile(self, qinfo): - ccnt = self.check_child(None, 0, self.LOGWRN, -1) if self.PGBACK['mproc'] > 1 else 0 + ccnt = self.wait_child_slot() if self.PGSIG['QUIT']: self.quit_dsquasar(qinfo) # prepare for backup one tar file dsids = qinfo['dsids'] @@ -1883,6 +1972,7 @@ def process_dataset_dsidfiles(self, dsfiles, backflag): qinfo = {'backflag' : backflag, 'bid' : 0, 'dsids' : [], 'fcnt' : 0, 'size' : 0, 'qdsids' : [], 'qfcnt' : 0, 'qsize' : 0, 'qcnt' : 0} for bid in bfiles: + self.check_batch_deadline() qinfo['bid'] = bid binfo = bfiles[bid] for bkey in binfo: qinfo[bkey] = binfo[bkey] @@ -1981,6 +2071,7 @@ def process_dataset_mcsfiles(self, dsfiles, backflag): qinfo = {'dsids' : [], 'fcnt' : 0, 'mcnt' : 0, 'fsize' : 0, 'msize' : 0} qcnt = 0 for bid in bfiles: + self.check_batch_deadline() qcnt += self.process_one_quasar_mcsfile(bid, bfiles[bid], qinfo) if qcnt > 0: s = 's' if qcnt > 1 else ''