diff --git a/python/pacemaker/_cts/audits.py b/python/pacemaker/_cts/audits.py index 2634d6a0878..7d89b5cfb76 100644 --- a/python/pacemaker/_cts/audits.py +++ b/python/pacemaker/_cts/audits.py @@ -4,7 +4,10 @@ __copyright__ = "Copyright 2000-2026 the Pacemaker project contributors" __license__ = "GNU General Public License version 2 or later (GPLv2+) WITHOUT ANY WARRANTY" +import glob +import os import re +import subprocess import time import uuid @@ -127,11 +130,14 @@ def _test_logging(self): kinds = [LogKind.LOCAL_FILE] if self._cm.env["have_systemd"]: kinds.append(LogKind.JOURNAL) + kinds.append(LogKind.REMOTE_FILE) for k in kinds: watch[k] = self._create_watcher(patterns, k) + logging.log(f"Logging test message with identifier {suffix}") + else: watch[watch_pref] = self._create_watcher(patterns, watch_pref) @@ -144,10 +150,12 @@ def _test_logging(self): logging.log(f"Checking for test message in {k} logs") w.look_for_all(silent=True) + if not w.unmatched: if watch_pref is None: logging.log(f"Found test message in {k} logs") self._cm.env["log_kind"] = k + return True for regex in w.unmatched: @@ -159,7 +167,6 @@ def __call__(self): """Perform the audit action.""" max_attempts = 3 attempt = 0 - passed = True self._cm.ns.wait_for_all_nodes(self._cm.env["nodes"]) while attempt <= max_attempts and not self._test_logging(): @@ -169,16 +176,13 @@ def __call__(self): if attempt > max_attempts: logging.log("ERROR: Cluster logging unrecoverable.") - passed = False + return False - return passed + return True def is_applicable(self): """Return True if this audit is applicable in the current test configuration.""" - if self._cm.env["LogAuditDisabled"]: - return False - - return True + return not self._cm.env["LogAuditDisabled"] class DiskAudit(ClusterAudit): @@ -207,7 +211,7 @@ def __call__(self): passed = True # @TODO Use directory of PCMK_logfile if set on host - dfcmd = "df -BM %s | tail -1 | awk '{print $(NF-1)\" \"$(NF-2)}' | tr -d 'M%%'" % BuildOptions.LOG_DIR + dfcmd = f"df -BM {BuildOptions.LOG_DIR} | tail -1 | awk '{{print $(NF-1)\" \"$(NF-2)}}' | tr -d 'M%%'" self._cm.ns.wait_for_all_nodes(self._cm.env["nodes"]) for node in self._cm.env["nodes"]: @@ -223,7 +227,7 @@ def __call__(self): used_percent = int(used) remaining_mb = int(remain) except (ValueError, TypeError): - logging.log(f"Warning: df output '{dfout}' from {node} was invalid [{used}, {remain}]") + logging.log(f"Warning: df output '{dfout}' from {node} was invalid") else: if remaining_mb < 10 or used_percent > 95: logging.log(f"CRIT: Out of log disk space on {node} ({used_percent}% / {remaining_mb}MB)") @@ -280,6 +284,9 @@ def _output_has_core(self, output, node): def _find_core_with_coredumpctl(self, node): """Use coredumpctl to find core dumps on the given node.""" + if not self._cm.env["have_systemd"]: + return False + (_, lsout) = self._cm.rsh.call(node, "coredumpctl --no-legend --no-pager") return self._output_has_core(lsout, node) @@ -296,13 +303,9 @@ def __call__(self): self._cm.ns.wait_for_all_nodes(self._cm.env["nodes"]) for node in self._cm.env["nodes"]: - found = False - - # If systemd is present, first see if coredumpctl logged any core dumps. - if self._cm.env["have_systemd"]: - found = self._find_core_with_coredumpctl(node) - if found: - passed = False + # First, try to use coredumpctl to find any core dumps. + if self._find_core_with_coredumpctl(node): + passed = False # If we didn't find any core dumps, it's for one of three reasons: # (1) Nothing crashed @@ -310,31 +313,27 @@ def __call__(self): # (3) systemd is present but coredumpctl is not enabled # # To handle the last two cases, check the other filesystem locations. - if not found: - found = self._find_core_on_fs(node, ["/var/lib/pacemaker/cores/*", - "/var/lib/corosync"]) - if found: - passed = False + elif self._find_core_on_fs(node, ["/var/lib/pacemaker/cores/*", + "/var/lib/corosync"]): + passed = False - if self._cm.expected_status.get(node) == "down": - clean = False - (_, lsout) = self._cm.rsh.call(node, "ls -al /dev/shm | grep qb-", verbose=1) + if self._cm.expected_status.get(node) != "down": + logging.debug(f"Skipping {node}") + continue + + (_, lsout) = self._cm.rsh.call(node, "ls -al /dev/shm | grep qb-", verbose=1) + if lsout: + passed = False for line in lsout: - passed = False - clean = True logging.log(f"Warning: Stale IPC file on {node}: {line}") - if clean: - (_, lsout) = self._cm.rsh.call(node, "ps axf | grep -e pacemaker -e corosync", verbose=1) + (_, lsout) = self._cm.rsh.call(node, "ps axf | grep -e pacemaker -e corosync", verbose=1) - for line in lsout: - logging.debug(f"ps[{node}]: {line}") + for line in lsout: + logging.debug(f"ps[{node}]: {line}") - self._cm.rsh.call(node, "rm -rf /dev/shm/qb-*") - - else: - logging.debug(f"Skipping {node}") + self._cm.rsh.call(node, "rm -rf /dev/shm/qb-*") return passed @@ -367,7 +366,7 @@ def __init__(self, cm, line): self.rclass = fields[6] self.rtype = fields[7] self.host = fields[8] - self.needs_quorum = fields[9] + self.needs_quorum = fields[9] == "1" self.flags = int(fields[10]) self.flags_s = fields[11] @@ -377,17 +376,17 @@ def __init__(self, cm, line): @property def unique(self): """Return True if this resource is unique.""" - return self.flags & 0x20 + return bool(self.flags & 0x20) @property def orphan(self): """Return True if this resource is an orphan.""" - return self.flags & 0x01 + return bool(self.flags & 0x01) @property def managed(self): """Return True if this resource is managed by the cluster.""" - return self.flags & 0x02 + return bool(self.flags & 0x02) class AuditConstraint: @@ -440,12 +439,7 @@ def __init__(self, cm): """ ClusterAudit.__init__(self, cm) self.name = "PrimitiveAudit" - - self._active_nodes = [] - self._constraints = [] - self._inactive_nodes = [] - self._resources = [] - self._target = None + self._reset() def _audit_resource(self, resource, quorum): """Perform the audit of a single resource.""" @@ -456,7 +450,7 @@ def _audit_resource(self, resource, quorum): if quorum: self.debug(f"Resource {resource.id} active on {active!r}") - elif resource.needs_quorum == 1: + elif resource.needs_quorum: logging.log(f"Resource {resource.id} active without quorum: {active!r}") rc = False @@ -487,22 +481,32 @@ def _audit_resource(self, resource, quorum): return rc + # pylint: disable=attribute-defined-outside-init + def _reset(self): + """Reset internal lists.""" + self._active_nodes = [] + self._constraints = [] + self._inactive_nodes = [] + self._resources = [] + self._target = None + def _setup(self): """ Verify cluster nodes are active. Collect resource and colocation information used for performing the audit. """ + self._reset() + for node in self._cm.env["nodes"]: if self._cm.expected_status[node] == "up": self._active_nodes.append(node) + + if self._target is None: + self._target = node else: self._inactive_nodes.append(node) - for node in self._cm.env["nodes"]: - if self._target is None and self._cm.expected_status[node] == "up": - self._target = node - if not self._target: # TODO: In Pacemaker 1.0 clusters we'll be able to run crm_resource # with CIB_file=/path/to/cib.xml even when the cluster isn't running @@ -513,9 +517,9 @@ def _setup(self): verbose=1) for line in lines: - if re.search("^Resource", line): + if line.startswith("Resource"): self._resources.append(AuditResource(self._cm, line)) - elif re.search("^Constraint", line): + elif line.startswith("Constraint"): self._constraints.append(AuditConstraint(self._cm, line)) else: logging.log(f"Unknown entry: {line}") @@ -540,13 +544,7 @@ def __call__(self): def is_applicable(self): """Return True if this audit is applicable in the current test configuration.""" - # @TODO Due to long-ago refactoring, this name test would never match, - # so this audit (and those derived from it) would never run. - # Uncommenting the next lines fixes the name test, but that then - # exposes pre-existing bugs that need to be fixed. - # if self._cm.name == "crm-corosync": - # return True - return False + return self._cm.name == "crm-corosync" class GroupAudit(PrimitiveAudit): @@ -671,12 +669,14 @@ def _crm_location(self, resource): (rc, lines) = self._cm.rsh.call(self._target, f"crm_resource --locate -r {resource} -Q", verbose=1) + if rc != 0: + return [] + hosts = [] - if rc == 0: - for line in lines: - fields = line.split() - hosts.append(fields[0]) + for line in lines: + fields = line.split() + hosts.append(fields[0]) return hosts @@ -695,15 +695,16 @@ def __call__(self): if not source: self.debug(f"Colocation audit ({coloc.id}): {coloc.rsc} not running") - else: - for node in source: - if node not in target: - passed = False - logging.log(f"Colocation audit ({coloc.id}): {coloc.rsc} running " - f"on {node} (not in {target!r})") - else: - self.debug(f"Colocation audit ({coloc.id}): {coloc.rsc} running " - f"on {node} (in {target!r})") + continue + + for node in source: + if node not in target: + passed = False + logging.log(f"Colocation audit ({coloc.id}): {coloc.rsc} running " + f"on {node} (not in {target!r})") + else: + self.debug(f"Colocation audit ({coloc.id}): {coloc.rsc} running " + f"on {node} (in {target!r})") return passed @@ -760,13 +761,7 @@ def __call__(self): def is_applicable(self): """Return True if this audit is applicable in the current test configuration.""" - # @TODO Due to long-ago refactoring, this name test would never match, - # so this audit (and those derived from it) would never run. - # Uncommenting the next lines fixes the name test, but that then - # exposes pre-existing bugs that need to be fixed. - # if self._cm.name == "crm-corosync": - # return True - return False + return self._cm.name == "crm-corosync" class CIBAudit(ClusterAudit): @@ -793,85 +788,72 @@ def __call__(self): for partition in ccm_partitions: self.debug(f"\tAuditing CIB consistency for: {partition}") - if self._audit_cib_contents(partition) == 0: + if not self._audit_cib_contents(partition): passed = False return passed + def _cleanup_cibs(self): + """Remove any fetched CIB files.""" + for f in glob.glob("/tmp/ctsaudit.*.xml"): + os.remove(f) + def _audit_cib_contents(self, hostlist): """Perform the CIB audit on the given hosts.""" passed = True - node0 = None - node0_xml = None partition_hosts = hostlist.split() for node in partition_hosts: - node_xml = self._store_remote_cib(node, node0) + node_xml = self._get_remote_cib(node) if node_xml is None: + # If we failed to fetch the CIB from a single node, the audit + # will fail. Clean up anything we did fetch and return. logging.log(f"Could not perform audit: No configuration from {node}") passed = False + self._cleanup_cibs() + return passed - elif node0 is None: - node0 = node - node0_xml = node_xml + with open(f"/tmp/ctsaudit.{node}.xml", "w", encoding="utf-8") as f: + for line in node_xml: + line = re.sub(r'cib-last-written="[^"]+"', 'cib-last-written=""', line) + f.write(line) - elif node0_xml is None: - logging.log(f"Could not perform audit: No configuration from {node0}") - passed = False + (first, rest) = (partition_hosts[0], partition_hosts[1:]) + first_xml = f"/tmp/ctsaudit.{first}.xml" - else: - (rc, result) = self._cm.rsh.call( - node0, f"crm_diff -VV -cf --new {node_xml} --original {node0_xml}", verbose=1) + for node in rest: + node_xml = f"/tmp/ctsaudit.{node}.xml" + proc = subprocess.run(["crm_diff", "-VV", "-c", "--new", node_xml, + "--original", first_xml], + check=False, capture_output=True, universal_newlines=True) - if rc != 0: - logging.log(f"Diff between {node0_xml} and {node_xml} failed: {rc}") - passed = False + if proc.returncode != 0: + logging.log(f"Diff between {first_xml} and {node_xml} failed: {proc.returncode}") + passed = False - for line in result: - if not re.search("", line): - passed = False - self.debug(f"CibDiff[{node0}-{node}]: {line}") - else: - self.debug(f"CibDiff[{node0}-{node}] Ignoring: {line}") + for line in proc.stdout.splitlines(): + if "" in line: + self.debug(f"CibDiff[{first}-{node}] Ignoring: {line}") + else: + passed = False + self.debug(f"CibDiff[{first}-{node}]: {line}") + self._cleanup_cibs() return passed - def _store_remote_cib(self, node, target): - """ - Store a copy of the given node's CIB on the given target node. - - If no target is given, store the CIB on the given node. - """ - filename = f"/tmp/ctsaudit.{node}.xml" - - if not target: - target = node - + def _get_remote_cib(self, node): + """Fetch a copy of the given node's CIB and return it as a list.""" (rc, lines) = self._cm.rsh.call(node, self._cm.templates["CibQuery"], verbose=1) if rc != 0: logging.log("Could not retrieve configuration") return None - self._cm.rsh.call("localhost", f"rm -f {filename}") - for line in lines: - self._cm.rsh.call("localhost", f"echo \'{line[:-1]}\' >> {filename}", verbose=0) - - if self._cm.rsh.copy(filename, f"root@{target}:{filename}") != 0: - logging.log("Could not store configuration") - return None - - return filename + return lines def is_applicable(self): """Return True if this audit is applicable in the current test configuration.""" - # @TODO Due to long-ago refactoring, this name test would never match, - # so this audit (and those derived from it) would never run. - # Uncommenting the next lines fixes the name test, but that then - # exposes pre-existing bugs that need to be fixed. - # if self._cm.name == "crm-corosync": - # return True - return False + return self._cm.name == "crm-corosync" class PartitionAudit(ClusterAudit): @@ -883,7 +865,7 @@ class PartitionAudit(ClusterAudit): * The number of partitions and the nodes in each is as expected * Each node is active when it should be active and inactive when it should be inactive - * The status and epoch of each node is as expected + * The status of each node is as expected * A partition has quorum * A partition has a DC when expected """ @@ -898,7 +880,6 @@ def __init__(self, cm): ClusterAudit.__init__(self, cm) self.name = "PartitionAudit" - self._node_epoch = {} self._node_state = {} self._node_quorum = {} @@ -919,13 +900,19 @@ def __call__(self): logging.log(f"\t {partition}") for partition in ccm_partitions: - if self._audit_partition(partition) == 0: + if not self._audit_partition(partition): passed = False + if not any(v == "1" for v in self._node_quorum.values()): + logging.log(f"ERROR: No node has quorum") + passed = False + return passed def _trim_string(self, avalue): """Remove the last character from a multi-character string.""" + avalue = avalue.strip() + if not avalue: return None @@ -934,20 +921,10 @@ def _trim_string(self, avalue): return avalue - def _trim2int(self, avalue): - """Remove the last character from a multi-character string and convert the result to an int.""" - trimmed = self._trim_string(avalue) - if trimmed: - return int(trimmed) - - return None - def _audit_partition(self, partition): """Perform the audit of a single partition.""" passed = True dc_found = [] - dc_allowed_list = [] - lowest_epoch = None node_list = partition.split() self.debug(f"Auditing partition: {partition}") @@ -959,30 +936,20 @@ def _audit_partition(self, partition): # checking for in this audit) (_, out) = self._cm.rsh.call(node, self._cm.templates["StatusCmd"] % node, verbose=1) - self._node_state[node] = out[0].strip() + if not out: + logging.log(f"ERROR: Could not determine status for node {node}") + passed = False + return passed - (_, out) = self._cm.rsh.call(node, self._cm.templates["EpochCmd"], verbose=1) - self._node_epoch[node] = out[0].strip() + self._node_state[node] = self._trim_string(out[0]) (_, out) = self._cm.rsh.call(node, self._cm.templates["QuorumCmd"], verbose=1) - self._node_quorum[node] = out[0].strip() - - self.debug(f"Node {node}: {self._node_state[node]} - {self._node_epoch[node]} - {self._node_quorum[node]}.") - self._node_state[node] = self._trim_string(self._node_state[node]) - self._node_epoch[node] = self._trim2int(self._node_epoch[node]) - self._node_quorum[node] = self._trim_string(self._node_quorum[node]) + if not out: + logging.log(f"ERROR: Could not determine quorum on node {node}") + passed = False + return passed - if not self._node_epoch[node]: - logging.log(f"Warn: Node {node} disappeared: can't determine epoch") - self._cm.expected_status[node] = "down" - # not in itself a reason to fail the audit (not what we're - # checking for in this audit) - elif lowest_epoch is None or self._node_epoch[node] < lowest_epoch: - lowest_epoch = self._node_epoch[node] - - if not lowest_epoch: - logging.log(f"Lowest epoch not determined in {partition}") - passed = False + self._node_quorum[node] = self._trim_string(out[0]) for node in node_list: if self._cm.expected_status[node] != "up": @@ -990,41 +957,19 @@ def _audit_partition(self, partition): if self._cm.is_node_dc(node, self._node_state[node]): dc_found.append(node) - if self._node_epoch[node] == lowest_epoch: - self.debug(f"{node}: OK") - elif not self._node_epoch[node]: - self.debug(f"Check on {node} ignored: no node epoch") - elif not lowest_epoch: - self.debug(f"Check on {node} ignored: no lowest epoch") - else: - logging.log(f"DC {node} is not the oldest node " - f"({self._node_epoch[node]} vs. {lowest_epoch})") - passed = False if not dc_found: - logging.log(f"DC not found on any of the {len(dc_allowed_list)} allowed " - f"nodes: {dc_allowed_list} (of {node_list})") + logging.log("DC not found on any node") elif len(dc_found) > 1: logging.log(f"{len(dc_found)} DCs ({dc_found}) found in cluster partition: {node_list}") passed = False - if not passed: - for node in node_list: - if self._cm.expected_status[node] == "up": - logging.log(f"epoch {self._node_epoch[node]} : {self._node_state[node]}") - return passed def is_applicable(self): """Return True if this audit is applicable in the current test configuration.""" - # @TODO Due to long-ago refactoring, this name test would never match, - # so this audit (and those derived from it) would never run. - # Uncommenting the next lines fixes the name test, but that then - # exposes pre-existing bugs that need to be fixed. - # if self._cm.name == "crm-corosync": - # return True - return False + return self._cm.name == "crm-corosync" # pylint: disable=invalid-name diff --git a/python/pacemaker/_cts/cibxml.py b/python/pacemaker/_cts/cibxml.py index b0b1a4a1481..7c5711dc822 100644 --- a/python/pacemaker/_cts/cibxml.py +++ b/python/pacemaker/_cts/cibxml.py @@ -546,7 +546,7 @@ def _constraints(self): for (k, kargs) in self._coloc.items(): attrs = {"id": f"{self.name}-with-{k}", "rsc": self.name, "with-rsc": k} - text += element("rsc_colocation", **attrs) + text += element("rsc_colocation", **attrs, **kargs) text += "" return text diff --git a/python/pacemaker/_cts/patterns.py b/python/pacemaker/_cts/patterns.py index 963f460f739..f737ef10a30 100644 --- a/python/pacemaker/_cts/patterns.py +++ b/python/pacemaker/_cts/patterns.py @@ -146,7 +146,6 @@ def __init__(self): "StartCmd": "service corosync start && service pacemaker start", "StopCmd": "service pacemaker stop; [ ! -e /usr/sbin/pacemaker-remoted ] || service pacemaker_remote stop; service corosync stop", - "EpochCmd": "crm_node -e", "QuorumCmd": "crm_node -q", "PartitionCmd": "crm_node -p", }) diff --git a/python/pacemaker/_cts/tests/splitbraintest.py b/python/pacemaker/_cts/tests/splitbraintest.py index d0d9fafa8cb..90709bb6306 100644 --- a/python/pacemaker/_cts/tests/splitbraintest.py +++ b/python/pacemaker/_cts/tests/splitbraintest.py @@ -185,6 +185,22 @@ def __call__(self, node): return self.failure("See previous errors") + def audit(self): + """Perform all the relevant audits (see ClusterAudit), returning whether or not they all passed.""" + passed = True + + for audit in self.audits: + # These audits don't work well on the split brain test + if audit.name in ["GroupAudit", "PrimitiveAudit"]: + continue + + if not audit(): + logging.log(f"Internal {self.name} Audit {audit.name} FAILED.") + self.incr("auditfail") + passed = False + + return passed + @property def errors_to_ignore(self): """Return a list of errors which should be ignored."""