diff --git a/README.md b/README.md index db86493..ab7e4aa 100644 --- a/README.md +++ b/README.md @@ -72,7 +72,7 @@ PgLOG.pglog("hello", PgLOG.LOGWRN) python -c "import rda_python_common; print(rda_python_common.__version__)" ``` -You should see the installed version (currently `2.1.10`). If the import +You should see the installed version (currently `2.1.11`). 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 09be99e..798f2a6 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "rda_python_common" -version = "2.1.10" +version = "2.1.11" authors = [ { name="Zaihua Ji", email="zji@ucar.edu" }, ] @@ -18,8 +18,8 @@ classifiers = [ "Development Status :: 5 - Production/Stable", ] dependencies = [ - "hvac", - "psycopg2==2.9.10", + "psycopg2-binary", + "psutil", "rda-python-globus", "unidecode", "hvac" diff --git a/requirements.txt b/requirements.txt index 31ae92d..7f4a1fe 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,7 +1,8 @@ iniconfig==2.1.0 packaging==24.2 pluggy==1.5.0 -psycopg2==2.9.10 +psycopg2-binary==2.9.10 +psutil pytest==8.3.5 rda-python-globus unidecode diff --git a/src/rda_python_common/__init__.py b/src/rda_python_common/__init__.py index 7ff6e1f..8c0a536 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__ = "2.1.10" +__version__ = "2.1.11" __all__ = [ "PgLOG", diff --git a/src/rda_python_common/pg_file.py b/src/rda_python_common/pg_file.py index d320a1c..068d8ef 100644 --- a/src/rda_python_common/pg_file.py +++ b/src/rda_python_common/pg_file.py @@ -19,6 +19,7 @@ import time import glob import json +import hashlib from .pg_util import PgUtil from .pg_sig import PgSIG @@ -2874,24 +2875,32 @@ def get_md5sum(self, file, count = 0, logact = 0): (with None for missing files) for multiple files, or None on failure. """ - cmd = 'md5sum ' if count > 0: checksum = [None]*count for i in range(count): if op.isfile(file[i]): - chksm = self.pgsystem(cmd + file[i], logact, 20) - if chksm: - ms = re.search(r'(\w{32})', chksm) - if ms: checksum[i] = ms.group(1) + checksum[i] = self._file_md5(file[i], logact) else: checksum = None if op.isfile(file): - chksm = self.pgsystem(cmd + file, logact, 20) - if chksm: - ms = re.search(r'(\w{32})', chksm) - if ms: checksum = ms.group(1) + checksum = self._file_md5(file, logact) return checksum + def _file_md5(self, path, logact=0): + """Compute MD5 hex digest of *path*, reading in 1 MiB chunks. + + Returns the hex digest string, or None on read error. + """ + try: + h = hashlib.md5() + with open(path, 'rb') as fh: + for chunk in iter(lambda: fh.read(1048576), b''): + h.update(chunk) + return h.hexdigest() + except OSError as e: + self.pglog("Error md5sum {}: {}".format(path, str(e)), logact) + return None + # Evaluate md5 checksums and compare them for two given files # file1, file2: file names # Return: 0 if same and 1 if not diff --git a/src/rda_python_common/pg_log.py b/src/rda_python_common/pg_log.py index b96e44d..ffc916f 100644 --- a/src/rda_python_common/pg_log.py +++ b/src/rda_python_common/pg_log.py @@ -18,6 +18,7 @@ import shlex import smtplib from email.message import EmailMessage +import subprocess from subprocess import Popen, PIPE from os import path as op import time @@ -769,7 +770,7 @@ def show_usage(self, progname, opts=None): nilcnt = 0 if begin: sys.stdout.write(line) else: - os.system("more " + usgname) + subprocess.run(['more', usgname]) self.pgexit(0) def err2std(self, line): @@ -1600,15 +1601,12 @@ def set_specialist_home(self, specialist): os.environ['MAIL'] = re.sub(self.PGLOG['CURUID'], specialist, os.environ['MAIL']) home = "{}/{}".format(self.PGLOG['USRHOME'], specialist) shell = "tcsh" - buf = self.pgsystem("grep ^{}: /etc/passwd".format(specialist), self.LOGWRN, 20) - if buf: - lines = buf.split('\n') - for line in lines: - ms = re.search(r':(/.+):(/.+)', line) - if ms: - home = ms.group(1) - shell = op.basename(ms.group(2)) - break + try: + pwent = pwd.getpwnam(specialist) + home = pwent.pw_dir + shell = op.basename(pwent.pw_shell) + except KeyError: + pass if home != os.environ['HOME'] and op.exists(home): os.environ['HOME'] = home return shell diff --git a/src/rda_python_common/pg_sig.py b/src/rda_python_common/pg_sig.py index 1b83da1..d84eef1 100644 --- a/src/rda_python_common/pg_sig.py +++ b/src/rda_python_common/pg_sig.py @@ -14,6 +14,8 @@ import errno import signal import time +import subprocess +import psutil from contextlib import contextmanager from .pg_dbi import PgDBI @@ -204,6 +206,42 @@ def stop_daemon(self, msg): self.PGLOG['LOGMASK'] |= self.MSGLOG # turn on logging before daemon stops self.pglog("{} Started at {}, Stopped gracefully{} by {}".format(self.PGSIG['DSTR'], self.PGSIG['STRTM'], msg, self.current_datetime()), self.LOGWRN) + # scan running processes via psutil and return a list of dicts matching aname/uname + # Mirrors the semantics of the previous "ps -u U -f | grep A" / "ps -C A -f" pipelines: + # - with uname: substring match of aname anywhere in cmdline (looser, like grep) + # - without uname: aname matches the executable basename (with optional .ext) + # Each entry has keys: 'pid', 'ppid', 'args' ('args' is cmdline[1:] joined). + def _scan_app_processes(self, aname, uname=None): + """Scan running processes matching an application name, optionally by user. + + Args: + aname (str): Application name (basename or substring of cmdline). + uname (str, optional): Username filter. + + Returns: + list[dict]: Each dict has keys ``pid``, ``ppid``, ``args``. + """ + results = [] + for proc in psutil.process_iter(['pid', 'ppid', 'username', 'cmdline']): + try: + info = proc.info + cmdline = info.get('cmdline') or [] + if not cmdline: continue + if uname is not None and info.get('username') != uname: continue + if uname: + if not any(aname in arg for arg in cmdline): continue + else: + exe = os.path.basename(cmdline[0]) + if exe != aname and re.sub(r'\.\w+$', '', exe) != aname: continue + results.append({ + 'pid': info['pid'], + 'ppid': info['ppid'], + 'args': ' '.join(cmdline[1:]), + }) + except (psutil.NoSuchProcess, psutil.AccessDenied): + continue + return results + # check if a daemon is running already # aname - application name for the daemon # uname - user login name who started the daemon @@ -218,21 +256,11 @@ def check_daemon(self, aname, uname=None): Returns: int: The process ID of the running daemon, or 0 if not running. """ - if uname: - self.check_vuser(uname, aname) - pcmd = "ps -u {} -f | grep {} | grep ' 1 '".format(uname, aname) - mp = r"^\s*{}\s+(\d+)\s+1\s+".format(uname) - else: - pcmd = "ps -C {} -f | grep ' 1 '".format(aname) - mp = r"^\s*\w+\s+(\d+)\s+1\s+" - buf = self.pgsystem(pcmd, self.LOGWRN, 20+1024) - if buf: - cpid = os.getpid() - lines = buf.split('\n') - for line in lines: - ms = re.match(mp, line) - pid = int(ms.group(1)) if ms else 0 - if pid > 0 and pid != cpid: return pid + if uname: self.check_vuser(uname, aname) + cpid = os.getpid() + for p in self._scan_app_processes(aname, uname): + if p['ppid'] != 1: continue + if p['pid'] != cpid: return p['pid'] return 0 # check if an application is running already; other than the current processs @@ -251,31 +279,21 @@ def check_application(self, aname, uname=None, sargv=None): Returns: int: The process ID of the running instance, or 0 if not found. """ - if uname: - self.check_vuser(uname, aname) - pcmd = "ps -u {} -f | grep {} | grep -v ' grep '".format(uname, aname) - mp = r"^\s*{}\s+(\d+)\s+(\d+)\s+.*{}\S*\s+(.*)$".format(uname, aname) - else: - pcmd = "ps -C {} -f".format(aname) - mp = r"^\s*\w+\s+(\d+)\s+(\d+)\s+.*{}\S*\s+(.*)$".format(aname) - buf = self.pgsystem(pcmd, self.LOGWRN, 20+1024) - if not buf: return 0 + if uname: self.check_vuser(uname, aname) + procs = self._scan_app_processes(aname, uname) + if not procs: return 0 cpids = [os.getpid(), os.getppid()] pids = [] ppids = [] astrs = [] - lines = buf.split('\n') - for line in lines: - ms = re.match(mp, line) - if not ms: continue - pid = int(ms.group(1)) - ppid = int(ms.group(2)) + for p in procs: + pid, ppid = p['pid'], p['ppid'] if pid in cpids: if ppid not in cpids: cpids.append(ppid) continue pids.append(pid) ppids.append(ppid) - if sargv: astrs.append(ms.group(3)) + if sargv: astrs.append(p['args']) pcnt = len(pids) if not pcnt: return 0 i = 0 @@ -329,25 +347,15 @@ def check_multiple_application(self, aname, uname=None, sargv=None): Returns: int: Number of running instances (excluding the current process). """ - if uname: - self.check_vuser(uname, aname) - pcmd = "ps -u {} -f | grep {} | grep -v ' grep '".format(uname, aname) - mp = r"^\s*{}\s+(\d+)\s+(\d+)\s+.*{}\S*\s+(.*)$".format(uname, aname) - else: - pcmd = "ps -C {} -f".format(aname) - mp = r"^\s*\w+\s+(\d+)\s+(\d+)\s+.*{}\S*\s+(.*)$".format(aname) - buf = self.pgsystem(pcmd, self.LOGWRN, 20+1024) - if not buf: return 0 + if uname: self.check_vuser(uname, aname) + procs = self._scan_app_processes(aname, uname) + if not procs: return 0 dpids = [os.getpid(), os.getppid()] pids = [] ppids = [] astrs = [] - lines = buf.split('\n') - for line in lines: - ms = re.match(mp, line) - if not ms: continue - pid = int(ms.group(1)) - ppid = int(ms.group(2)) + for p in procs: + pid, ppid = p['pid'], p['ppid'] if pid in dpids: if ppid > 1 and ppid not in dpids: dpids.append(ppid) continue @@ -356,7 +364,7 @@ def check_multiple_application(self, aname, uname=None, sargv=None): continue pids.append(pid) ppids.append(ppid) - if sargv: astrs.append(ms.group(3)) + if sargv: astrs.append(p['args']) pcnt = len(pids) if not pcnt: return 0 i = 0 @@ -627,18 +635,20 @@ def kill_children(self, pid, logact=None): list: PIDs of processes that were successfully killed. """ if logact is None: logact = self.LOGWRN - buf = self.pgsystem("ps --ppid {} -o pid".format(pid), logact, 20) pids = [] - if buf: - lines = buf.split('\n') - for line in lines: - ms = re.match(r'^\s*(\d+)', line) - if not ms: continue - cid = int(ms.group(1)) - if not self.check_process(cid): continue - cids = self.kill_children(cid, logact) - if cids: pids = cids + pids - if self.kill_process(cid, signal.SIGKILL, logact) == self.SUCCESS: pids.insert(0, cid) + try: + children = psutil.Process(pid).children() + except psutil.NoSuchProcess: + children = [] + except Exception as e: + self.pglog("Error listing children of pid {}: {}".format(pid, str(e)), logact) + children = [] + for child in children: + cid = child.pid + if not self.check_process(cid): continue + cids = self.kill_children(cid, logact) + if cids: pids = cids + pids + if self.kill_process(cid, signal.SIGKILL, logact) == self.SUCCESS: pids.insert(0, cid) if logact and len(pids): self.pglog("Process({}) Killed".format(','.join(map(str, pids))), logact) return pids @@ -808,13 +818,11 @@ def check_process(self, pid): Returns: int: 1 if the process is running, 0 otherwise. """ - buf = self.pgsystem("ps -p {} -o pid".format(pid), self.LGWNEX, 20) - if buf: - mp = r'^\s*{}$'.format(pid) - lines = buf.split('\n') - for line in lines: - if re.match(mp, line): return 1 - return 0 + try: + os.kill(pid, 0) + except OSError: + return 0 + return 1 # check a process id on give host def check_host_pid(self, host, pid, pmsg=None, logact=None): @@ -1092,7 +1100,10 @@ def start_background(self, cmd, logact=None, cmdopt=5, dowait=0): self.PGLOG['ERRFILE'] = re.sub(r'\.log$', '.err', self.PGLOG['LOGFILE'], 1) bckcmd += " 2>> {}/{}".format(self.PGLOG['LOGPATH'], self.PGLOG['ERRFILE']) bckcmd += " &" - os.system(bckcmd) + # shell=True is required for the redirections (>> / 2>>) and trailing '&'; + # the '&' makes the shell fork the command and exit, so the command gets + # reparented to init (ppid=1), matching the lookup logic in record_background(). + subprocess.Popen(bckcmd, shell=True) return self.record_background(cmd, logact) # get background process id for given bcmd @@ -1170,27 +1181,31 @@ def record_background(self, bcmd, logact=None): """ if logact is None: logact = self.LOGWRN ms = re.match(r'^(\S+)', bcmd) - if ms: - aname = ms.group(1) - else: - aname = bcmd - mp = r"^\s*(\S+)\s+(\d+)\s+1\s+.*{}(.*)$".format(aname) - pc = "ps -u {},{} -f | grep ' 1 ' | grep {}".format(self.PGLOG['CURUID'], self.PGLOG['GDEXUSER'], aname) + aname = ms.group(1) if ms else bcmd + curuid = self.PGLOG['CURUID'] + gdexuser = self.PGLOG['GDEXUSER'] for i in range(2): - buf = self.pgsystem(pc, logact, 20+1024) - if buf: - lines = buf.split('\n') - for line in lines: - ms = re.match(mp, line) - if not ms: continue - (uid, sbid, acmd) = ms.groups() - bid = int(sbid) + for proc in psutil.process_iter(['pid', 'ppid', 'username', 'cmdline']): + try: + info = proc.info + if info.get('ppid') != 1: continue + uid = info.get('username') + if uid != curuid and uid != gdexuser: continue + cmdline = info.get('cmdline') or [] + if not cmdline: continue + line = ' '.join(cmdline) + idx = line.find(aname) + if idx < 0: continue + bid = info['pid'] if bid in self.CBIDS: return -1 - if uid == self.PGLOG['GDEXUSER']: + acmd = line[idx+len(aname):] + if uid == gdexuser: acmd = re.sub(r'^\.(pl|py)\s+', '', acmd, 1) if re.match(r'^{}{}'.format(aname, acmd), bcmd): continue self.CBIDS[bid] = bcmd return 1 + except (psutil.NoSuchProcess, psutil.AccessDenied): + continue time.sleep(2) return 0