From 5846a75a10449bdf97f684bcb1948c363a88e7d1 Mon Sep 17 00:00:00 2001 From: zaihuaji Date: Wed, 30 Sep 2026 09:01:27 -0500 Subject: [PATCH] report one error per failed Globus transfer and hold a restart until the parked report is mailed; bump version to 3.0.19 pg_file.py: a cancelled transfer was logged twice - once as the full 17-line dsglobus get-task dump, once as the terse status line - so three failures read as six errors in the email. The cancel reason is now condensed to one line, cached in QCANCEL and appended to the single error that names the file. pg_cmd.py: init_dscheck restarted a record whose einfo still held a progress report, stranding it since the dscheck daemon skips running records. It now leaves the record for the daemon to mail first. Co-Authored-By: Claude Opus 4.6 --- README.md | 2 +- pyproject.toml | 2 +- src/rda_python_common/__init__.py | 2 +- src/rda_python_common/pg_cmd.py | 8 ++++++++ src/rda_python_common/pg_file.py | 26 +++++++++++++++++++++----- 5 files changed, 32 insertions(+), 8 deletions(-) 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