Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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" },
]
Expand Down
155 changes: 123 additions & 32 deletions src/rda_python_dsquasar/dsquasar.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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')
Expand Down Expand Up @@ -791,15 +795,15 @@ 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)
qcnt = qinfo['qcnt']
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)
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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']:
Expand All @@ -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)
Expand All @@ -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):
Expand Down Expand Up @@ -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']
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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']:
Expand Down Expand Up @@ -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:
Expand All @@ -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']
Expand Down Expand Up @@ -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]
Expand Down Expand Up @@ -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 ''
Expand Down
Loading