diff --git a/README.md b/README.md index 93bc659..c980135 100644 --- a/README.md +++ b/README.md @@ -165,7 +165,7 @@ PgLOG.pglog("hello", PgLOG.LOGWRN) python -c "import rda_python_common; print(rda_python_common.__version__)" ``` -You should see the installed version (currently `3.0.18`). If the import +You should see the installed version (currently `3.0.19`). If the import fails, double-check that the active Python environment is the one where you ran `pip install`. diff --git a/pyproject.toml b/pyproject.toml index 7d6b86f..801693c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "rda_python_common" -version = "3.0.18" +version = "3.0.19" authors = [ { name="Zaihua Ji", email="zji@ucar.edu" }, ] diff --git a/src/rda_python_common/__init__.py b/src/rda_python_common/__init__.py index 30843dc..c170f19 100644 --- a/src/rda_python_common/__init__.py +++ b/src/rda_python_common/__init__.py @@ -22,7 +22,7 @@ from . import PgLOG, PgUtil, PgDBI, PgFile, PgLock, PgCMD, PgSIG, PgOPT, PgSplit -__version__ = "3.0.18" +__version__ = "3.0.19" __all__ = [ "PgLOG", diff --git a/src/rda_python_common/pg_cmd.py b/src/rda_python_common/pg_cmd.py index 28cea85..2473236 100644 --- a/src/rda_python_common/pg_cmd.py +++ b/src/rda_python_common/pg_cmd.py @@ -207,6 +207,14 @@ def init_dscheck(self, oindex, otype, cmd, dsid, action, workdir=None, specialis self.pglog(cmsg + "is Running, No restart", self.LOGWRN) sys.exit(0) if cidx > 0: + # a report parked in einfo by the previous run is only mailed by the dscheck + # daemon, which skips any record that is already running. restarting the record + # here would strand that report until the next run overwrites it, so leave the + # record alone and let the daemon mail it first - it clears einfo when it does + if pgrec['einfo']: + self.pglog(cmsg + "has a report to email, No restart", self.LOGWRN) + self.lock_dscheck(cidx, 0, logact) + sys.exit(0) if not hosts and pgrec['hostname']: hosts = pgrec['hostname'] self.set_one_boption('hostname', hosts, 0) diff --git a/src/rda_python_common/pg_file.py b/src/rda_python_common/pg_file.py index e63cf7c..a94a811 100644 --- a/src/rda_python_common/pg_file.py +++ b/src/rda_python_common/pg_file.py @@ -105,7 +105,8 @@ def __init__(self): } self.TARSTR = '|'.join(self.PGTARS) self.DELDIRS = {} - self.TASKIDS = {} # cache unfinished + self.TASKIDS = {} # cache unfinished + self.QCANCEL = {} # taskid -> why a task got cancelled, reported by the waiting caller self.LHOST = "localhost" self.OHOST = self.PGLOG['OBJCTSTR'] self.BHOST = self.PGLOG['BACKUPNM'] @@ -532,7 +533,11 @@ def submit_globus_task(self, cmd, endpoint, logact = 0, qstr = None): time.sleep(self.PGSIG['ETIME']) lp += 1 if task['stat'] == 'S' or task['stat'] == 'A': break - if task['stat'] == 'F' and not syserr: break + if task['stat'] == 'F' and not syserr: + # nothing waits on this task, so report the cancellation here instead + if task['id'] in self.QCANCEL: + self.errlog("{}: Cancel Task due to {}".format(task['id'], self.QCANCEL.pop(task['id'])), 'B', 1, logact) + break errmsg = "Error Execute: " + cmd if qstr: errmsg += " with stdin:\n" + qstr if syserr: @@ -581,8 +586,16 @@ def check_globus_status(self, taskid, endpoint = None, logact = 0): detail = ms.group(1) if detail not in astats: if logact&self.NOWAIT: - errmsg = "{}: Cancel Task due to {}:\n{}".format(taskid, detail, buf) - self.errlog(errmsg, 'B', 1, logact) + # record why, and let the caller waiting on this task report it as a + # single error naming the file; dumping the whole get-task output here + # doubled every failure into two error entries, the first 17 lines long + reason = detail + ms = re.search(r'Bytes Transferred:\s+(\d+)', buf) + if ms: reason += " after " + self.format_float_value(ms.group(1)) + ms = re.search(r'Files:\s+(\d+)', buf) + if ms: reason += " of {} file(s)".format(ms.group(1)) + self.QCANCEL[taskid] = reason + self.pglog("{}: Cancel Task due to {}".format(taskid, reason), self.LOGWRN) ccmd = f"{bcmd} cancel-task {taskid}" self.pgsystem(ccmd, logact, 7) else: @@ -640,7 +653,10 @@ def check_globus_finished(self, tofile, topoint, logact = 0): del self.TASKIDS[ckey] else: status = self.QSTATS[stat] if stat in self.QSTATS else 'UNKNOWN' - self.errlog("{}: Status '{}' for Task {}".format(ckey, status, taskid), 'B', 1, logact) + errmsg = "{}: Status '{}' for Task {}".format(ckey, status, taskid) + if taskid in self.QCANCEL: + errmsg += " - " + self.QCANCEL.pop(taskid) + self.errlog(errmsg, 'B', 1, logact) ret = self.FAILURE break return ret