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.20"
version = "3.0.21"
authors = [
{ name="Zaihua Ji", email="zji@ucar.edu" },
]
Expand Down
56 changes: 21 additions & 35 deletions src/rda_python_dsquasar/dsquasar.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@
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 @@ -203,7 +202,7 @@ def start_actions(self):
self.cmdlog(self.PGBACK['cmd'])
self.PGBACK['starttime'] = tm()
self.set_walltime_deadline()
self.catch_batch_termination()
self.catch_term_signal(self.batch_term_report)
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 @@ -852,42 +851,29 @@ def set_walltime_deadline(self):
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.
# park a progress report for a run PBS is about to kill. catch_term_signal() arms this
# for a batch run only, restores the default handler before calling it, and re-raises
# the signal afterwards, so all this has to do is build the report.
# it runs on borrowed time (SIGKILL follows SIGTERM by the MoM's kill_delay, 10 seconds
# by default), 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 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)
def batch_term_report(self, signum, frame):
if self.PGBACK['einfo'] or not (self.PGBACK['doemail'] and self.PGLOG['DSCHECK']): return
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)

# 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
Expand Down
Loading