From 668269b982a3e6878cf9d703657508a797d2c30d Mon Sep 17 00:00:00 2001 From: Anson Yim Date: Thu, 24 Sep 2026 16:59:28 -0400 Subject: [PATCH 1/7] concurrent oob_redfish plugin --- nodescraper/base/oobanddataplugin.py | 216 +++++++++++++++++- nodescraper/connection/redfish/__init__.py | 7 +- .../connection/redfish/redfish_manager.py | 55 ++++- .../connection/redfish/redfish_params.py | 51 ++++- 4 files changed, 317 insertions(+), 12 deletions(-) diff --git a/nodescraper/base/oobanddataplugin.py b/nodescraper/base/oobanddataplugin.py index c88ffc91..d19a383e 100644 --- a/nodescraper/base/oobanddataplugin.py +++ b/nodescraper/base/oobanddataplugin.py @@ -23,14 +23,17 @@ # SOFTWARE. # ############################################################################### -from typing import Generic +from concurrent.futures import ThreadPoolExecutor, as_completed +from typing import Any, Generic, Optional, Union from nodescraper.connection.redfish import ( RedfishConnectionManager, RedfishConnectionParams, ) +from nodescraper.enums import EventPriority, ExecutionStatus, SystemInteractionLevel from nodescraper.generictypes import TAnalyzeArg, TCollectArg, TDataModel from nodescraper.interfaces import DataPlugin +from nodescraper.models import TaskResult class OOBandDataPlugin( @@ -43,6 +46,215 @@ class OOBandDataPlugin( ], Generic[TDataModel, TCollectArg, TAnalyzeArg], ): - """Base class for out-of-band (OOB) plugins that use Redfish connection.""" + """Base class for OOB plugins (Redfish). Adds multi-target support. + + When ``connection_args.targets`` is non-empty, runs all configured collectors + concurrently across targets (one thread per target), keying results by + ``target_key``. Single-target behaviour and all IB plugins are unaffected. + """ CONNECTION_TYPE = RedfishConnectionManager + + def _collect_for_target( + self, + target_key: str, + conn: Any, + collector_classes: tuple, + collection_args: Any, + system_interaction_level: Any, + max_event_priority_level: Any, + ) -> tuple[str, list[TaskResult], Any]: + """Run all configured collectors for one target. Executed in a worker thread.""" + self.logger.info("Starting collection for target %r", target_key) + target_data = None + results: list[TaskResult] = [] + parent = f"{self.__class__.__name__}[{target_key}]" + for collector_cls in collector_classes: + resolved_args = self._resolve_collector_args(collector_cls, collection_args) + task = collector_cls( + system_info=self.system_info.model_copy(), + connection=conn, + logger=self.logger, + system_interaction_level=system_interaction_level, + max_event_priority_level=max_event_priority_level, + parent=parent, + task_result_hooks=self.task_result_hooks, + event_reporter=self.event_reporter, + session_id=self.session_id, + log_path=self.log_path, + ) + result, data = task.collect_data(resolved_args) + results.append(result) + target_data = self._merge_collected_data(target_data, data) + self.logger.info("Finished collection for target %r", target_key) + return target_key, results, target_data + + def analyze( + self, + max_event_priority_level: Optional[Union[EventPriority, str]] = EventPriority.CRITICAL, + analysis_args: Optional[Union[TAnalyzeArg, dict]] = None, + data: Optional[Any] = None, + ) -> TaskResult: + """Analyze collected data. Single-target delegates to DataPlugin.analyze(). + Multi-target runs the analyzer once per target and aggregates results.""" + multi_target_data: dict = getattr(self, "multi_target_data", None) or {} + + if not multi_target_data: + return super().analyze( + max_event_priority_level=max_event_priority_level, + analysis_args=analysis_args, + data=data, + ) + + if self.ANALYZER is None: + self.analysis_result = TaskResult( + status=ExecutionStatus.NOT_RAN, + parent=self.__class__.__name__, + message=f"Data analysis not supported for {self.__class__.__name__}", + ) + return self.analysis_result + + if ( + analysis_args is not None + and isinstance(analysis_args, dict) + and hasattr(self, "ANALYZER_ARGS") + and self.ANALYZER_ARGS is not None + ): + analysis_args = self.ANALYZER_ARGS.model_validate(analysis_args) # type: ignore[assignment] + + analysis_results: list[TaskResult] = [] + for target_key, target_data in multi_target_data.items(): + parent = f"{self.__class__.__name__}[{target_key}]" + analyzer_task = self.ANALYZER( + system_info=self.system_info.model_copy(), + logger=self.logger, + max_event_priority_level=max_event_priority_level or EventPriority.CRITICAL, + parent=parent, + task_result_hooks=self.task_result_hooks, + event_reporter=self.event_reporter, + session_id=self.session_id, + ) + analysis_results.append(analyzer_task.analyze_data(target_data, analysis_args)) + + self.analysis_result = self._aggregate_collection_results( + self.__class__.__name__, analysis_results + ) + return self.analysis_result + + def collect( + self, + max_event_priority_level: Optional[Union[EventPriority, str]] = EventPriority.CRITICAL, + system_interaction_level: Optional[ + Union[SystemInteractionLevel, str] + ] = SystemInteractionLevel.INTERACTIVE, + preserve_connection: bool = False, + collection_args: Optional[TCollectArg] = None, + ) -> TaskResult: + """Run collectors. Single-target delegates to DataPlugin.collect(). + Multi-target runs all targets concurrently and merges results per target_key.""" + collector_classes = self.get_collector_classes() + + # Ensure the connection manager exists so we can inspect its params before connecting. + if not self.connection_manager: + if self.CONNECTION_TYPE is None: + self.collection_result = TaskResult( + parent=self.__class__.__name__, + status=ExecutionStatus.NOT_RAN, + message=f"No connection type configured for {self.__class__.__name__}", + ) + return self.collection_result + self.connection_manager = self.CONNECTION_TYPE( + system_info=self.system_info.model_copy(), + logger=self.logger, + parent=self.__class__.__name__, + task_result_hooks=self.task_result_hooks, + event_reporter=self.event_reporter, + session_id=self.session_id, + ) + + cm = self.connection_manager + params = cm.connection_args + + # Single-target: delegate the full lifecycle to the parent. + if not (isinstance(params, RedfishConnectionParams) and params.is_multi_target): + return super().collect( + max_event_priority_level=max_event_priority_level, + system_interaction_level=system_interaction_level, + preserve_connection=preserve_connection, + collection_args=collection_args, + ) + + # Multi-target path. + if not collector_classes: + self.collection_result = TaskResult( + parent=self.__class__.__name__, + status=ExecutionStatus.NOT_RAN, + message=f"Data collection not supported for {self.__class__.__name__}", + ) + return self.collection_result + + try: + if cm.result.status == ExecutionStatus.UNSET: + cm.connect() + + target_connections: dict = getattr(cm, "target_connections", None) or {} + + if not target_connections: + self.collection_result = TaskResult( + parent=self.__class__.__name__, + status=ExecutionStatus.EXECUTION_FAILURE, + message="No Redfish target connections were established", + ) + return self.collection_result + + max_w = params.max_workers or min(len(target_connections), 32) + all_results: list[TaskResult] = [] + merged_data: dict[str, Any] = {} + + with ThreadPoolExecutor(max_workers=max_w) as executor: + futures = { + executor.submit( + self._collect_for_target, + target_key, + conn, + collector_classes, + collection_args, + system_interaction_level, + max_event_priority_level, + ): target_key + for target_key, conn in target_connections.items() + } + for future in as_completed(futures): + target_key = futures[future] + try: + key, results, data = future.result() + all_results.extend(results) + if data is not None: + merged_data[key] = data + except Exception as exc: + self.logger.error("Collection failed for target %r: %s", target_key, exc) + + self.collection_result = self._aggregate_collection_results( + self.__class__.__name__, all_results + ) + # Multi-target results are keyed by target — not a single DataModel. + # Store them separately so analysis gracefully returns NOT_RAN (rather + # than crashing when it receives a dict instead of a DataModel). + self.multi_target_data: dict[str, Any] = merged_data + self._data = None + + except Exception as e: + self.logger.exception( + "Unhandled exception in multi-target collection for %s", + self.__class__.__name__, + ) + self.collection_result = TaskResult( + parent=self.__class__.__name__, + status=ExecutionStatus.EXECUTION_FAILURE, + message=f"Unhandled exception running multi-target collection: {e}", + ) + finally: + if not preserve_connection: + cm.disconnect() + + return self.collection_result diff --git a/nodescraper/connection/redfish/__init__.py b/nodescraper/connection/redfish/__init__.py index 12b5af16..20302470 100644 --- a/nodescraper/connection/redfish/__init__.py +++ b/nodescraper/connection/redfish/__init__.py @@ -39,7 +39,11 @@ collect_oem_diagnostic_data, get_oem_diagnostic_allowable_values, ) -from .redfish_params import RedfishConnectionParams, redfish_params_to_ssh +from .redfish_params import ( + RedfishConnectionParams, + RedfishTargetParams, + redfish_params_to_ssh, +) from .redfish_path import RedfishPath __all__ = [ @@ -48,6 +52,7 @@ "RedfishGetResult", "RedfishConnectionManager", "RedfishConnectionParams", + "RedfishTargetParams", "redfish_params_to_ssh", "RedfishPath", "collect_oem_diagnostic_data", diff --git a/nodescraper/connection/redfish/redfish_manager.py b/nodescraper/connection/redfish/redfish_manager.py index cd2a9b00..7c44c71e 100644 --- a/nodescraper/connection/redfish/redfish_manager.py +++ b/nodescraper/connection/redfish/redfish_manager.py @@ -83,7 +83,7 @@ def __init__( ) def connect(self) -> TaskResult: - """Connect to the Redfish service and perform a simple GET to verify.""" + """Connect to Redfish (single-target or multi-target).""" if not self.connection_args: self._log_event( category=EventCategory.RUNTIME, @@ -94,7 +94,6 @@ def connect(self) -> TaskResult: self.result.status = ExecutionStatus.EXECUTION_FAILURE return self.result - # Accept dict from JSON config; convert to RedfishConnectionParams raw = self.connection_args if isinstance(raw, dict): params = RedfishConnectionParams.model_validate(raw) @@ -110,14 +109,57 @@ def connect(self) -> TaskResult: self.result.status = ExecutionStatus.EXECUTION_FAILURE return self.result + if params.is_multi_target: + self.target_connections: dict[str, RedfishConnection] = {} + for target in params.targets: # type: ignore[union-attr] + key, conn = self._connect_target(target) + if conn is not None: + self.target_connections[key] = conn + if not self.target_connections: + self.result.status = ExecutionStatus.EXECUTION_FAILURE + return self.result + + return self._connect_single(params) + + def _connect_target( + self, target: RedfishConnectionParams + ) -> tuple[str, Optional[RedfishConnection]]: + """Connect one target; returns (key, connection) or (key, None) on failure.""" + key = target.target_key or str(target.host) + password = target.password.get_secret_value() if target.password else None + base_url = _build_base_url(str(target.host), target.port, target.use_https) + try: + self.logger.info("Connecting to Redfish at %s (target=%r)", base_url, key) + conn = RedfishConnection( + base_url=base_url, + username=target.username or "", + password=password, + timeout=target.timeout_seconds, + use_session_auth=target.use_session_auth, + verify_ssl=target.verify_ssl, + api_root=target.api_root, + ) + conn._ensure_session() + conn.get_service_root() + return key, conn + except (RedfishConnectionError, Exception) as exc: # noqa: BLE001 + self._log_event( + category=EventCategory.RUNTIME, + description=f"Redfish connection failed for target {key!r}: {exc}", + priority=EventPriority.CRITICAL, + console_log=True, + ) + return key, None + + def _connect_single(self, params: RedfishConnectionParams) -> TaskResult: + """Connect in single-target mode (original connect logic).""" password = params.password.get_secret_value() if params.password else None base_url = _build_base_url(str(params.host), params.port, params.use_https) - try: self.logger.info("Connecting to Redfish at %s", base_url) self.connection = RedfishConnection( base_url=base_url, - username=params.username, + username=params.username or "", password=password, timeout=params.timeout_seconds, use_session_auth=params.use_session_auth, @@ -147,7 +189,10 @@ def connect(self) -> TaskResult: return self.result def disconnect(self) -> None: - """Disconnect and release the Redfish session.""" + """Disconnect all Redfish sessions.""" if self.connection is not None: self.connection.close() + for conn in getattr(self, "target_connections", {}).values(): + conn.close() + self.target_connections = {} super().disconnect() diff --git a/nodescraper/connection/redfish/redfish_params.py b/nodescraper/connection/redfish/redfish_params.py index 4eb70a96..6df829e8 100644 --- a/nodescraper/connection/redfish/redfish_params.py +++ b/nodescraper/connection/redfish/redfish_params.py @@ -27,7 +27,10 @@ from typing import Optional, Union -from pydantic import BaseModel, ConfigDict, Field, SecretStr +from pydantic import BaseModel, ConfigDict, Field, SecretStr, model_validator + +# RedfishTargetParams is an alias kept for callers that import it by name. +# Internally, both the top-level config and per-target entries use the same model. from pydantic.networks import IPvAnyAddress from nodescraper.connection.inband.sshparams import SSHConnectionParams @@ -36,12 +39,33 @@ class RedfishConnectionParams(BaseModel): - """Connection parameters for a Redfish (BMC) API endpoint.""" + """Connection parameters for a Redfish (BMC) API endpoint. + + Single-target mode: supply ``host`` (and optionally ``username``, ``password``, etc.). + Multi-target mode: supply ``targets`` — a list of ``RedfishConnectionParams`` entries, + each with ``host`` set and an optional ``target_key`` identifier. + """ model_config = ConfigDict(arbitrary_types_allowed=True) - host: Union[IPvAnyAddress, str] - username: str + target_key: Optional[str] = Field( + default=None, + description="Identifier used when this entry appears inside a 'targets' list.", + ) + name: Optional[str] = Field(default=None) + targets: Optional[list[RedfishConnectionParams]] = Field( + default=None, description="List of OOB Redfish targets (multi-target mode)." + ) + max_workers: Optional[int] = Field( + default=None, + ge=1, + description=( + "Max concurrent threads for multi-target collection. " + "Defaults to min(len(targets), 32) when not set." + ), + ) + host: Optional[Union[IPvAnyAddress, str]] = None + username: Optional[str] = None password: Optional[SecretStr] = None port: Optional[int] = Field(default=None, ge=1, le=65535) use_https: bool = True @@ -56,6 +80,25 @@ class RedfishConnectionParams(BaseModel): description="Redfish API path (e.g. 'redfish/v1'). Override for a different API version.", ) + @model_validator(mode="after") + def _validate_target_config(self) -> "RedfishConnectionParams": + if not self.targets and self.host is None: + raise ValueError( + "Either 'targets' (multi-target mode) or 'host' (single-target mode) must be provided." + ) + return self + + @property + def is_multi_target(self) -> bool: + """True when one or more targets are configured via the ``targets`` list.""" + return bool(self.targets) + + +RedfishConnectionParams.model_rebuild() + +# Backward-compatible alias so existing callers of RedfishTargetParams still work. +RedfishTargetParams = RedfishConnectionParams + def redfish_params_to_ssh( params: Union[RedfishConnectionParams, dict], From 0e752043e9982676f51ce0c39ccb56061b2a02eb Mon Sep 17 00:00:00 2001 From: Anson Yim Date: Mon, 28 Sep 2026 15:02:24 -0400 Subject: [PATCH 2/7] Refactor --- nodescraper/base/oobanddataplugin.py | 164 +----------------- nodescraper/base/redfishcollectortask.py | 141 ++++++++++++++- nodescraper/connection/redfish/__init__.py | 3 +- .../connection/redfish/redfish_manager.py | 27 ++- 4 files changed, 169 insertions(+), 166 deletions(-) diff --git a/nodescraper/base/oobanddataplugin.py b/nodescraper/base/oobanddataplugin.py index d19a383e..3d7aeb74 100644 --- a/nodescraper/base/oobanddataplugin.py +++ b/nodescraper/base/oobanddataplugin.py @@ -23,14 +23,13 @@ # SOFTWARE. # ############################################################################### -from concurrent.futures import ThreadPoolExecutor, as_completed from typing import Any, Generic, Optional, Union from nodescraper.connection.redfish import ( RedfishConnectionManager, RedfishConnectionParams, ) -from nodescraper.enums import EventPriority, ExecutionStatus, SystemInteractionLevel +from nodescraper.enums import EventPriority, ExecutionStatus from nodescraper.generictypes import TAnalyzeArg, TCollectArg, TDataModel from nodescraper.interfaces import DataPlugin from nodescraper.models import TaskResult @@ -48,47 +47,13 @@ class OOBandDataPlugin( ): """Base class for OOB plugins (Redfish). Adds multi-target support. - When ``connection_args.targets`` is non-empty, runs all configured collectors - concurrently across targets (one thread per target), keying results by - ``target_key``. Single-target behaviour and all IB plugins are unaffected. + Multi-target collection is handled transparently by RedfishDataCollector's + __init_subclass__ wrapper; this class only needs to override analyze() to run + the analyzer once per target after collection completes. """ CONNECTION_TYPE = RedfishConnectionManager - def _collect_for_target( - self, - target_key: str, - conn: Any, - collector_classes: tuple, - collection_args: Any, - system_interaction_level: Any, - max_event_priority_level: Any, - ) -> tuple[str, list[TaskResult], Any]: - """Run all configured collectors for one target. Executed in a worker thread.""" - self.logger.info("Starting collection for target %r", target_key) - target_data = None - results: list[TaskResult] = [] - parent = f"{self.__class__.__name__}[{target_key}]" - for collector_cls in collector_classes: - resolved_args = self._resolve_collector_args(collector_cls, collection_args) - task = collector_cls( - system_info=self.system_info.model_copy(), - connection=conn, - logger=self.logger, - system_interaction_level=system_interaction_level, - max_event_priority_level=max_event_priority_level, - parent=parent, - task_result_hooks=self.task_result_hooks, - event_reporter=self.event_reporter, - session_id=self.session_id, - log_path=self.log_path, - ) - result, data = task.collect_data(resolved_args) - results.append(result) - target_data = self._merge_collected_data(target_data, data) - self.logger.info("Finished collection for target %r", target_key) - return target_key, results, target_data - def analyze( self, max_event_priority_level: Optional[Union[EventPriority, str]] = EventPriority.CRITICAL, @@ -97,7 +62,8 @@ def analyze( ) -> TaskResult: """Analyze collected data. Single-target delegates to DataPlugin.analyze(). Multi-target runs the analyzer once per target and aggregates results.""" - multi_target_data: dict = getattr(self, "multi_target_data", None) or {} + cm = self.connection_manager + multi_target_data: dict = getattr(cm, "_multi_target_data", None) or {} if not multi_target_data: return super().analyze( @@ -140,121 +106,3 @@ def analyze( self.__class__.__name__, analysis_results ) return self.analysis_result - - def collect( - self, - max_event_priority_level: Optional[Union[EventPriority, str]] = EventPriority.CRITICAL, - system_interaction_level: Optional[ - Union[SystemInteractionLevel, str] - ] = SystemInteractionLevel.INTERACTIVE, - preserve_connection: bool = False, - collection_args: Optional[TCollectArg] = None, - ) -> TaskResult: - """Run collectors. Single-target delegates to DataPlugin.collect(). - Multi-target runs all targets concurrently and merges results per target_key.""" - collector_classes = self.get_collector_classes() - - # Ensure the connection manager exists so we can inspect its params before connecting. - if not self.connection_manager: - if self.CONNECTION_TYPE is None: - self.collection_result = TaskResult( - parent=self.__class__.__name__, - status=ExecutionStatus.NOT_RAN, - message=f"No connection type configured for {self.__class__.__name__}", - ) - return self.collection_result - self.connection_manager = self.CONNECTION_TYPE( - system_info=self.system_info.model_copy(), - logger=self.logger, - parent=self.__class__.__name__, - task_result_hooks=self.task_result_hooks, - event_reporter=self.event_reporter, - session_id=self.session_id, - ) - - cm = self.connection_manager - params = cm.connection_args - - # Single-target: delegate the full lifecycle to the parent. - if not (isinstance(params, RedfishConnectionParams) and params.is_multi_target): - return super().collect( - max_event_priority_level=max_event_priority_level, - system_interaction_level=system_interaction_level, - preserve_connection=preserve_connection, - collection_args=collection_args, - ) - - # Multi-target path. - if not collector_classes: - self.collection_result = TaskResult( - parent=self.__class__.__name__, - status=ExecutionStatus.NOT_RAN, - message=f"Data collection not supported for {self.__class__.__name__}", - ) - return self.collection_result - - try: - if cm.result.status == ExecutionStatus.UNSET: - cm.connect() - - target_connections: dict = getattr(cm, "target_connections", None) or {} - - if not target_connections: - self.collection_result = TaskResult( - parent=self.__class__.__name__, - status=ExecutionStatus.EXECUTION_FAILURE, - message="No Redfish target connections were established", - ) - return self.collection_result - - max_w = params.max_workers or min(len(target_connections), 32) - all_results: list[TaskResult] = [] - merged_data: dict[str, Any] = {} - - with ThreadPoolExecutor(max_workers=max_w) as executor: - futures = { - executor.submit( - self._collect_for_target, - target_key, - conn, - collector_classes, - collection_args, - system_interaction_level, - max_event_priority_level, - ): target_key - for target_key, conn in target_connections.items() - } - for future in as_completed(futures): - target_key = futures[future] - try: - key, results, data = future.result() - all_results.extend(results) - if data is not None: - merged_data[key] = data - except Exception as exc: - self.logger.error("Collection failed for target %r: %s", target_key, exc) - - self.collection_result = self._aggregate_collection_results( - self.__class__.__name__, all_results - ) - # Multi-target results are keyed by target — not a single DataModel. - # Store them separately so analysis gracefully returns NOT_RAN (rather - # than crashing when it receives a dict instead of a DataModel). - self.multi_target_data: dict[str, Any] = merged_data - self._data = None - - except Exception as e: - self.logger.exception( - "Unhandled exception in multi-target collection for %s", - self.__class__.__name__, - ) - self.collection_result = TaskResult( - parent=self.__class__.__name__, - status=ExecutionStatus.EXECUTION_FAILURE, - message=f"Unhandled exception running multi-target collection: {e}", - ) - finally: - if not preserve_connection: - cm.disconnect() - - return self.collection_result diff --git a/nodescraper/base/redfishcollectortask.py b/nodescraper/base/redfishcollectortask.py index 48c11ae2..44a654d4 100644 --- a/nodescraper/base/redfishcollectortask.py +++ b/nodescraper/base/redfishcollectortask.py @@ -24,14 +24,20 @@ # ############################################################################### import logging -from typing import Generic, Optional, Union +from concurrent.futures import ThreadPoolExecutor, as_completed +from functools import wraps +from typing import Any, Callable, Generic, Optional, Union -from nodescraper.connection.redfish import RedfishConnection, RedfishGetResult +from nodescraper.connection.redfish import ( + MultiTargetRedfishConnection, + RedfishConnection, + RedfishGetResult, +) from nodescraper.constants import DEFAULT_EVENT_REPORTER -from nodescraper.enums import EventPriority +from nodescraper.enums import EventPriority, ExecutionStatus from nodescraper.generictypes import TCollectArg, TDataModel from nodescraper.interfaces import DataCollector, TaskResultHook -from nodescraper.models import SystemInfo +from nodescraper.models import SystemInfo, TaskResult class RedfishDataCollector( @@ -76,6 +82,133 @@ def __init__( **kwargs, ) + def __init_subclass__(cls, **kwargs: Any) -> None: + """Wrap each concrete collector's collect_data with multi-target detection. + + Execution order when collect_data is called: + 1. _multi_target_wrapper (outermost) — checks for MultiTargetRedfishConnection + 2. collect_decorator (from DataCollector) — handles hooks, finalization + 3. Concrete collect_data implementation (innermost) + + In multi-target mode a fresh collector instance is created per target so that + concurrent threads never share mutable state. + """ + super().__init_subclass__(**kwargs) # applies collect_decorator via DataCollector + if "collect_data" not in vars(cls): + return + inner = cls.collect_data # already wrapped by collect_decorator at this point + + @wraps(inner) + def _multi_target_wrapper( + collector: "RedfishDataCollector", + args: Any = None, + *, + _fn: Any = inner, + ) -> tuple[TaskResult, Any]: + if not isinstance(collector.connection, MultiTargetRedfishConnection): + return _fn(collector, args) + + multi_conn: MultiTargetRedfishConnection = collector.connection # type: ignore[assignment] + parent_name = type(collector).__name__ + all_results: list[TaskResult] = [] + max_workers = min(len(multi_conn.target_connections), 32) + + def _run_for_target( + target_key: str, conn: RedfishConnection + ) -> tuple[str, TaskResult, Any]: + # Each thread gets its own collector instance to avoid shared state. + target_collector = type(collector)( + system_info=collector.system_info, + connection=conn, + logger=collector.logger, + max_event_priority_level=getattr( + collector, "max_event_priority_level", EventPriority.CRITICAL + ), + parent=f"{parent_name}[{target_key}]", + task_result_hooks=collector.task_result_hooks, + event_reporter=collector.event_reporter, + session_id=collector.session_id, + log_path=getattr(collector, "log_path", None), + system_interaction_level=getattr(collector, "system_interaction_level", None), + ) + collector.logger.info( + "Starting collection for target %r using %s", target_key, parent_name + ) + result, data = _fn(target_collector, args) + collector.logger.info("Finished collection for target %r", target_key) + return target_key, result, data + + with ThreadPoolExecutor(max_workers=max_workers) as executor: + futures = { + executor.submit(_run_for_target, tk, conn): tk + for tk, conn in multi_conn.target_connections.items() + } + for future in as_completed(futures): + target_key = futures[future] + try: + tk, result, data = future.result() + all_results.append(result) + if data is not None: + multi_conn.multi_target_data[tk] = data + except Exception as exc: + collector.logger.error( + "Collection failed for target %r: %s", target_key, exc + ) + + if not all_results: + return TaskResult(status=ExecutionStatus.NOT_RAN, parent=parent_name), None + combined_status = max(r.status for r in all_results) + return TaskResult(status=combined_status, parent=parent_name), None + + cls.collect_data = _multi_target_wrapper # type: ignore[method-assign, assignment] + + @staticmethod + def collect_for_target( + target_key: str, + conn: RedfishConnection, + collector_classes: tuple, + collection_args: Any, + *, + system_info: SystemInfo, + logger: logging.Logger, + system_interaction_level: Any, + max_event_priority_level: Any, + parent: str, + task_result_hooks: list, + event_reporter: str, + session_id: Optional[str], + log_path: Optional[str], + resolve_args: Callable, + merge_data: Callable, + ) -> tuple[str, list[TaskResult], Any]: + """Run all configured collectors for one Redfish target. + + Intended to be submitted to a ThreadPoolExecutor by OOBandDataPlugin. + Returns (target_key, per-collector TaskResults, merged data model). + """ + logger.info("Starting collection for target %r", target_key) + target_data = None + results: list[TaskResult] = [] + for collector_cls in collector_classes: + resolved = resolve_args(collector_cls, collection_args) + task = collector_cls( + system_info=system_info.model_copy(), + connection=conn, + logger=logger, + system_interaction_level=system_interaction_level, + max_event_priority_level=max_event_priority_level, + parent=parent, + task_result_hooks=task_result_hooks, + event_reporter=event_reporter, + session_id=session_id, + log_path=log_path, + ) + result, data = task.collect_data(resolved) + results.append(result) + target_data = merge_data(target_data, data) + logger.info("Finished collection for target %r", target_key) + return target_key, results, target_data + def _run_redfish_get( self, path: str, diff --git a/nodescraper/connection/redfish/__init__.py b/nodescraper/connection/redfish/__init__.py index 20302470..61776d68 100644 --- a/nodescraper/connection/redfish/__init__.py +++ b/nodescraper/connection/redfish/__init__.py @@ -34,7 +34,7 @@ RF_MEMBERS_NEXT_LINK, RF_ODATA_ID, ) -from .redfish_manager import RedfishConnectionManager +from .redfish_manager import MultiTargetRedfishConnection, RedfishConnectionManager from .redfish_oem_diag import ( collect_oem_diagnostic_data, get_oem_diagnostic_allowable_values, @@ -50,6 +50,7 @@ "RedfishConnection", "RedfishConnectionError", "RedfishGetResult", + "MultiTargetRedfishConnection", "RedfishConnectionManager", "RedfishConnectionParams", "RedfishTargetParams", diff --git a/nodescraper/connection/redfish/redfish_manager.py b/nodescraper/connection/redfish/redfish_manager.py index 7c44c71e..5ad14f99 100644 --- a/nodescraper/connection/redfish/redfish_manager.py +++ b/nodescraper/connection/redfish/redfish_manager.py @@ -26,7 +26,7 @@ from __future__ import annotations from logging import Logger -from typing import Optional, Union +from typing import Any, Optional, Union from nodescraper.enums import EventCategory, EventPriority, ExecutionStatus from nodescraper.interfaces.connectionmanager import ConnectionManager @@ -45,6 +45,19 @@ def _build_base_url(host: str, port: Optional[int], use_https: bool) -> str: return f"{scheme}://{host_str}" +class MultiTargetRedfishConnection: + """Sentinel placed on RedfishConnectionManager.connection in multi-target mode. + + Carrying this object (rather than leaving connection=None) allows the standard + DataPlugin.collect() guard to pass, while RedfishDataCollector's __init_subclass__ + wrapper intercepts collect_data and iterates the individual target connections. + """ + + def __init__(self, target_connections: dict[str, RedfishConnection]) -> None: + self.target_connections = target_connections + self.multi_target_data: dict[str, Any] = {} + + class RedfishConnectionManager(ConnectionManager[RedfishConnection, RedfishConnectionParams]): """Connection manager for Redfish (BMC) API.""" @@ -117,6 +130,9 @@ def connect(self) -> TaskResult: self.target_connections[key] = conn if not self.target_connections: self.result.status = ExecutionStatus.EXECUTION_FAILURE + else: + # Set self.connection so DataPlugin.collect() does not short-circuit. + self.connection = MultiTargetRedfishConnection(self.target_connections) # type: ignore[assignment] return self.result return self._connect_single(params) @@ -189,8 +205,13 @@ def _connect_single(self, params: RedfishConnectionParams) -> TaskResult: return self.result def disconnect(self) -> None: - """Disconnect all Redfish sessions.""" - if self.connection is not None: + """Disconnect all Redfish sessions, preserving multi-target data on the manager.""" + if isinstance(self.connection, MultiTargetRedfishConnection): + # Persist collected data so analyze() can access it after disconnect. + self._multi_target_data: dict[str, Any] = dict(self.connection.multi_target_data) + for conn in self.connection.target_connections.values(): + conn.close() + elif self.connection is not None: self.connection.close() for conn in getattr(self, "target_connections", {}).values(): conn.close() From ad62d695bfc7418c25fcff1428d5cc909f14ceac Mon Sep 17 00:00:00 2001 From: Anson Yim Date: Mon, 28 Sep 2026 16:42:30 -0400 Subject: [PATCH 3/7] update readme --- README.md | 44 +++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 43 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 14fa516a..1c039e32 100644 --- a/README.md +++ b/README.md @@ -28,6 +28,9 @@ system debug. For details on what data is collected and analyzed, see the [plugi - [Global args](#global-args) - [Plugin config: **'--plugin-configs' command**](#plugin-config---plugin-configs-command) - [Post-action plugins](#post-action-plugins) + - [Config structure](#config-structure) + - [Condition fields](#condition-fields) + - [Example: run OsPlugin if DmesgPlugin finds error-level events](#example-run-osplugin-if-dmesgplugin-finds-error-level-events) - [Reference config: **'gen-reference-config' command**](#reference-config-gen-reference-config-command) ## Installation @@ -210,6 +213,43 @@ Redfish (BMC) connection for Redfish-only plugins: - `api_root` (optional): Redfish API path (e.g. `redfish/v1`). If omitted, the default `redfish/v1` is used. Override this when your BMC uses a different API version path. +Redfish **multi-target** connection (collect from multiple BMCs concurrently): + +```json +{ + "RedfishConnectionManager": { + "targets": [ + { + "target_key": "node-a", + "host": "bmc-node-a.example.com", + "username": "admin", + "password": "secret", + "use_https": true, + "verify_ssl": false, + "timeout_seconds": 30 + }, + { + "target_key": "node-b", + "host": "bmc-node-b.example.com", + "username": "admin", + "password": "secret", + "use_https": true, + "verify_ssl": false, + "timeout_seconds": 30 + } + ], + "max_workers": 32 + } +} +``` + +Multi-target mode is supported by all OOB Redfish plugins (`RedfishEndpointPlugin`, `RedfishOemDiagPlugin`, etc.). All targets are collected **concurrently** — wall-clock time is bounded by the slowest individual target, not the sum. + +- `targets`: list of per-target connection parameters. Each entry accepts the same fields as the single-target config plus an optional `target_key` (used as the result key; defaults to the host string). +- `max_workers` (optional): maximum concurrent collection threads. Defaults to `min(len(targets), 32)`. + +Per-target results are written to separate log subdirectories named `[]/`. + **Notes:** - If using SSH keys, specify `key_filename` instead of `password`. - The remote user must have permissions to run the requested plugins and access required files. If needed, use the `--skip-sudo` argument to skip plugins requiring sudo. @@ -472,9 +512,11 @@ Use a plugin config that points at your LogService and lists the types to collec The RedfishEndpointPlugin collects Redfish URIs (GET responses) and optionally runs checks on the returned JSON. It requires a Redfish connection config (same as RedfishOemDiagPlugin). +**Multi-target support:** `RedfishEndpointPlugin` supports concurrent collection from multiple BMCs simultaneously. Use the `targets` list format in your connection config (see [multi-target example above](#example-connection_configjson)) — the same `uris`/`checks` plugin config is applied to every target in parallel. Per-target results are written to `redfish_endpoint_plugin[]/` under the run log directory. + **How to run** -1. Create a connection config (e.g. `connection-config.json`) with `RedfishConnectionManager` and your BMC host, credentials, and API root. +1. Create a connection config (e.g. `connection-config.json`) with `RedfishConnectionManager` and your BMC host (single target) or `targets` list (multi-target). 2. Create a plugin config with `uris` to collect and optional `checks` for analysis (see example below). For example save as `plugin_config_redfish_endpoint.json`. 3. Run: ```sh From a9e190ed4c233a6259ef6c821157760167459707 Mon Sep 17 00:00:00 2001 From: Alexandra Bara Date: Tue, 29 Sep 2026 15:40:40 -0500 Subject: [PATCH 4/7] fix for multi target analysis / collection --- README.md | 20 +- config/connection-config_inband.example.json | 8 + config/connection-config_oob.example.json | 12 + ...n-config_redfish_multi_target.example.json | 25 ++ nodescraper/base/oobanddataplugin.py | 15 +- nodescraper/base/redfishcollectortask.py | 169 +++++++----- .../oob_ssh/oob_ssh_connection_manager.py | 8 + nodescraper/connection/redfish/__init__.py | 7 +- .../connection/redfish/redfish_manager.py | 76 ++++-- .../connection/redfish/redfish_params.py | 5 + .../helpers/plugin_execution_target.py | 10 + nodescraper/interfaces/dataplugin.py | 55 ++-- .../redfish/test_redfish_multi_target.py | 257 ++++++++++++++++++ 13 files changed, 542 insertions(+), 125 deletions(-) create mode 100644 config/connection-config_inband.example.json create mode 100644 config/connection-config_oob.example.json create mode 100644 config/connection-config_redfish_multi_target.example.json create mode 100644 test/unit/connection/redfish/test_redfish_multi_target.py diff --git a/README.md b/README.md index 1c039e32..4ff50d3c 100644 --- a/README.md +++ b/README.md @@ -195,6 +195,8 @@ In-band (SSH) connection: } ``` +A sample is in `config/connection-config_inband.example.json`. Use `password` or `key_filename`. + Redfish (BMC) connection for Redfish-only plugins: ```json @@ -213,7 +215,11 @@ Redfish (BMC) connection for Redfish-only plugins: - `api_root` (optional): Redfish API path (e.g. `redfish/v1`). If omitted, the default `redfish/v1` is used. Override this when your BMC uses a different API version path. -Redfish **multi-target** connection (collect from multiple BMCs concurrently): +OOB SSH plugins use this same single-host `RedfishConnectionManager` block and open SSH to that BMC. A sample is in `config/connection-config_oob.example.json`. + +#### Redfish multi-target + +Redfish plugins can collect from multiple BMCs concurrently. In-band plugins and OOB SSH plugins stay single-host. If this config has no top-level `host`, OOB SSH plugins are skipped. Add a top-level `host` when those plugins should still run against one BMC. ```json { @@ -243,12 +249,16 @@ Redfish **multi-target** connection (collect from multiple BMCs concurrently): } ``` -Multi-target mode is supported by all OOB Redfish plugins (`RedfishEndpointPlugin`, `RedfishOemDiagPlugin`, etc.). All targets are collected **concurrently** — wall-clock time is bounded by the slowest individual target, not the sum. +Multi-target mode applies to Redfish plugins (`RedfishEndpointPlugin`, `RedfishOemDiagPlugin`, and other plugins based on `OOBandDataPlugin`). Targets are collected concurrently. Wall-clock time follows the slowest target. + +A target that fails to connect or collect does not fail the run when another target succeeds. That plugin result is a warning, and analysis still runs for the targets that returned data. The run fails when every target fails, or when analysis of collected data reports an error. - `targets`: list of per-target connection parameters. Each entry accepts the same fields as the single-target config plus an optional `target_key` (used as the result key; defaults to the host string). -- `max_workers` (optional): maximum concurrent collection threads. Defaults to `min(len(targets), 32)`. +- `max_workers` (optional): maximum concurrent collection threads. Defaults to `min(len(targets), 32)` and is capped at 32. + +Per-target results are written to `[]//`. -Per-target results are written to separate log subdirectories named `[]/`. +A ready-to-edit sample is in `config/connection-config_redfish_multi_target.example.json`. Replace the example hosts and password, then pass it with `--connection-config`. **Notes:** - If using SSH keys, specify `key_filename` instead of `password`. @@ -512,7 +522,7 @@ Use a plugin config that points at your LogService and lists the types to collec The RedfishEndpointPlugin collects Redfish URIs (GET responses) and optionally runs checks on the returned JSON. It requires a Redfish connection config (same as RedfishOemDiagPlugin). -**Multi-target support:** `RedfishEndpointPlugin` supports concurrent collection from multiple BMCs simultaneously. Use the `targets` list format in your connection config (see [multi-target example above](#example-connection_configjson)) — the same `uris`/`checks` plugin config is applied to every target in parallel. Per-target results are written to `redfish_endpoint_plugin[]/` under the run log directory. +**Multi-target support:** `RedfishEndpointPlugin` collects from each BMC in the `targets` list at the same time. Use the [Redfish multi-target](#redfish-multi-target) connection config. The same `uris` and `checks` apply to every target. Per-target results are written to `redfish_endpoint_plugin[]/redfish_endpoint_collector/` under the run log directory. A BMC that cannot be reached is reported as a warning when another target succeeds. **How to run** diff --git a/config/connection-config_inband.example.json b/config/connection-config_inband.example.json new file mode 100644 index 00000000..cc05b2ff --- /dev/null +++ b/config/connection-config_inband.example.json @@ -0,0 +1,8 @@ +{ + "InBandConnectionManager": { + "hostname": "host.example.com", + "port": 22, + "username": "admin", + "password": "placeholder" + } +} diff --git a/config/connection-config_oob.example.json b/config/connection-config_oob.example.json new file mode 100644 index 00000000..1edd627b --- /dev/null +++ b/config/connection-config_oob.example.json @@ -0,0 +1,12 @@ +{ + "RedfishConnectionManager": { + "host": "bmc.example.com", + "port": 443, + "username": "admin", + "password": "placeholder", + "use_https": true, + "verify_ssl": false, + "timeout_seconds": 30, + "api_root": "redfish/v1" + } +} diff --git a/config/connection-config_redfish_multi_target.example.json b/config/connection-config_redfish_multi_target.example.json new file mode 100644 index 00000000..34cb845e --- /dev/null +++ b/config/connection-config_redfish_multi_target.example.json @@ -0,0 +1,25 @@ +{ + "RedfishConnectionManager": { + "max_workers": 2, + "targets": [ + { + "target_key": "node-a", + "host": "bmc-node-a.example.com", + "username": "admin", + "password": "placeholder", + "use_https": true, + "verify_ssl": false, + "timeout_seconds": 30 + }, + { + "target_key": "node-b", + "host": "bmc-node-b.example.com", + "username": "admin", + "password": "placeholder", + "use_https": true, + "verify_ssl": false, + "timeout_seconds": 30 + } + ] + } +} diff --git a/nodescraper/base/oobanddataplugin.py b/nodescraper/base/oobanddataplugin.py index 3d7aeb74..4d90858c 100644 --- a/nodescraper/base/oobanddataplugin.py +++ b/nodescraper/base/oobanddataplugin.py @@ -28,6 +28,7 @@ from nodescraper.connection.redfish import ( RedfishConnectionManager, RedfishConnectionParams, + collected_multi_target_data, ) from nodescraper.enums import EventPriority, ExecutionStatus from nodescraper.generictypes import TAnalyzeArg, TCollectArg, TDataModel @@ -60,10 +61,18 @@ def analyze( analysis_args: Optional[Union[TAnalyzeArg, dict]] = None, data: Optional[Any] = None, ) -> TaskResult: - """Analyze collected data. Single-target delegates to DataPlugin.analyze(). - Multi-target runs the analyzer once per target and aggregates results.""" + """Analyze collected data for one BMC or once per Redfish target. + + Args: + max_event_priority_level: Priority limit for events. + analysis_args: Analyzer arguments. + data: Pre-collected data for a single-target run. + + Returns: + TaskResult: Analysis result for the targets that returned data. + """ cm = self.connection_manager - multi_target_data: dict = getattr(cm, "_multi_target_data", None) or {} + multi_target_data: dict = collected_multi_target_data(cm) if not multi_target_data: return super().analyze( diff --git a/nodescraper/base/redfishcollectortask.py b/nodescraper/base/redfishcollectortask.py index 44a654d4..8b28b2c2 100644 --- a/nodescraper/base/redfishcollectortask.py +++ b/nodescraper/base/redfishcollectortask.py @@ -24,9 +24,11 @@ # ############################################################################### import logging +import re from concurrent.futures import ThreadPoolExecutor, as_completed from functools import wraps -from typing import Any, Callable, Generic, Optional, Union +from pathlib import Path +from typing import Any, Generic, Optional, Union from nodescraper.connection.redfish import ( MultiTargetRedfishConnection, @@ -37,8 +39,46 @@ from nodescraper.enums import EventPriority, ExecutionStatus from nodescraper.generictypes import TCollectArg, TDataModel from nodescraper.interfaces import DataCollector, TaskResultHook +from nodescraper.interfaces.dataplugin import DataPlugin from nodescraper.models import SystemInfo, TaskResult +_HARD_FAIL = {ExecutionStatus.ERROR, ExecutionStatus.EXECUTION_FAILURE} +_TARGET_SUCCESS = {ExecutionStatus.OK, ExecutionStatus.WARNING} + + +def _target_dir_name(target_key: str) -> str: + """Return a filesystem-safe directory name for a target key. + + Args: + target_key: Target identifier from the connection config. + + Returns: + str: Directory name containing only letters, numbers, dot, underscore, and dash. + """ + # Collapse characters that are unsafe in a log directory name. + cleaned = re.sub(r"[^A-Za-z0-9._-]+", "_", target_key).strip("._") + return cleaned or "target" + + +def combine_target_results(parent: str, results: list[TaskResult]) -> TaskResult: + """Merge per-target task results. + + Args: + parent: Parent name stored on the combined result. + results: Per-target task results. + + Returns: + TaskResult: Combined result. Mixed success and failure stays at warning. + """ + aggregated = DataPlugin._aggregate_collection_results(parent, results) + hard_fail = any(result.status in _HARD_FAIL for result in results) + succeeded = any(result.status in _TARGET_SUCCESS for result in results) + if hard_fail and succeeded: + aggregated.status = ExecutionStatus.WARNING + note = "One or more targets failed; continuing with the rest." + aggregated.message = f"{note} {aggregated.message}".strip() if aggregated.message else note + return aggregated + class RedfishDataCollector( DataCollector[RedfishConnection, TDataModel, TCollectArg], @@ -109,28 +149,42 @@ def _multi_target_wrapper( return _fn(collector, args) multi_conn: MultiTargetRedfishConnection = collector.connection # type: ignore[assignment] - parent_name = type(collector).__name__ - all_results: list[TaskResult] = [] - max_workers = min(len(multi_conn.target_connections), 32) + parent_name = collector.parent or type(collector).__name__ + all_results: list[TaskResult] = [ + TaskResult( + status=ExecutionStatus.WARNING, + parent=f"{parent_name}[{target_key}]", + message=message, + ) + for target_key, message in multi_conn.failed_targets.items() + ] + target_count = len(multi_conn.target_connections) + configured_workers = multi_conn.max_workers or target_count + max_workers = min(configured_workers, target_count, 32) if target_count else 1 def _run_for_target( target_key: str, conn: RedfishConnection ) -> tuple[str, TaskResult, Any]: - # Each thread gets its own collector instance to avoid shared state. + target_parent = f"{parent_name}[{target_key}]" + target_log = getattr(collector, "log_path", None) + if target_log: + target_log = str(Path(target_log) / _target_dir_name(target_key)) target_collector = type(collector)( - system_info=collector.system_info, + system_info=collector.system_info.model_copy(), connection=conn, logger=collector.logger, max_event_priority_level=getattr( collector, "max_event_priority_level", EventPriority.CRITICAL ), - parent=f"{parent_name}[{target_key}]", + parent=target_parent, task_result_hooks=collector.task_result_hooks, event_reporter=collector.event_reporter, session_id=collector.session_id, - log_path=getattr(collector, "log_path", None), + log_path=target_log, system_interaction_level=getattr(collector, "system_interaction_level", None), ) + if target_log: + target_collector.log_path = target_log collector.logger.info( "Starting collection for target %r using %s", target_key, parent_name ) @@ -138,77 +192,44 @@ def _run_for_target( collector.logger.info("Finished collection for target %r", target_key) return target_key, result, data - with ThreadPoolExecutor(max_workers=max_workers) as executor: - futures = { - executor.submit(_run_for_target, tk, conn): tk - for tk, conn in multi_conn.target_connections.items() - } - for future in as_completed(futures): - target_key = futures[future] - try: - tk, result, data = future.result() - all_results.append(result) - if data is not None: - multi_conn.multi_target_data[tk] = data - except Exception as exc: - collector.logger.error( - "Collection failed for target %r: %s", target_key, exc - ) + if multi_conn.target_connections: + with ThreadPoolExecutor(max_workers=max_workers) as executor: + futures = { + executor.submit(_run_for_target, tk, conn): tk + for tk, conn in multi_conn.target_connections.items() + } + for future in as_completed(futures): + target_key = futures[future] + try: + tk, result, data = future.result() + all_results.append(result) + if data is not None: + multi_conn.multi_target_data[tk] = data + except Exception as exc: + collector.logger.error( + "Collection failed for target %r: %s", target_key, exc + ) + all_results.append( + TaskResult( + status=ExecutionStatus.EXECUTION_FAILURE, + parent=f"{parent_name}[{target_key}]", + message=f"Collection failed for target {target_key!r}: {exc}", + ) + ) if not all_results: - return TaskResult(status=ExecutionStatus.NOT_RAN, parent=parent_name), None - combined_status = max(r.status for r in all_results) - return TaskResult(status=combined_status, parent=parent_name), None + return ( + TaskResult( + status=ExecutionStatus.NOT_RAN, + parent=parent_name, + message="No Redfish targets were collected", + ), + None, + ) + return combine_target_results(parent_name, all_results), None cls.collect_data = _multi_target_wrapper # type: ignore[method-assign, assignment] - @staticmethod - def collect_for_target( - target_key: str, - conn: RedfishConnection, - collector_classes: tuple, - collection_args: Any, - *, - system_info: SystemInfo, - logger: logging.Logger, - system_interaction_level: Any, - max_event_priority_level: Any, - parent: str, - task_result_hooks: list, - event_reporter: str, - session_id: Optional[str], - log_path: Optional[str], - resolve_args: Callable, - merge_data: Callable, - ) -> tuple[str, list[TaskResult], Any]: - """Run all configured collectors for one Redfish target. - - Intended to be submitted to a ThreadPoolExecutor by OOBandDataPlugin. - Returns (target_key, per-collector TaskResults, merged data model). - """ - logger.info("Starting collection for target %r", target_key) - target_data = None - results: list[TaskResult] = [] - for collector_cls in collector_classes: - resolved = resolve_args(collector_cls, collection_args) - task = collector_cls( - system_info=system_info.model_copy(), - connection=conn, - logger=logger, - system_interaction_level=system_interaction_level, - max_event_priority_level=max_event_priority_level, - parent=parent, - task_result_hooks=task_result_hooks, - event_reporter=event_reporter, - session_id=session_id, - log_path=log_path, - ) - result, data = task.collect_data(resolved) - results.append(result) - target_data = merge_data(target_data, data) - logger.info("Finished collection for target %r", target_key) - return target_key, results, target_data - def _run_redfish_get( self, path: str, diff --git a/nodescraper/connection/oob_ssh/oob_ssh_connection_manager.py b/nodescraper/connection/oob_ssh/oob_ssh_connection_manager.py index 823c00cb..1289ad85 100644 --- a/nodescraper/connection/oob_ssh/oob_ssh_connection_manager.py +++ b/nodescraper/connection/oob_ssh/oob_ssh_connection_manager.py @@ -88,6 +88,14 @@ def connect(self) -> TaskResult: self.result.status = ExecutionStatus.EXECUTION_FAILURE return self.result + if params.is_multi_target and params.host is None: + self.result.status = ExecutionStatus.NOT_RAN + self.result.message = ( + "OOB SSH plugins use a single BMC host. " + "A targets list applies only to Redfish plugins." + ) + return self.result + try: ssh_params = redfish_params_to_ssh(params) self.logger.info("Initializing OOB SSH to BMC host %s", ssh_params.hostname) diff --git a/nodescraper/connection/redfish/__init__.py b/nodescraper/connection/redfish/__init__.py index 61776d68..a56874ae 100644 --- a/nodescraper/connection/redfish/__init__.py +++ b/nodescraper/connection/redfish/__init__.py @@ -34,7 +34,11 @@ RF_MEMBERS_NEXT_LINK, RF_ODATA_ID, ) -from .redfish_manager import MultiTargetRedfishConnection, RedfishConnectionManager +from .redfish_manager import ( + MultiTargetRedfishConnection, + RedfishConnectionManager, + collected_multi_target_data, +) from .redfish_oem_diag import ( collect_oem_diagnostic_data, get_oem_diagnostic_allowable_values, @@ -52,6 +56,7 @@ "RedfishGetResult", "MultiTargetRedfishConnection", "RedfishConnectionManager", + "collected_multi_target_data", "RedfishConnectionParams", "RedfishTargetParams", "redfish_params_to_ssh", diff --git a/nodescraper/connection/redfish/redfish_manager.py b/nodescraper/connection/redfish/redfish_manager.py index 5ad14f99..5e7bb5f5 100644 --- a/nodescraper/connection/redfish/redfish_manager.py +++ b/nodescraper/connection/redfish/redfish_manager.py @@ -53,9 +53,31 @@ class MultiTargetRedfishConnection: wrapper intercepts collect_data and iterates the individual target connections. """ - def __init__(self, target_connections: dict[str, RedfishConnection]) -> None: + def __init__( + self, + target_connections: dict[str, RedfishConnection], + max_workers: Optional[int] = None, + failed_targets: Optional[dict[str, str]] = None, + ) -> None: self.target_connections = target_connections self.multi_target_data: dict[str, Any] = {} + self.max_workers = max_workers + self.failed_targets = dict(failed_targets or {}) + + +def collected_multi_target_data(connection_manager: Any) -> dict[str, Any]: + """Return per-target data from an open connection or from the copy made at disconnect. + + Args: + connection_manager: Redfish connection manager for this plugin. + + Returns: + dict[str, Any]: Target key to collected data. Empty for single-target runs. + """ + conn = getattr(connection_manager, "connection", None) + if isinstance(conn, MultiTargetRedfishConnection) and conn.multi_target_data: + return conn.multi_target_data + return getattr(connection_manager, "_multi_target_data", None) or {} class RedfishConnectionManager(ConnectionManager[RedfishConnection, RedfishConnectionParams]): @@ -124,23 +146,47 @@ def connect(self) -> TaskResult: if params.is_multi_target: self.target_connections: dict[str, RedfishConnection] = {} + failed_targets: dict[str, str] = {} for target in params.targets: # type: ignore[union-attr] - key, conn = self._connect_target(target) - if conn is not None: - self.target_connections[key] = conn + key, conn, error = self._connect_target(target) + if conn is None: + failed_targets[key] = error or f"Redfish connection failed for target {key!r}" + continue + if key in self.target_connections: + self._log_event( + category=EventCategory.RUNTIME, + description=f"Duplicate Redfish target key {key!r}; keeping the first connection", + priority=EventPriority.WARNING, + console_log=True, + ) + conn.close() + continue + self.target_connections[key] = conn if not self.target_connections: self.result.status = ExecutionStatus.EXECUTION_FAILURE + self.result.message = "Redfish connection failed for every target" else: # Set self.connection so DataPlugin.collect() does not short-circuit. - self.connection = MultiTargetRedfishConnection(self.target_connections) # type: ignore[assignment] + self.connection = MultiTargetRedfishConnection( # type: ignore[assignment] + self.target_connections, + max_workers=params.max_workers, + failed_targets=failed_targets, + ) return self.result return self._connect_single(params) def _connect_target( self, target: RedfishConnectionParams - ) -> tuple[str, Optional[RedfishConnection]]: - """Connect one target; returns (key, connection) or (key, None) on failure.""" + ) -> tuple[str, Optional[RedfishConnection], Optional[str]]: + """Connect one target. + + Args: + target: Connection parameters for a single BMC. + + Returns: + tuple[str, Optional[RedfishConnection], Optional[str]]: Target key, connection, and error text. + """ key = target.target_key or str(target.host) password = target.password.get_secret_value() if target.password else None base_url = _build_base_url(str(target.host), target.port, target.use_https) @@ -157,15 +203,16 @@ def _connect_target( ) conn._ensure_session() conn.get_service_root() - return key, conn + return key, conn, None except (RedfishConnectionError, Exception) as exc: # noqa: BLE001 + description = f"Redfish connection failed for target {key!r}: {exc}" self._log_event( category=EventCategory.RUNTIME, - description=f"Redfish connection failed for target {key!r}: {exc}", - priority=EventPriority.CRITICAL, + description=description, + priority=EventPriority.WARNING, console_log=True, ) - return key, None + return key, None, description def _connect_single(self, params: RedfishConnectionParams) -> TaskResult: """Connect in single-target mode (original connect logic).""" @@ -205,15 +252,12 @@ def _connect_single(self, params: RedfishConnectionParams) -> TaskResult: return self.result def disconnect(self) -> None: - """Disconnect all Redfish sessions, preserving multi-target data on the manager.""" + """Disconnect Redfish sessions and keep collected multi-target data on the manager.""" if isinstance(self.connection, MultiTargetRedfishConnection): - # Persist collected data so analyze() can access it after disconnect. self._multi_target_data: dict[str, Any] = dict(self.connection.multi_target_data) for conn in self.connection.target_connections.values(): conn.close() + self.target_connections = {} elif self.connection is not None: self.connection.close() - for conn in getattr(self, "target_connections", {}).values(): - conn.close() - self.target_connections = {} super().disconnect() diff --git a/nodescraper/connection/redfish/redfish_params.py b/nodescraper/connection/redfish/redfish_params.py index 6df829e8..9827f2e8 100644 --- a/nodescraper/connection/redfish/redfish_params.py +++ b/nodescraper/connection/redfish/redfish_params.py @@ -116,6 +116,11 @@ def redfish_params_to_ssh( if ssh_port is None: ssh_port = 22 + if params.host is None: + raise ValueError( + "SSH mapping requires a Redfish host. Multi-target configs apply only to Redfish plugins." + ) + return SSHConnectionParams( hostname=str(params.host), username=params.username, diff --git a/nodescraper/helpers/plugin_execution_target.py b/nodescraper/helpers/plugin_execution_target.py index eac1c546..0b8daffb 100644 --- a/nodescraper/helpers/plugin_execution_target.py +++ b/nodescraper/helpers/plugin_execution_target.py @@ -40,6 +40,16 @@ def _host_from_connection_args(raw: Optional[Union[dict[str, Any], BaseModel]]) if raw is None: return None if isinstance(raw, dict): + targets = raw.get("targets") + if isinstance(targets, list) and targets: + hosts = [] + for item in targets: + if isinstance(item, dict): + host = item.get("target_key") or item.get("host") + if host: + hosts.append(str(host)) + if hosts: + return ", ".join(hosts) for key in ("host", "hostname", "ip"): value = raw.get(key) if value: diff --git a/nodescraper/interfaces/dataplugin.py b/nodescraper/interfaces/dataplugin.py index da46e47b..9241ff1e 100644 --- a/nodescraper/interfaces/dataplugin.py +++ b/nodescraper/interfaces/dataplugin.py @@ -594,34 +594,37 @@ def find_datamodel_path_in_run(cls, run_path: str) -> Optional[str]: data_model_cls = getattr(cls, "DATA_MODEL", None) if not data_model_cls: return None + plugin_dir = resolve_log_dir_name(cls.__name__) + plugin_roots = [ + entry + for entry in os.listdir(run_path) + if entry == plugin_dir or entry.startswith(plugin_dir + "[") + ] for collector_cls in cls.get_collector_classes(): - collector_dir = os.path.join( - run_path, - resolve_log_dir_name(cls.__name__), - resolve_log_dir_name(collector_cls.__name__), - ) - if not os.path.isdir(collector_dir): - continue - result_path = os.path.join(collector_dir, "result.json") - if not os.path.isfile(result_path): - continue - try: - res_payload = json.loads(Path(result_path).read_text(encoding="utf-8")) - if res_payload.get("parent") != cls.__name__: + collector_name = resolve_log_dir_name(collector_cls.__name__) + for plugin_root in plugin_roots: + collector_dir = os.path.join(run_path, plugin_root, collector_name) + if not os.path.isdir(collector_dir): + continue + result_path = os.path.join(collector_dir, "result.json") + if not os.path.isfile(result_path): + continue + try: + res_payload = json.loads(Path(result_path).read_text(encoding="utf-8")) + parent = res_payload.get("parent") or "" + if parent != cls.__name__ and not str(parent).startswith(cls.__name__ + "["): + continue + except (json.JSONDecodeError, OSError): continue - except (json.JSONDecodeError, OSError): - continue - want_json = data_model_cls.__name__.lower() + ".json" - # First search all files for the json - for fname in os.listdir(collector_dir): - low = fname.lower() - if low.endswith(f"{data_model_cls.__name__.lower()}.json") or low == want_json: - return os.path.join(collector_dir, fname) - # Then search for log since that is valid in some cases - for fname in os.listdir(collector_dir): - low = fname.lower() - if low.endswith(".log"): - return os.path.join(collector_dir, fname) + want_json = data_model_cls.__name__.lower() + ".json" + for fname in os.listdir(collector_dir): + low = fname.lower() + if low.endswith(f"{data_model_cls.__name__.lower()}.json") or low == want_json: + return os.path.join(collector_dir, fname) + for fname in os.listdir(collector_dir): + low = fname.lower() + if low.endswith(".log"): + return os.path.join(collector_dir, fname) return None @classmethod diff --git a/test/unit/connection/redfish/test_redfish_multi_target.py b/test/unit/connection/redfish/test_redfish_multi_target.py new file mode 100644 index 00000000..e440158b --- /dev/null +++ b/test/unit/connection/redfish/test_redfish_multi_target.py @@ -0,0 +1,257 @@ +############################################################################### +# +# MIT License +# +# Copyright (c) 2026 Advanced Micro Devices, Inc. +# +# Permission is hereby granted, free of charge, to any person obtaining a copy +# of this software and associated documentation files (the "Software"), to deal +# in the Software without restriction, including without limitation the rights +# to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +# copies of the Software, and to permit persons to whom the Software is +# furnished to do so, subject to the following conditions: +# +# The above copyright notice and this permission notice shall be included in all +# copies or substantial portions of the Software. +# +# THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +# IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +# FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +# AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +# LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +# OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +# SOFTWARE. +# +############################################################################### +import logging +from typing import cast +from unittest.mock import MagicMock, patch + +from nodescraper.base import OOBandDataPlugin, RedfishDataCollector +from nodescraper.connection.oob_ssh.oob_ssh_connection_manager import ( + OobSshConnectionManager, +) +from nodescraper.connection.redfish import ( + MultiTargetRedfishConnection, + RedfishConnection, + RedfishConnectionError, + RedfishConnectionManager, + RedfishConnectionParams, +) +from nodescraper.enums import ExecutionStatus +from nodescraper.interfaces import DataAnalyzer +from nodescraper.models import DataModel, SystemInfo + + +class SampleModel(DataModel): + label: str = "" + + +class SampleCollector(RedfishDataCollector[SampleModel, None]): + DATA_MODEL = SampleModel + + def collect_data(self, args=None): + """Collect a label from the fake connection. + + Args: + args: Unused collection arguments. + + Returns: + tuple[TaskResult, Optional[SampleModel]]: Collector result and data model. + """ + marker = getattr(self.connection, "marker", "") + if marker == "boom": + raise RuntimeError("bmc down") + self.result.status = ExecutionStatus.OK + self.result.message = marker + return self.result, SampleModel(label=marker) + + +class SampleAnalyzer(DataAnalyzer[SampleModel, None]): + DATA_MODEL = SampleModel + + def analyze_data(self, data, args=None): + """Flag a bad label as an analysis error. + + Args: + data: Collected sample model. + args: Unused analysis arguments. + + Returns: + TaskResult: Analysis result for this target. + """ + if data.label == "bad": + self.result.status = ExecutionStatus.ERROR + self.result.message = "check failed" + else: + self.result.status = ExecutionStatus.OK + self.result.message = "ok" + return self.result + + +class SamplePlugin(OOBandDataPlugin[SampleModel, None, None]): + DATA_MODEL = SampleModel + COLLECTOR = SampleCollector + ANALYZER = SampleAnalyzer + + +def _manager( + system_info: SystemInfo, targets: list[RedfishConnectionParams], **extra +) -> RedfishConnectionManager: + return RedfishConnectionManager( + system_info=system_info, + connection_args=RedfishConnectionParams(targets=targets, **extra), + ) + + +@patch("nodescraper.connection.redfish.redfish_manager.RedfishConnection") +def test_partial_connect_keeps_reachable_bmc(mock_conn_cls, system_info: SystemInfo) -> None: + """One unreachable BMC leaves the reachable BMC connected at warning.""" + + def build(**kwargs): + conn = MagicMock() + if "down" in kwargs["base_url"]: + conn._ensure_session.side_effect = RedfishConnectionError("unreachable") + return conn + + mock_conn_cls.side_effect = build + manager = _manager( + system_info, + [ + RedfishConnectionParams(target_key="good", host="bmc-good.example", username="admin"), + RedfishConnectionParams(target_key="bad", host="bmc-down.example", username="admin"), + ], + max_workers=2, + ) + + result = manager.connect() + + assert result.status == ExecutionStatus.WARNING + assert isinstance(manager.connection, MultiTargetRedfishConnection) + assert set(manager.connection.target_connections) == {"good"} + assert "bad" in manager.connection.failed_targets + assert manager.connection.max_workers == 2 + + +@patch("nodescraper.connection.redfish.redfish_manager.RedfishConnection") +def test_every_target_down_fails_connect(mock_conn_cls, system_info: SystemInfo) -> None: + """Every target failing is an execution failure and does not set a connection.""" + mock_conn_cls.side_effect = lambda **kwargs: MagicMock( + _ensure_session=MagicMock(side_effect=RedfishConnectionError("unreachable")) + ) + manager = _manager( + system_info, + [ + RedfishConnectionParams(target_key="a", host="bmc-a.example", username="admin"), + RedfishConnectionParams(target_key="b", host="bmc-b.example", username="admin"), + ], + ) + + result = manager.connect() + + assert result.status == ExecutionStatus.EXECUTION_FAILURE + assert manager.connection is None + + +def test_one_bmc_collection_failure_stays_warning(system_info: SystemInfo) -> None: + """A collect failure and a missed connect stay a warning when another target succeeds.""" + good = MagicMock() + good.marker = "good" + bad = MagicMock() + bad.marker = "boom" + multi = MultiTargetRedfishConnection( + {"good": good, "bad": bad}, + max_workers=2, + failed_targets={"offline": "Redfish connection failed for target 'offline'"}, + ) + collector = SampleCollector( + system_info=system_info, + connection=cast(RedfishConnection, multi), + logger=logging.getLogger("multi-target-test"), + parent="SamplePlugin", + ) + + result, data = collector.collect_data() + + assert data is None + assert result.status == ExecutionStatus.WARNING + assert result.status <= ExecutionStatus.WARNING + assert multi.multi_target_data["good"].label == "good" + assert "bad" not in multi.multi_target_data + assert "offline" in result.message + + +def test_analyze_reads_data_before_disconnect(system_info: SystemInfo) -> None: + """Analysis uses data still held on the open multi-target connection.""" + multi = MultiTargetRedfishConnection({"good": MagicMock()}) + multi.multi_target_data["good"] = SampleModel(label="good") + connection_manager = MagicMock() + connection_manager.connection = multi + connection_manager._multi_target_data = {} + plugin = SamplePlugin( + system_info=system_info, + logger=logging.getLogger("multi-target-test"), + connection_manager=connection_manager, + ) + + result = plugin.analyze() + + assert result.status == ExecutionStatus.OK + + +def test_analyze_error_on_collected_data_fails(system_info: SystemInfo) -> None: + """A check failure on data that was collected still fails analysis.""" + multi = MultiTargetRedfishConnection({"bad": MagicMock(), "good": MagicMock()}) + multi.multi_target_data["bad"] = SampleModel(label="bad") + multi.multi_target_data["good"] = SampleModel(label="good") + connection_manager = MagicMock() + connection_manager.connection = multi + plugin = SamplePlugin( + system_info=system_info, + logger=logging.getLogger("multi-target-test"), + connection_manager=connection_manager, + ) + + result = plugin.analyze() + + assert result.status == ExecutionStatus.ERROR + + +def test_oob_ssh_skips_targets_without_a_single_host(system_info: SystemInfo) -> None: + """OOB SSH does not fan out across a Redfish targets list.""" + manager = OobSshConnectionManager( + system_info=system_info, + connection_args=RedfishConnectionParams( + targets=[ + RedfishConnectionParams(host="bmc-a.example", username="admin"), + RedfishConnectionParams(host="bmc-b.example", username="admin"), + ] + ), + ) + + result = manager.connect() + + assert result.status == ExecutionStatus.NOT_RAN + assert manager.connection is None + + +@patch("nodescraper.connection.oob_ssh.oob_ssh_connection_manager.RemoteShell") +def test_oob_ssh_uses_top_level_host_beside_targets(mock_shell, system_info: SystemInfo) -> None: + """A top-level host keeps OOB SSH on one BMC while Redfish uses targets.""" + mock_shell.return_value = MagicMock() + manager = OobSshConnectionManager( + system_info=system_info, + connection_args=RedfishConnectionParams( + host="bmc-ssh.example", + username="admin", + targets=[ + RedfishConnectionParams(host="bmc-a.example", username="admin"), + RedfishConnectionParams(host="bmc-b.example", username="admin"), + ], + ), + ) + + result = manager.connect() + + assert result.status == ExecutionStatus.OK + assert mock_shell.call_args.args[0].hostname == "bmc-ssh.example" From 5f5ff434b37c5c652883fbf7f808b5b9c8eb0f1c Mon Sep 17 00:00:00 2001 From: Anson Yim Date: Wed, 30 Sep 2026 14:45:27 -0400 Subject: [PATCH 5/7] fix logging and summary --- nodescraper/base/oobanddataplugin.py | 78 ++++++++++++++++++++++-- nodescraper/base/redfishcollectortask.py | 55 +++++++++++++---- 2 files changed, 116 insertions(+), 17 deletions(-) diff --git a/nodescraper/base/oobanddataplugin.py b/nodescraper/base/oobanddataplugin.py index 4d90858c..bfe877d9 100644 --- a/nodescraper/base/oobanddataplugin.py +++ b/nodescraper/base/oobanddataplugin.py @@ -23,17 +23,22 @@ # SOFTWARE. # ############################################################################### +import shutil +from pathlib import Path from typing import Any, Generic, Optional, Union +from nodescraper.base.redfishcollectortask import _target_dir_name from nodescraper.connection.redfish import ( RedfishConnectionManager, RedfishConnectionParams, collected_multi_target_data, ) -from nodescraper.enums import EventPriority, ExecutionStatus +from nodescraper.enums import EventPriority, ExecutionStatus, SystemInteractionLevel from nodescraper.generictypes import TAnalyzeArg, TCollectArg, TDataModel from nodescraper.interfaces import DataPlugin from nodescraper.models import TaskResult +from nodescraper.taskresulthooks.filesystemloghook import FileSystemLogHook +from nodescraper.utils import resolve_log_dir_name class OOBandDataPlugin( @@ -55,6 +60,55 @@ class OOBandDataPlugin( CONNECTION_TYPE = RedfishConnectionManager + def run( # type: ignore[override] + self, + collection: bool = True, + analysis: bool = True, + max_event_priority_level: Union[EventPriority, str] = EventPriority.CRITICAL, + system_interaction_level: Union[ + SystemInteractionLevel, str + ] = SystemInteractionLevel.INTERACTIVE, + preserve_connection: bool = False, + data: Optional[Any] = None, + collection_args: Optional[Any] = None, + analysis_args: Optional[Any] = None, + ): + """Run plugin. For multi-target OK runs, the summary includes per-target collection detail.""" + result = super().run( + collection=collection, + analysis=analysis, + max_event_priority_level=max_event_priority_level, + system_interaction_level=system_interaction_level, + preserve_connection=preserve_connection, + data=data, + collection_args=collection_args, + analysis_args=analysis_args, + ) + cm = self.connection_manager + if collected_multi_target_data(cm): + # DataPlugin.run() replaces the message with "Plugin tasks completed successfully" + # for OK status, discarding per-target detail. Restore it for multi-target runs. + if result.status == ExecutionStatus.OK and getattr( + self.collection_result, "message", None + ): + result.message = self.collection_result.message + + # DataPlugin.collect() unconditionally creates a shared collector directory + # (e.g. redfish_endpoint_plugin/redfish_endpoint_collector/) and writes a + # combined result.json there. In multi-target mode this directory is redundant + # and confusing alongside the per-target subdirectories, so remove it. + if self.log_path: + for collector_cls in self.get_collector_classes(): + shared_dir = ( + Path(self.log_path) + / resolve_log_dir_name(self.__class__.__name__) + / resolve_log_dir_name(collector_cls.__name__) + ) + if shared_dir.is_dir(): + shutil.rmtree(str(shared_dir), ignore_errors=True) + + return result + def analyze( self, max_event_priority_level: Optional[Union[EventPriority, str]] = EventPriority.CRITICAL, @@ -97,15 +151,31 @@ def analyze( ): analysis_args = self.ANALYZER_ARGS.model_validate(analysis_args) # type: ignore[assignment] + # Replace FileSystemLogHook instances so analyzer logs land at: + # /// + # (same directory tree used by the collector for that target) + plugin_log_dir: Optional[str] = None + if self.log_path: + plugin_log_dir = str( + Path(self.log_path) / resolve_log_dir_name(self.__class__.__name__) + ) + per_target_hooks = [ + ( + FileSystemLogHook(log_base_path=plugin_log_dir) + if (isinstance(hook, FileSystemLogHook) and plugin_log_dir) + else hook + ) + for hook in self.task_result_hooks + ] + analysis_results: list[TaskResult] = [] for target_key, target_data in multi_target_data.items(): - parent = f"{self.__class__.__name__}[{target_key}]" analyzer_task = self.ANALYZER( system_info=self.system_info.model_copy(), logger=self.logger, max_event_priority_level=max_event_priority_level or EventPriority.CRITICAL, - parent=parent, - task_result_hooks=self.task_result_hooks, + parent=_target_dir_name(target_key), # → // + task_result_hooks=per_target_hooks, event_reporter=self.event_reporter, session_id=self.session_id, ) diff --git a/nodescraper/base/redfishcollectortask.py b/nodescraper/base/redfishcollectortask.py index 8b28b2c2..1e03ef58 100644 --- a/nodescraper/base/redfishcollectortask.py +++ b/nodescraper/base/redfishcollectortask.py @@ -73,10 +73,19 @@ def combine_target_results(parent: str, results: list[TaskResult]) -> TaskResult aggregated = DataPlugin._aggregate_collection_results(parent, results) hard_fail = any(result.status in _HARD_FAIL for result in results) succeeded = any(result.status in _TARGET_SUCCESS for result in results) - if hard_fail and succeeded: - aggregated.status = ExecutionStatus.WARNING - note = "One or more targets failed; continuing with the rest." - aggregated.message = f"{note} {aggregated.message}".strip() if aggregated.message else note + + # Display one line per target so the summary table is easy to read. + messages = [r.message for r in results if r.message] + if messages: + prefix = ( + "One or more targets failed; continuing with the rest.\n" + if (hard_fail and succeeded) + else "" + ) + aggregated.message = prefix + "\n".join(messages) + if hard_fail and succeeded: + aggregated.status = ExecutionStatus.WARNING + return aggregated @@ -165,10 +174,27 @@ def _multi_target_wrapper( def _run_for_target( target_key: str, conn: RedfishConnection ) -> tuple[str, TaskResult, Any]: - target_parent = f"{parent_name}[{target_key}]" - target_log = getattr(collector, "log_path", None) - if target_log: - target_log = str(Path(target_log) / _target_dir_name(target_key)) + from nodescraper.taskresulthooks.filesystemloghook import ( + FileSystemLogHook, + ) + + safe_key = _target_dir_name(target_key) + + # Replace FileSystemLogHook instances so per-target logs land at: + # /// + # instead of the old /// + collector_log = getattr(collector, "log_path", None) + per_target_hooks: list = [] + if collector_log: + plugin_dir = str(Path(collector_log).parent) + for hook in collector.task_result_hooks: + if isinstance(hook, FileSystemLogHook): + per_target_hooks.append(FileSystemLogHook(log_base_path=plugin_dir)) + else: + per_target_hooks.append(hook) + else: + per_target_hooks = list(collector.task_result_hooks) + target_collector = type(collector)( system_info=collector.system_info.model_copy(), connection=conn, @@ -176,20 +202,23 @@ def _run_for_target( max_event_priority_level=getattr( collector, "max_event_priority_level", EventPriority.CRITICAL ), - parent=target_parent, - task_result_hooks=collector.task_result_hooks, + parent=safe_key, # → /// + task_result_hooks=per_target_hooks, event_reporter=collector.event_reporter, session_id=collector.session_id, - log_path=target_log, + log_path=None, # avoid duplicate FileSystemLogHook from DataCollector system_interaction_level=getattr(collector, "system_interaction_level", None), ) - if target_log: - target_collector.log_path = target_log collector.logger.info( "Starting collection for target %r using %s", target_key, parent_name ) result, data = _fn(target_collector, args) collector.logger.info("Finished collection for target %r", target_key) + + # Prefix result message with target key for attribution in the summary table. + if result.message: + result.message = f"[{target_key}] {result.message}" + return target_key, result, data if multi_conn.target_connections: From 8a548db574d785ae4812c5cd4b743421c8254646 Mon Sep 17 00:00:00 2001 From: Alexandra Bara Date: Thu, 1 Oct 2026 11:28:51 -0500 Subject: [PATCH 6/7] multi targets updates for ssh proxy connection + enhancement for err from redfish --- README.md | 4 +- nodescraper/base/oobanddataplugin.py | 50 ++-- nodescraper/base/redfishcollectortask.py | 29 +-- nodescraper/connection/inband/inbandremote.py | 4 + .../connection/redfish/redfish_manager.py | 6 +- .../connection/redfish/redfish_oem_diag.py | 20 +- .../redfish/ssh_proxy_connection.py | 12 +- .../connection/redfish/ssh_proxy_manager.py | 239 ++++++++++++++---- .../connection/redfish/ssh_proxy_params.py | 56 +++- nodescraper/interfaces/dataplugin.py | 55 +++- .../amc_redfish_diag/amc_diag_collector.py | 4 +- .../amc_redfish_diag/amc_diag_plugin.py | 3 +- .../redfish_oem_diag/oem_diag_collector.py | 10 +- .../taskresulthooks/filesystemloghook.py | 32 ++- .../connection/inband/test_shellcommand.py | 13 + .../redfish/test_redfish_multi_target.py | 45 ++++ .../redfish/test_redfish_oem_diag.py | 23 ++ .../redfish/test_ssh_proxy_connection.py | 115 +++++++++ test/unit/plugin/test_amc_redfish_diag.py | 77 +++++- .../plugin/test_redfish_oem_diag_collector.py | 2 +- 20 files changed, 643 insertions(+), 156 deletions(-) diff --git a/README.md b/README.md index 4ff50d3c..158a5845 100644 --- a/README.md +++ b/README.md @@ -256,7 +256,7 @@ A target that fails to connect or collect does not fail the run when another tar - `targets`: list of per-target connection parameters. Each entry accepts the same fields as the single-target config plus an optional `target_key` (used as the result key; defaults to the host string). - `max_workers` (optional): maximum concurrent collection threads. Defaults to `min(len(targets), 32)` and is capped at 32. -Per-target results are written to `[]//`. +Per-target results are written to `///`. Single-target results stay in `//`. The same layout is used for the analyzer directory and for the AMC SSH-proxy plugin. A ready-to-edit sample is in `config/connection-config_redfish_multi_target.example.json`. Replace the example hosts and password, then pass it with `--connection-config`. @@ -522,7 +522,7 @@ Use a plugin config that points at your LogService and lists the types to collec The RedfishEndpointPlugin collects Redfish URIs (GET responses) and optionally runs checks on the returned JSON. It requires a Redfish connection config (same as RedfishOemDiagPlugin). -**Multi-target support:** `RedfishEndpointPlugin` collects from each BMC in the `targets` list at the same time. Use the [Redfish multi-target](#redfish-multi-target) connection config. The same `uris` and `checks` apply to every target. Per-target results are written to `redfish_endpoint_plugin[]/redfish_endpoint_collector/` under the run log directory. A BMC that cannot be reached is reported as a warning when another target succeeds. +**Multi-target support:** `RedfishEndpointPlugin` collects from each BMC in the `targets` list at the same time. Use the [Redfish multi-target](#redfish-multi-target) connection config. The same `uris` and `checks` apply to every target. Per-target results are written to `redfish_endpoint_plugin/redfish_endpoint_collector//` under the run log directory. A BMC that cannot be reached is reported as a warning when another target succeeds. **How to run** diff --git a/nodescraper/base/oobanddataplugin.py b/nodescraper/base/oobanddataplugin.py index bfe877d9..c1d9905c 100644 --- a/nodescraper/base/oobanddataplugin.py +++ b/nodescraper/base/oobanddataplugin.py @@ -23,7 +23,6 @@ # SOFTWARE. # ############################################################################### -import shutil from pathlib import Path from typing import Any, Generic, Optional, Union @@ -37,7 +36,7 @@ from nodescraper.generictypes import TAnalyzeArg, TCollectArg, TDataModel from nodescraper.interfaces import DataPlugin from nodescraper.models import TaskResult -from nodescraper.taskresulthooks.filesystemloghook import FileSystemLogHook +from nodescraper.taskresulthooks.filesystemloghook import hooks_for_fixed_directory from nodescraper.utils import resolve_log_dir_name @@ -93,20 +92,6 @@ def run( # type: ignore[override] ): result.message = self.collection_result.message - # DataPlugin.collect() unconditionally creates a shared collector directory - # (e.g. redfish_endpoint_plugin/redfish_endpoint_collector/) and writes a - # combined result.json there. In multi-target mode this directory is redundant - # and confusing alongside the per-target subdirectories, so remove it. - if self.log_path: - for collector_cls in self.get_collector_classes(): - shared_dir = ( - Path(self.log_path) - / resolve_log_dir_name(self.__class__.__name__) - / resolve_log_dir_name(collector_cls.__name__) - ) - if shared_dir.is_dir(): - shutil.rmtree(str(shared_dir), ignore_errors=True) - return result def analyze( @@ -151,31 +136,28 @@ def analyze( ): analysis_args = self.ANALYZER_ARGS.model_validate(analysis_args) # type: ignore[assignment] - # Replace FileSystemLogHook instances so analyzer logs land at: - # /// - # (same directory tree used by the collector for that target) - plugin_log_dir: Optional[str] = None - if self.log_path: - plugin_log_dir = str( - Path(self.log_path) / resolve_log_dir_name(self.__class__.__name__) - ) - per_target_hooks = [ - ( - FileSystemLogHook(log_base_path=plugin_log_dir) - if (isinstance(hook, FileSystemLogHook) and plugin_log_dir) - else hook - ) - for hook in self.task_result_hooks - ] + plugin_log_dir = ( + Path(self.log_path) / resolve_log_dir_name(self.__class__.__name__) + if self.log_path + else None + ) + analyzer_name = resolve_log_dir_name(self.ANALYZER.__name__) analysis_results: list[TaskResult] = [] for target_key, target_data in multi_target_data.items(): + safe_key = _target_dir_name(target_key) + target_hooks = self.task_result_hooks + if plugin_log_dir is not None: + target_hooks = hooks_for_fixed_directory( + self.task_result_hooks, + str(plugin_log_dir / analyzer_name / safe_key), + ) analyzer_task = self.ANALYZER( system_info=self.system_info.model_copy(), logger=self.logger, max_event_priority_level=max_event_priority_level or EventPriority.CRITICAL, - parent=_target_dir_name(target_key), # → // - task_result_hooks=per_target_hooks, + parent=safe_key, + task_result_hooks=target_hooks, event_reporter=self.event_reporter, session_id=self.session_id, ) diff --git a/nodescraper/base/redfishcollectortask.py b/nodescraper/base/redfishcollectortask.py index 1e03ef58..b24e0874 100644 --- a/nodescraper/base/redfishcollectortask.py +++ b/nodescraper/base/redfishcollectortask.py @@ -41,6 +41,7 @@ from nodescraper.interfaces import DataCollector, TaskResultHook from nodescraper.interfaces.dataplugin import DataPlugin from nodescraper.models import SystemInfo, TaskResult +from nodescraper.taskresulthooks.filesystemloghook import hooks_for_fixed_directory _HARD_FAIL = {ExecutionStatus.ERROR, ExecutionStatus.EXECUTION_FAILURE} _TARGET_SUCCESS = {ExecutionStatus.OK, ExecutionStatus.WARNING} @@ -174,26 +175,14 @@ def _multi_target_wrapper( def _run_for_target( target_key: str, conn: RedfishConnection ) -> tuple[str, TaskResult, Any]: - from nodescraper.taskresulthooks.filesystemloghook import ( - FileSystemLogHook, - ) - safe_key = _target_dir_name(target_key) - - # Replace FileSystemLogHook instances so per-target logs land at: - # /// - # instead of the old /// collector_log = getattr(collector, "log_path", None) - per_target_hooks: list = [] - if collector_log: - plugin_dir = str(Path(collector_log).parent) - for hook in collector.task_result_hooks: - if isinstance(hook, FileSystemLogHook): - per_target_hooks.append(FileSystemLogHook(log_base_path=plugin_dir)) - else: - per_target_hooks.append(hook) - else: - per_target_hooks = list(collector.task_result_hooks) + target_log = str(Path(collector_log) / safe_key) if collector_log else None + per_target_hooks = ( + hooks_for_fixed_directory(collector.task_result_hooks, target_log) + if target_log + else list(collector.task_result_hooks) + ) target_collector = type(collector)( system_info=collector.system_info.model_copy(), @@ -202,11 +191,11 @@ def _run_for_target( max_event_priority_level=getattr( collector, "max_event_priority_level", EventPriority.CRITICAL ), - parent=safe_key, # → /// + parent=safe_key, task_result_hooks=per_target_hooks, event_reporter=collector.event_reporter, session_id=collector.session_id, - log_path=None, # avoid duplicate FileSystemLogHook from DataCollector + log_path=target_log, system_interaction_level=getattr(collector, "system_interaction_level", None), ) collector.logger.info( diff --git a/nodescraper/connection/inband/inbandremote.py b/nodescraper/connection/inband/inbandremote.py index 282e3c47..c79d5f03 100644 --- a/nodescraper/connection/inband/inbandremote.py +++ b/nodescraper/connection/inband/inbandremote.py @@ -172,6 +172,10 @@ def run_command( stderr_str = "Command timed out" stdout_str = "" exit_code = 124 + except SSHException as exc: + stderr_str = str(exc) or "SSH channel failed" + stdout_str = "" + exit_code = 124 return CommandArtifact( command=cmd_str, diff --git a/nodescraper/connection/redfish/redfish_manager.py b/nodescraper/connection/redfish/redfish_manager.py index 5e7bb5f5..fa278c3a 100644 --- a/nodescraper/connection/redfish/redfish_manager.py +++ b/nodescraper/connection/redfish/redfish_manager.py @@ -26,7 +26,7 @@ from __future__ import annotations from logging import Logger -from typing import Any, Optional, Union +from typing import Any, Mapping, Optional, Union from nodescraper.enums import EventCategory, EventPriority, ExecutionStatus from nodescraper.interfaces.connectionmanager import ConnectionManager @@ -55,11 +55,11 @@ class MultiTargetRedfishConnection: def __init__( self, - target_connections: dict[str, RedfishConnection], + target_connections: Mapping[str, RedfishConnection], max_workers: Optional[int] = None, failed_targets: Optional[dict[str, str]] = None, ) -> None: - self.target_connections = target_connections + self.target_connections = dict(target_connections) self.multi_target_data: dict[str, Any] = {} self.max_workers = max_workers self.failed_targets = dict(failed_targets or {}) diff --git a/nodescraper/connection/redfish/redfish_oem_diag.py b/nodescraper/connection/redfish/redfish_oem_diag.py index 3cd418e6..a04fd6f3 100644 --- a/nodescraper/connection/redfish/redfish_oem_diag.py +++ b/nodescraper/connection/redfish/redfish_oem_diag.py @@ -207,7 +207,10 @@ def _poll_task_resource( while True: if time.time() - start > timeout_s: return None, f"Task did not complete within {timeout_s}s" - poll_resp = conn.get_response(task_path) + try: + poll_resp = conn.get_response(task_path) + except RedfishConnectionError as exc: + return None, str(exc) if poll_resp.status_code == codes.ok: try: body = poll_resp.json() @@ -551,7 +554,10 @@ def collect_oem_diagnostic_data( if time.time() - start > task_timeout_s: return None, None, f"Task did not complete within {task_timeout_s}s" monitor_path = _get_path_from_connection(conn, task_monitor) - poll_resp = conn.get_response(monitor_path) + try: + poll_resp = conn.get_response(monitor_path) + except RedfishConnectionError as exc: + return None, None, str(exc) if poll_resp.status_code == codes.not_found: return None, None, f"TaskMonitor GET failed: status {codes.not_found}" if poll_resp.status_code != codes.accepted: @@ -578,7 +584,10 @@ def collect_oem_diagnostic_data( if poll_err: return None, None, poll_err else: - task_resp = conn.get_response(follow_path) + try: + task_resp = conn.get_response(follow_path) + except RedfishConnectionError as exc: + return None, None, str(exc) if task_resp.status_code != codes.ok: return None, None, f"Task GET failed: {task_resp.status_code}" task_json = task_resp.json() @@ -624,5 +633,8 @@ def collect_oem_diagnostic_data( return None, None, f"LogEntry GET failed: {err} (GET {log_entry_path})" file_stem = oem_type or diag_type - log_bytes = _download_log_and_save(conn, log_entry_json, file_stem, output_dir, log) + try: + log_bytes = _download_log_and_save(conn, log_entry_json, file_stem, output_dir, log) + except RedfishConnectionError as exc: + return None, None, str(exc) return log_bytes, log_entry_json, None diff --git a/nodescraper/connection/redfish/ssh_proxy_connection.py b/nodescraper/connection/redfish/ssh_proxy_connection.py index 77c9f87c..a58f86da 100644 --- a/nodescraper/connection/redfish/ssh_proxy_connection.py +++ b/nodescraper/connection/redfish/ssh_proxy_connection.py @@ -31,6 +31,7 @@ from typing import Any, Optional, Union from urllib.parse import urljoin, urlparse +from paramiko.ssh_exception import SSHException from requests.structures import CaseInsensitiveDict from nodescraper.connection.inband.inband import BinaryFileArtifact, CommandArtifact @@ -140,7 +141,8 @@ def _mktemp(self) -> str: artifact = self._shell.run_command("mktemp", timeout=self._cmd_timeout()) path = (artifact.stdout or "").strip() if artifact.exit_code != 0 or not path: - raise RedfishConnectionError("Failed to create remote temp file via mktemp") + detail = (artifact.stderr or "").strip() + raise RedfishConnectionError(detail or "Failed to create remote temp file via mktemp") return path def _rm(self, *paths: str) -> None: @@ -149,9 +151,11 @@ def _rm(self, *paths: str) -> None: self._shell.run_command(f"rm -f {quoted}", timeout=self._cmd_timeout()) def _curl(self, url: str, extra_flags: str = "") -> CurlResponse: - header_path = self._mktemp() - body_path = self._mktemp() + header_path = "" + body_path = "" try: + header_path = self._mktemp() + body_path = self._mktemp() max_time = max(int(self.timeout), 1) cmd = ( f"curl -sS --max-time {max_time} -D {shlex.quote(header_path)} " @@ -186,6 +190,8 @@ def _curl(self, url: str, extra_flags: str = "") -> CurlResponse: content=body, headers=parse_curl_headers(header_text), ) + except SSHException as exc: + raise RedfishConnectionError(f"SSH channel failed: {exc}") from exc finally: self._rm(header_path, body_path) diff --git a/nodescraper/connection/redfish/ssh_proxy_manager.py b/nodescraper/connection/redfish/ssh_proxy_manager.py index 0c601af4..78a06c20 100644 --- a/nodescraper/connection/redfish/ssh_proxy_manager.py +++ b/nodescraper/connection/redfish/ssh_proxy_manager.py @@ -36,7 +36,7 @@ from ..inband.inbandremote import RemoteShell, SSHConnectionError from .redfish_connection import RedfishConnectionError -from .redfish_manager import _build_base_url +from .redfish_manager import MultiTargetRedfishConnection, _build_base_url from .ssh_proxy_connection import SshProxyRedfishConnection from .ssh_proxy_params import RedfishSshProxyConnectionParams @@ -70,7 +70,20 @@ def connect(self) -> TaskResult: """Open BMC SSH and verify AMC Redfish via curl. Returns: - TaskResult for the connection attempt. + TaskResult: Connection result. One failed target is a warning when another connects. + """ + params = self._validated_params() + if params is None: + return self.result + if params.is_multi_target: + return self._connect_multi(params) + return self._connect_single(params) + + def _validated_params(self) -> Optional[RedfishSshProxyConnectionParams]: + """Return connection params or record a failure when they are missing. + + Returns: + Optional[RedfishSshProxyConnectionParams]: Parsed params, or None when invalid. """ if not self.connection_args: self._log_event( @@ -80,44 +93,36 @@ def connect(self) -> TaskResult: console_log=True, ) self.result.status = ExecutionStatus.EXECUTION_FAILURE - return self.result + return None raw = self.connection_args if isinstance(raw, dict): - params = RedfishSshProxyConnectionParams.model_validate(raw) - elif isinstance(raw, RedfishSshProxyConnectionParams): - params = raw - else: - self._log_event( - category=EventCategory.RUNTIME, - description=( - "Redfish SSH-proxy connection_args must be dict or " - "RedfishSshProxyConnectionParams" - ), - priority=EventPriority.CRITICAL, - console_log=True, - ) - self.result.status = ExecutionStatus.EXECUTION_FAILURE - return self.result + return RedfishSshProxyConnectionParams.model_validate(raw) + if isinstance(raw, RedfishSshProxyConnectionParams): + return raw + self._log_event( + category=EventCategory.RUNTIME, + description=( + "Redfish SSH-proxy connection_args must be dict or " + "RedfishSshProxyConnectionParams" + ), + priority=EventPriority.CRITICAL, + console_log=True, + ) + self.result.status = ExecutionStatus.EXECUTION_FAILURE + return None - base_url = _build_base_url(str(params.host), params.port, params.use_https) - shell: Optional[RemoteShell] = None + def _connect_single(self, params: RedfishSshProxyConnectionParams) -> TaskResult: + """Connect one BMC SSH proxy. + + Args: + params: Single-target SSH-proxy parameters. + + Returns: + TaskResult: Connection result. + """ try: - self.logger.info( - "Connecting SSH proxy Redfish: ssh=%s curl=%s", - params.ssh.hostname, - base_url, - ) - shell = RemoteShell(params.ssh) - shell.connect_ssh() - conn = SshProxyRedfishConnection( - shell=shell, - base_url=base_url, - timeout=params.timeout_seconds, - api_root=params.api_root, - ) - conn.get_service_root() - self.connection = conn + self.connection = self._open_proxy(params) except SSHConnectionError as exc: self._log_event( category=EventCategory.SSH, @@ -127,11 +132,6 @@ def connect(self) -> TaskResult: ) self.result.status = ExecutionStatus.EXECUTION_FAILURE self.connection = None - if shell is not None: - try: - shell.client.close() - except Exception: - pass except RedfishConnectionError as exc: self._log_event( category=EventCategory.RUNTIME, @@ -141,11 +141,6 @@ def connect(self) -> TaskResult: ) self.result.status = ExecutionStatus.EXECUTION_FAILURE self.connection = None - if shell is not None: - try: - shell.client.close() - except Exception: - pass except Exception as exc: self._log_event( category=EventCategory.RUNTIME, @@ -156,20 +151,156 @@ def connect(self) -> TaskResult: ) self.result.status = ExecutionStatus.EXECUTION_FAILURE self.connection = None - if shell is not None: + return self.result + + def _connect_multi(self, params: RedfishSshProxyConnectionParams) -> TaskResult: + """Connect each SSH-proxy target. One failure does not drop the others. + + Args: + params: Parameters whose targets list is set. + + Returns: + TaskResult: Warning when any target connects, execution failure when none do. + """ + self.target_connections: dict[str, SshProxyRedfishConnection] = {} + failed_targets: dict[str, str] = {} + for target in params.targets or []: + key, conn, error = self._connect_target(target) + if conn is None: + failed_targets[key] = error or f"SSH-proxy connection failed for target {key!r}" + continue + if key in self.target_connections: + self._log_event( + category=EventCategory.RUNTIME, + description=( + f"Duplicate SSH-proxy target key {key!r}; keeping the first connection" + ), + priority=EventPriority.WARNING, + console_log=True, + ) + self._close_proxy(conn) + continue + self.target_connections[key] = conn + if not self.target_connections: + self.result.status = ExecutionStatus.EXECUTION_FAILURE + self.result.message = "SSH-proxy Redfish connection failed for every target" + return self.result + self.connection = MultiTargetRedfishConnection( # type: ignore[assignment] + self.target_connections, + max_workers=params.max_workers, + failed_targets=failed_targets, + ) + return self.result + + def _connect_target( + self, target: RedfishSshProxyConnectionParams + ) -> tuple[str, Optional[SshProxyRedfishConnection], Optional[str]]: + """Connect one target and keep going when it fails. + + Args: + target: SSH-proxy parameters for a single BMC. + + Returns: + tuple[str, Optional[SshProxyRedfishConnection], Optional[str]]: Key, connection, and error. + """ + key = _proxy_target_key(target) + try: + return key, self._open_proxy(target), None + except Exception as exc: # noqa: BLE001 + description = f"SSH-proxy Redfish connection failed for target {key!r}: {exc}" + self._log_event( + category=EventCategory.RUNTIME, + description=description, + priority=EventPriority.WARNING, + console_log=True, + ) + return key, None, description + + def _open_proxy(self, params: RedfishSshProxyConnectionParams) -> SshProxyRedfishConnection: + """Open one BMC SSH session and verify the AMC Redfish root. + + Args: + params: Single-target SSH-proxy parameters. + + Returns: + SshProxyRedfishConnection: Open connection. + + Raises: + ValueError: When host or ssh is missing. + """ + if params.host is None or params.ssh is None: + raise ValueError("SSH-proxy target requires host and ssh") + base_url = _build_base_url(str(params.host), params.port, params.use_https) + self.logger.info( + "Connecting SSH proxy Redfish: ssh=%s curl=%s", + params.ssh.hostname, + base_url, + ) + shell = RemoteShell(params.ssh) + conn: Optional[SshProxyRedfishConnection] = None + try: + shell.connect_ssh() + conn = SshProxyRedfishConnection( + shell=shell, + base_url=base_url, + timeout=params.timeout_seconds, + api_root=params.api_root, + ) + conn.get_service_root() + return conn + except Exception: + if conn is not None: + self._close_proxy(conn) + else: try: shell.client.close() except Exception: pass - return self.result + raise + + @staticmethod + def _close_proxy(conn: SshProxyRedfishConnection) -> None: + """Close the curl wrapper and the BMC SSH session. + + Args: + conn: SSH-proxy Redfish connection to close. + + Returns: + None. + """ + conn.close() + try: + conn._shell.client.close() + except Exception: + pass def disconnect(self) -> None: - """Close the curl wrapper and the BMC SSH session.""" + """Close every SSH-proxy session and keep collected multi-target data.""" conn = self.connection - if isinstance(conn, SshProxyRedfishConnection): - conn.close() - try: - conn._shell.client.close() - except Exception: - pass + if isinstance(conn, MultiTargetRedfishConnection): + self._multi_target_data: dict = dict(conn.multi_target_data) + for target_conn in conn.target_connections.values(): + if isinstance(target_conn, SshProxyRedfishConnection): + self._close_proxy(target_conn) + else: + target_conn.close() + self.target_connections = {} + elif isinstance(conn, SshProxyRedfishConnection): + self._close_proxy(conn) super().disconnect() + + +def _proxy_target_key(target: RedfishSshProxyConnectionParams) -> str: + """Return the target key, falling back to the SSH hostname. + + Args: + target: One SSH-proxy target. + + Returns: + str: target_key, otherwise the SSH hostname, otherwise the AMC host. + """ + if target.target_key: + return target.target_key + if target.ssh is not None: + return str(target.ssh.hostname) + return str(target.host) diff --git a/nodescraper/connection/redfish/ssh_proxy_params.py b/nodescraper/connection/redfish/ssh_proxy_params.py index 1c7f1533..426275b5 100644 --- a/nodescraper/connection/redfish/ssh_proxy_params.py +++ b/nodescraper/connection/redfish/ssh_proxy_params.py @@ -27,7 +27,7 @@ from typing import Optional, Union -from pydantic import BaseModel, ConfigDict, Field +from pydantic import BaseModel, ConfigDict, Field, model_validator from pydantic.networks import IPvAnyAddress from nodescraper.connection.inband.sshparams import SSHConnectionParams @@ -36,15 +36,36 @@ class RedfishSshProxyConnectionParams(BaseModel): - """Redfish over SSH: curl on the BMC to an AMC (or other) address reachable from that host.""" + """Redfish over SSH: curl on the BMC to an AMC address reachable from that host. + + Single-target mode supplies host and ssh. Multi-target mode supplies targets, + and each entry has its own ssh block and AMC host. + """ model_config = ConfigDict(arbitrary_types_allowed=True) - host: Union[IPvAnyAddress, str] = Field( + target_key: Optional[str] = Field( + default=None, + description="Identifier used when this entry appears inside a targets list.", + ) + targets: Optional[list["RedfishSshProxyConnectionParams"]] = Field( + default=None, + description="BMC SSH-proxy targets. Each entry has its own ssh block and AMC host.", + ) + max_workers: Optional[int] = Field( + default=None, + ge=1, + description=( + "Max concurrent collection threads. Defaults to the target count, capped at 32." + ), + ) + host: Optional[Union[IPvAnyAddress, str]] = Field( + default=None, description="Redfish host as seen from the SSH target (internal AMC address).", ) - ssh: SSHConnectionParams = Field( - description="SSH parameters for the BMC (or other host) that can reach host.", + ssh: Optional[SSHConnectionParams] = Field( + default=None, + description="SSH parameters for the BMC that can reach host.", ) port: Optional[int] = Field(default=80, ge=1, le=65535) use_https: bool = Field( @@ -56,3 +77,28 @@ class RedfishSshProxyConnectionParams(BaseModel): default=DEFAULT_REDFISH_API_ROOT, description="Redfish API path (e.g. redfish/v1).", ) + + @model_validator(mode="after") + def _validate_target_config(self) -> "RedfishSshProxyConnectionParams": + """Require host and ssh unless a targets list is set. + + Returns: + RedfishSshProxyConnectionParams: This params object. + """ + if self.targets: + return self + if self.host is None or self.ssh is None: + raise ValueError("Either targets or both host and ssh must be provided.") + return self + + @property + def is_multi_target(self) -> bool: + """True when one or more targets are configured via the targets list. + + Returns: + bool: True when targets is non-empty. + """ + return bool(self.targets) + + +RedfishSshProxyConnectionParams.model_rebuild() diff --git a/nodescraper/interfaces/dataplugin.py b/nodescraper/interfaces/dataplugin.py index 9241ff1e..d81f3dc4 100644 --- a/nodescraper/interfaces/dataplugin.py +++ b/nodescraper/interfaces/dataplugin.py @@ -604,11 +604,40 @@ def find_datamodel_path_in_run(cls, run_path: str) -> Optional[str]: collector_name = resolve_log_dir_name(collector_cls.__name__) for plugin_root in plugin_roots: collector_dir = os.path.join(run_path, plugin_root, collector_name) - if not os.path.isdir(collector_dir): - continue - result_path = os.path.join(collector_dir, "result.json") - if not os.path.isfile(result_path): - continue + found = cls._datamodel_file_in_collector_tree( + collector_dir, data_model_cls.__name__ + ) + if found: + return found + return None + + @classmethod + def _datamodel_file_in_collector_tree( + cls, collector_dir: str, data_model_name: str + ) -> Optional[str]: + """Return a datamodel file in a collector directory or one target subdirectory. + + Args: + collector_dir: Plugin collector directory, which may contain target folders. + data_model_name: DATA_MODEL class name used to match the JSON file. + + Returns: + Optional[str]: Absolute path to the datamodel file, or None. + """ + if not os.path.isdir(collector_dir): + return None + search_dirs = [collector_dir] + search_dirs.extend( + os.path.join(collector_dir, child) + for child in sorted(os.listdir(collector_dir)) + if os.path.isdir(os.path.join(collector_dir, child)) + ) + want_json = data_model_name.lower() + ".json" + for search_dir in search_dirs: + result_path = os.path.join(search_dir, "result.json") + if not os.path.isfile(result_path): + continue + if search_dir == collector_dir: try: res_payload = json.loads(Path(result_path).read_text(encoding="utf-8")) parent = res_payload.get("parent") or "" @@ -616,15 +645,13 @@ def find_datamodel_path_in_run(cls, run_path: str) -> Optional[str]: continue except (json.JSONDecodeError, OSError): continue - want_json = data_model_cls.__name__.lower() + ".json" - for fname in os.listdir(collector_dir): - low = fname.lower() - if low.endswith(f"{data_model_cls.__name__.lower()}.json") or low == want_json: - return os.path.join(collector_dir, fname) - for fname in os.listdir(collector_dir): - low = fname.lower() - if low.endswith(".log"): - return os.path.join(collector_dir, fname) + for fname in os.listdir(search_dir): + low = fname.lower() + if low.endswith(f"{data_model_name.lower()}.json") or low == want_json: + return os.path.join(search_dir, fname) + for fname in os.listdir(search_dir): + if fname.lower().endswith(".log"): + return os.path.join(search_dir, fname) return None @classmethod diff --git a/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_collector.py b/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_collector.py index 4d2f57bb..41397fb2 100644 --- a/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_collector.py +++ b/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_collector.py @@ -92,7 +92,6 @@ def collect_data( if self.log_path: output_dir = (Path(self.log_path) / "diag_logs").resolve() - output_dir.mkdir(parents=True, exist_ok=True) self.logger.info( "(AmcRedfishDiagPlugin) Diagnostic archives will be written to: %s", output_dir, @@ -126,8 +125,9 @@ def collect_data( if collect_err: self._log_event( category=EventCategory.RUNTIME, - description=f"AMC diag {key}: {collect_err}", + description=f"AMC diag {key}: error", priority=EventPriority.WARNING, + data={"response_body": collect_err}, console_log=True, ) results[key] = OemDiagTypeResult(success=False, error=collect_err, metadata=None) diff --git a/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_plugin.py b/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_plugin.py index 24eeb877..8804b018 100644 --- a/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_plugin.py +++ b/nodescraper/plugins/ooband/amc_redfish_diag/amc_diag_plugin.py @@ -51,7 +51,8 @@ class AmcRedfishDiagPlugin( collector reports ERROR when no collection succeeded, OK when at least one succeeded. Configure RedfishSshProxyConnectionManager: ssh to the BMC, host/port of the AMC Redfish - URL reachable from that BMC. collection_args selects Managers Manager dump and Systems OEM AllLogs. + URL reachable from that BMC. A targets list collects once per BMC. collection_args selects + Managers Manager dump and Systems OEM AllLogs. """ CONNECTION_TYPE = RedfishSshProxyConnectionManager diff --git a/nodescraper/plugins/ooband/redfish_oem_diag/oem_diag_collector.py b/nodescraper/plugins/ooband/redfish_oem_diag/oem_diag_collector.py index 968eca51..3eca21c5 100644 --- a/nodescraper/plugins/ooband/redfish_oem_diag/oem_diag_collector.py +++ b/nodescraper/plugins/ooband/redfish_oem_diag/oem_diag_collector.py @@ -62,15 +62,12 @@ def collect_data( if self.log_path: output_dir = (Path(self.log_path) / "diag_logs").resolve() - else: - output_dir = None - - if output_dir is not None: - output_dir.mkdir(parents=True, exist_ok=True) self.logger.info( "(RedfishOemDiagPlugin) Diagnostic archives (e.g. *.tar.xz) will be written to: %s", output_dir, ) + else: + output_dir = None results: dict[str, OemDiagTypeResult] = {} validate = bool(args.oem_diagnostic_types_allowable) @@ -91,8 +88,9 @@ def collect_data( if err: self._log_event( category=EventCategory.RUNTIME, - description=f"OEM diag {oem_type!r}: {err}", + description=f"OEM diag {oem_type!r}: error", priority=EventPriority.WARNING, + data={"response_body": err}, console_log=True, ) results[oem_type] = OemDiagTypeResult(success=False, error=err, metadata=None) diff --git a/nodescraper/taskresulthooks/filesystemloghook.py b/nodescraper/taskresulthooks/filesystemloghook.py index bf407c26..7e90ab54 100644 --- a/nodescraper/taskresulthooks/filesystemloghook.py +++ b/nodescraper/taskresulthooks/filesystemloghook.py @@ -32,27 +32,49 @@ class FileSystemLogHook(TaskResultHook): - def __init__(self, log_base_path=None, **kwargs) -> None: + def __init__(self, log_base_path=None, include_task_path: bool = True, **kwargs) -> None: """Create a FileSystemLogHook Instance Args: log_base_path (Optional[str], optional): The base path where logs will be stored. Defaults to the current working directory. + include_task_path: Append parent and task directory names when True. **kwargs: Additional keyword arguments, which are not used. """ if log_base_path is None: log_base_path = os.getcwd() self.log_base_path = log_base_path + self.include_task_path = include_task_path def process_result(self, task_result: TaskResult, data: Optional[DataModel] = None, **kwargs): """Log task result to the filesystem (single events.json per directory).""" log_path = self.log_base_path - if task_result.parent: - log_path = os.path.join(log_path, resolve_log_dir_name(task_result.parent)) - if task_result.task: - log_path = os.path.join(log_path, resolve_log_dir_name(task_result.task)) + if self.include_task_path: + if task_result.parent: + log_path = os.path.join(log_path, resolve_log_dir_name(task_result.parent)) + if task_result.task: + log_path = os.path.join(log_path, resolve_log_dir_name(task_result.task)) task_result.log_result(log_path) if data: data.log_model(log_path) + + +def hooks_for_fixed_directory(hooks: list, directory: str) -> list: + """Point filesystem log hooks at one directory. + + Args: + hooks: Hooks copied from the parent task. + directory: Directory that should receive this task's logs. + + Returns: + list: Hooks with filesystem logs written directly into directory. + """ + rewritten = [] + for hook in hooks: + if isinstance(hook, FileSystemLogHook): + rewritten.append(FileSystemLogHook(log_base_path=directory, include_task_path=False)) + else: + rewritten.append(hook) + return rewritten diff --git a/test/unit/connection/inband/test_shellcommand.py b/test/unit/connection/inband/test_shellcommand.py index ead564a2..e969c1bd 100644 --- a/test/unit/connection/inband/test_shellcommand.py +++ b/test/unit/connection/inband/test_shellcommand.py @@ -26,6 +26,8 @@ import subprocess from unittest.mock import MagicMock, patch +from paramiko.ssh_exception import SSHException + from nodescraper.connection.inband.inbandlocal import LocalShell from nodescraper.connection.inband.inbandremote import RemoteShell from nodescraper.connection.inband.sshparams import SSHConnectionParams @@ -79,3 +81,14 @@ def test_remoteshell_sudo_password_is_newline_terminated(mock_client_cls): shell.run_command("cat /etc/shadow", sudo=True) stdin.write.assert_called_once_with("hunter2\n") + + +@patch("nodescraper.connection.inband.inbandremote.paramiko.SSHClient") +def test_remoteshell_ssh_exception_is_a_failed_command(mock_client_cls): + mock_client_cls.return_value.exec_command.side_effect = SSHException("Timeout opening channel.") + shell = RemoteShell( + SSHConnectionParams(hostname="127.0.0.1", username="user", password="secret") + ) + artifact = shell.run_command("mktemp") + assert artifact.exit_code == 124 + assert artifact.stderr == "Timeout opening channel." diff --git a/test/unit/connection/redfish/test_redfish_multi_target.py b/test/unit/connection/redfish/test_redfish_multi_target.py index e440158b..60510ff8 100644 --- a/test/unit/connection/redfish/test_redfish_multi_target.py +++ b/test/unit/connection/redfish/test_redfish_multi_target.py @@ -41,6 +41,7 @@ from nodescraper.enums import ExecutionStatus from nodescraper.interfaces import DataAnalyzer from nodescraper.models import DataModel, SystemInfo +from nodescraper.taskresulthooks.filesystemloghook import FileSystemLogHook class SampleModel(DataModel): @@ -181,6 +182,50 @@ def test_one_bmc_collection_failure_stays_warning(system_info: SystemInfo) -> No assert "offline" in result.message +def test_multi_target_logs_under_collector_directory(system_info: SystemInfo, tmp_path) -> None: + """Collector logs for each target are written under the collector directory.""" + good = MagicMock() + good.marker = "good" + other = MagicMock() + other.marker = "other" + multi = MultiTargetRedfishConnection({"node-a": good, "node-b": other}) + collector_dir = tmp_path / "sample_plugin" / "sample_collector" + collector_dir.mkdir(parents=True) + collector = SampleCollector( + system_info=system_info, + connection=cast(RedfishConnection, multi), + logger=logging.getLogger("multi-target-test"), + parent="SamplePlugin", + log_path=str(collector_dir), + task_result_hooks=[FileSystemLogHook(log_base_path=str(tmp_path))], + ) + + collector.collect_data() + + assert (collector_dir / "node-a" / "result.json").is_file() + assert (collector_dir / "node-b" / "result.json").is_file() + assert not (tmp_path / "sample_plugin" / "node-a").exists() + + +def test_multi_target_analyzer_logs_under_analyzer_directory( + system_info: SystemInfo, tmp_path +) -> None: + """Analyzer logs for each target are written under the analyzer directory.""" + plugin = SamplePlugin( + system_info=system_info, + log_path=str(tmp_path), + logger=logging.getLogger("multi-target-test"), + ) + multi = MultiTargetRedfishConnection({"node-a": MagicMock()}) + multi.multi_target_data["node-a"] = SampleModel(label="good") + plugin.connection_manager = MagicMock(connection=multi) + + plugin.analyze() + + assert (tmp_path / "sample_plugin" / "sample_analyzer" / "node-a" / "result.json").is_file() + assert not (tmp_path / "sample_plugin" / "node-a").exists() + + def test_analyze_reads_data_before_disconnect(system_info: SystemInfo) -> None: """Analysis uses data still held on the open multi-target connection.""" multi = MultiTargetRedfishConnection({"good": MagicMock()}) diff --git a/test/unit/connection/redfish/test_redfish_oem_diag.py b/test/unit/connection/redfish/test_redfish_oem_diag.py index 1a78c974..4a92d9a5 100644 --- a/test/unit/connection/redfish/test_redfish_oem_diag.py +++ b/test/unit/connection/redfish/test_redfish_oem_diag.py @@ -268,6 +268,29 @@ def test_collect_polls_task_resource_until_completed(): assert polled[1] == "redfish/v1/TaskService/Tasks/dummy-1" +def test_collect_ssh_failure_during_poll_returns_error(): + conn = MagicMock() + conn.base_url = "https://bmc.example.test" + post_resp = MagicMock() + post_resp.status_code = codes.accepted + post_resp.headers = {"Location": "/redfish/v1/TaskService/TaskMonitors/1"} + post_resp.text = "" + post_resp.json.return_value = { + "@odata.id": "/redfish/v1/TaskService/Tasks/dummy-1", + "TaskState": "Running", + } + conn.post.return_value = post_resp + conn.get_response.side_effect = RedfishConnectionError("Timeout opening channel.") + + _log_bytes, _metadata, err = collect_oem_diagnostic_data( + conn, + "redfish/v1/Systems/dummy-system/LogServices/DiagLogs", + oem_diagnostic_type="AllLogs", + task_timeout_s=30, + ) + assert err == "Timeout opening channel." + + def test_collect_taskmonitor_404_does_not_spin(): conn = MagicMock() conn.base_url = "https://bmc.example.test" diff --git a/test/unit/connection/redfish/test_ssh_proxy_connection.py b/test/unit/connection/redfish/test_ssh_proxy_connection.py index 6ad329fd..9c5b9289 100644 --- a/test/unit/connection/redfish/test_ssh_proxy_connection.py +++ b/test/unit/connection/redfish/test_ssh_proxy_connection.py @@ -24,12 +24,16 @@ # ############################################################################### import shlex +from typing import Optional from unittest.mock import MagicMock, patch import pytest from pydantic import ValidationError from nodescraper.connection.inband.inband import BaseFileArtifact, CommandArtifact +from nodescraper.connection.inband.sshparams import SSHConnectionParams +from nodescraper.connection.redfish.redfish_connection import RedfishConnectionError +from nodescraper.connection.redfish.redfish_manager import MultiTargetRedfishConnection from nodescraper.connection.redfish.ssh_proxy_connection import ( CurlResponse, SshProxyRedfishConnection, @@ -155,3 +159,114 @@ def test_ssh_proxy_manager_connect_success(system_info): def test_ssh_proxy_params_require_ssh(): with pytest.raises(ValidationError): RedfishSshProxyConnectionParams(host="192.0.2.10") + + +def _proxy_target( + target_key: Optional[str], amc_host: str, ssh_host: str +) -> RedfishSshProxyConnectionParams: + return RedfishSshProxyConnectionParams( + target_key=target_key, + host=amc_host, + ssh=SSHConnectionParams( + hostname=ssh_host, username="testuser", key_filename="/tmp/dummy_id" + ), + ) + + +def test_ssh_proxy_partial_connect_keeps_reachable_bmc(system_info): + """One unreachable BMC leaves the reachable BMC connected at warning.""" + shells: list[MagicMock] = [] + + def factory(_params): + shell = MagicMock() + shells.append(shell) + return shell + + def service_root(conn): + if "amc-down" in conn.base_url: + raise RedfishConnectionError("unreachable") + return {"RedfishVersion": "1.0"} + + params = RedfishSshProxyConnectionParams( + max_workers=2, + targets=[ + _proxy_target("good", "192.0.2.10", "bmc-good.example"), + _proxy_target("bad", "amc-down.example", "bmc-down.example"), + ], + ) + mgr = RedfishSshProxyConnectionManager(system_info=system_info, connection_args=params) + with ( + patch( + "nodescraper.connection.redfish.ssh_proxy_manager.RemoteShell", + side_effect=factory, + ), + patch.object(SshProxyRedfishConnection, "get_service_root", service_root), + ): + result = mgr.connect() + + assert result.status == ExecutionStatus.WARNING + assert isinstance(mgr.connection, MultiTargetRedfishConnection) + assert set(mgr.connection.target_connections) == {"good"} + assert "bad" in mgr.connection.failed_targets + assert mgr.connection.max_workers == 2 + assert shells[1].client.close.called + assert not shells[0].client.close.called + + mgr.connection.multi_target_data["good"] = {"label": "ok"} + mgr.disconnect() + assert mgr.connection is None + assert mgr._multi_target_data == {"good": {"label": "ok"}} + assert shells[0].client.close.called + + +def test_ssh_proxy_every_target_down_fails_connect(system_info): + """Every target failing is an execution failure and does not set a connection.""" + + def service_root(_conn): + raise RedfishConnectionError("unreachable") + + params = RedfishSshProxyConnectionParams( + targets=[ + _proxy_target("a", "amc-down-a.example", "bmc-a.example"), + _proxy_target("b", "amc-down-b.example", "bmc-b.example"), + ], + ) + mgr = RedfishSshProxyConnectionManager(system_info=system_info, connection_args=params) + with ( + patch( + "nodescraper.connection.redfish.ssh_proxy_manager.RemoteShell", + return_value=MagicMock(), + ), + patch.object(SshProxyRedfishConnection, "get_service_root", service_root), + ): + result = mgr.connect() + + assert result.status == ExecutionStatus.EXECUTION_FAILURE + assert mgr.connection is None + + +def test_ssh_proxy_target_key_defaults_to_ssh_hostname(system_info): + """Targets that share an AMC address stay distinct by SSH hostname.""" + + def service_root(_conn): + return {"RedfishVersion": "1.0"} + + params = RedfishSshProxyConnectionParams( + targets=[ + _proxy_target(None, "192.0.2.10", "bmc-one.example"), + _proxy_target(None, "192.0.2.10", "bmc-two.example"), + ], + ) + mgr = RedfishSshProxyConnectionManager(system_info=system_info, connection_args=params) + with ( + patch( + "nodescraper.connection.redfish.ssh_proxy_manager.RemoteShell", + return_value=MagicMock(), + ), + patch.object(SshProxyRedfishConnection, "get_service_root", service_root), + ): + result = mgr.connect() + + assert result.status in (ExecutionStatus.UNSET, ExecutionStatus.OK) + assert isinstance(mgr.connection, MultiTargetRedfishConnection) + assert set(mgr.connection.target_connections) == {"bmc-one.example", "bmc-two.example"} diff --git a/test/unit/plugin/test_amc_redfish_diag.py b/test/unit/plugin/test_amc_redfish_diag.py index c7d5a48e..5b5f019a 100644 --- a/test/unit/plugin/test_amc_redfish_diag.py +++ b/test/unit/plugin/test_amc_redfish_diag.py @@ -39,6 +39,7 @@ AmcRedfishDiagCollectorArgs, AmcRedfishDiagPlugin, ) +from nodescraper.taskresulthooks.filesystemloghook import FileSystemLogHook @pytest.fixture @@ -166,7 +167,7 @@ def test_amc_collector_output_dir_is_diag_logs( ) output_dir = mock_collect.call_args.kwargs["output_dir"] assert output_dir == (tmp_path / "diag_logs").resolve() - assert output_dir.is_dir() + assert not output_dir.exists() def _both_collections(): @@ -197,7 +198,8 @@ def side_effect(conn, log_service_path, diagnostic_data_type, **kwargs): assert data.results["Systems:OEM:AllLogs"].success is True warnings = [e for e in result.events if e.priority == EventPriority.WARNING] assert len(warnings) == 1 - assert "Managers:Manager" in warnings[0].description + assert warnings[0].description == "AMC diag Managers:Manager: error" + assert warnings[0].data["response_body"] == "LogEntry GET failed" @patch("nodescraper.plugins.ooband.amc_redfish_diag.amc_diag_collector.collect_oem_diagnostic_data") @@ -215,3 +217,74 @@ def test_amc_collector_all_failed_is_error_with_event_per_failure(mock_collect, assert all(not r.success for r in data.results.values()) warnings = [e for e in result.events if e.priority == EventPriority.WARNING] assert len(warnings) == 2 + + +@patch("nodescraper.plugins.ooband.amc_redfish_diag.amc_diag_collector.collect_oem_diagnostic_data") +def test_amc_collector_redfish_body_stays_on_event_not_description(mock_collect, amc_collector): + body = '{"error":{"code":"Base.1.21.ResourceInUse",' '"message":"The resource is in use."}}' + full = f"Unexpected status 503 for CollectDiagnosticData: {body}" + mock_collect.return_value = (None, None, full) + amc_collector.connection.run_get.side_effect = _get_side_effect + result, data = amc_collector.collect_data( + args=AmcRedfishDiagCollectorArgs( + manager_ids=["dummy-amc"], + collections=[ + AmcDiagCollectionSpec(root="Managers", diagnostic_data_type="Manager"), + ], + ) + ) + warning = result.events[0] + assert warning.description == "AMC diag Managers:Manager: error" + assert warning.data["response_body"] == full + assert data.results["Managers:Manager"].error == full + assert "{" not in result.message + + +@patch("nodescraper.plugins.ooband.amc_redfish_diag.amc_diag_collector.collect_oem_diagnostic_data") +def test_amc_collector_logs_event_for_ssh_failure(mock_collect, amc_collector): + mock_collect.return_value = (None, None, "Timeout opening channel.") + amc_collector.connection.run_get.side_effect = _get_side_effect + result, data = amc_collector.collect_data( + args=AmcRedfishDiagCollectorArgs( + manager_ids=["dummy-amc"], + collections=[ + AmcDiagCollectionSpec(root="Managers", diagnostic_data_type="Manager"), + ], + ) + ) + warning = result.events[0] + assert warning.description == "AMC diag Managers:Manager: error" + assert warning.data["response_body"] == "Timeout opening channel." + assert data.results["Managers:Manager"].error == "Timeout opening channel." + + +@patch("nodescraper.plugins.ooband.amc_redfish_diag.amc_diag_collector.collect_oem_diagnostic_data") +def test_failed_collection_skips_diag_logs_and_writes_events( + mock_collect, system_info, redfish_conn_mock, tmp_path +): + full = ( + "Unexpected status 503 for CollectDiagnosticData: " + '{"error":{"code":"Base.1.21.ResourceInUse","message":"The resource is in use."}}' + ) + mock_collect.return_value = (None, None, full) + redfish_conn_mock.api_root = "redfish/v1" + redfish_conn_mock.run_get.side_effect = _get_side_effect + collector = AmcRedfishDiagCollector( + system_info=system_info, + connection=redfish_conn_mock, + log_path=str(tmp_path), + task_result_hooks=[FileSystemLogHook(log_base_path=str(tmp_path), include_task_path=False)], + ) + result, _data = collector.collect_data( + args=AmcRedfishDiagCollectorArgs( + manager_ids=["dummy-amc"], + collections=[ + AmcDiagCollectionSpec(root="Managers", diagnostic_data_type="Manager"), + ], + ) + ) + assert result.status == ExecutionStatus.ERROR + assert not (tmp_path / "diag_logs").exists() + event_log = (tmp_path / "events.json").read_text(encoding="utf-8") + assert "ResourceInUse" in event_log + assert "The resource is in use." in event_log diff --git a/test/unit/plugin/test_redfish_oem_diag_collector.py b/test/unit/plugin/test_redfish_oem_diag_collector.py index a42fc43d..d41f09de 100644 --- a/test/unit/plugin/test_redfish_oem_diag_collector.py +++ b/test/unit/plugin/test_redfish_oem_diag_collector.py @@ -139,4 +139,4 @@ def test_redfish_oem_diag_collector_output_dir_is_diag_logs( collector.collect_data(args=RedfishOemDiagCollectorArgs(oem_diagnostic_types=["AllLogs"])) output_dir = mock_collect.call_args.kwargs["output_dir"] assert output_dir == (tmp_path / "diag_logs").resolve() - assert output_dir.is_dir() + assert not output_dir.exists() From 65b56ca68d6cb049b40eafbc568487eb1fe51261 Mon Sep 17 00:00:00 2001 From: Anson Yim Date: Thu, 1 Oct 2026 17:57:26 -0400 Subject: [PATCH 7/7] clean up logging, add oob target prefixes to stdout --- nodescraper/base/redfishcollectortask.py | 13 ++++++++++- .../connection/redfish/redfish_manager.py | 22 ++++++++++++------- .../connection/redfish/ssh_proxy_manager.py | 11 +++++++++- 3 files changed, 36 insertions(+), 10 deletions(-) diff --git a/nodescraper/base/redfishcollectortask.py b/nodescraper/base/redfishcollectortask.py index b24e0874..1e5e7d4d 100644 --- a/nodescraper/base/redfishcollectortask.py +++ b/nodescraper/base/redfishcollectortask.py @@ -47,6 +47,17 @@ _TARGET_SUCCESS = {ExecutionStatus.OK, ExecutionStatus.WARNING} +class _TargetPrefixAdapter(logging.LoggerAdapter): + """Prepend ``[target_key]`` to every log message emitted by a per-target collector thread.""" + + def __init__(self, logger: logging.Logger, extra: dict[str, Any]) -> None: + super().__init__(logger, extra) + self._target_key: str = extra["target_key"] + + def process(self, msg: str, kwargs: Any) -> tuple[str, Any]: + return f"[{self._target_key}] {msg}", kwargs + + def _target_dir_name(target_key: str) -> str: """Return a filesystem-safe directory name for a target key. @@ -187,7 +198,7 @@ def _run_for_target( target_collector = type(collector)( system_info=collector.system_info.model_copy(), connection=conn, - logger=collector.logger, + logger=_TargetPrefixAdapter(collector.logger, {"target_key": target_key}), # type: ignore[arg-type] max_event_priority_level=getattr( collector, "max_event_priority_level", EventPriority.CRITICAL ), diff --git a/nodescraper/connection/redfish/redfish_manager.py b/nodescraper/connection/redfish/redfish_manager.py index fa278c3a..55d0ff06 100644 --- a/nodescraper/connection/redfish/redfish_manager.py +++ b/nodescraper/connection/redfish/redfish_manager.py @@ -164,14 +164,20 @@ def connect(self) -> TaskResult: self.target_connections[key] = conn if not self.target_connections: self.result.status = ExecutionStatus.EXECUTION_FAILURE - self.result.message = "Redfish connection failed for every target" - else: - # Set self.connection so DataPlugin.collect() does not short-circuit. - self.connection = MultiTargetRedfishConnection( # type: ignore[assignment] - self.target_connections, - max_workers=params.max_workers, - failed_targets=failed_targets, - ) + lines = "\n".join(f"[{k}] {v}" for k, v in failed_targets.items()) + self.result.message = f"Redfish connection failed for every target\n{lines}" + self.result.events.clear() + return self.result + if failed_targets: + lines = "\n".join(f"[{k}] {v}" for k, v in failed_targets.items()) + self.result.status = ExecutionStatus.WARNING + self.result.message = f"Some Redfish targets failed to connect\n{lines}" + self.result.events.clear() + self.connection = MultiTargetRedfishConnection( # type: ignore[assignment] + self.target_connections, + max_workers=params.max_workers, + failed_targets=failed_targets, + ) return self.result return self._connect_single(params) diff --git a/nodescraper/connection/redfish/ssh_proxy_manager.py b/nodescraper/connection/redfish/ssh_proxy_manager.py index 78a06c20..e126fdaf 100644 --- a/nodescraper/connection/redfish/ssh_proxy_manager.py +++ b/nodescraper/connection/redfish/ssh_proxy_manager.py @@ -183,8 +183,17 @@ def _connect_multi(self, params: RedfishSshProxyConnectionParams) -> TaskResult: self.target_connections[key] = conn if not self.target_connections: self.result.status = ExecutionStatus.EXECUTION_FAILURE - self.result.message = "SSH-proxy Redfish connection failed for every target" + lines = "\n".join(f"[{k}] {v}" for k, v in failed_targets.items()) + self.result.message = f"SSH-proxy Redfish connection failed for every target\n{lines}" + self.result.events.clear() return self.result + if failed_targets: + lines = "\n".join(f"[{k}] {v}" for k, v in failed_targets.items()) + self.result.status = ExecutionStatus.WARNING + self.result.message = f"Some SSH-proxy Redfish targets failed to connect\n{lines}" + # Clear events so TaskResult.finalize() does not re-append them as + # "(N warnings: ...)" on the same line — the info is already in result.message. + self.result.events.clear() self.connection = MultiTargetRedfishConnection( # type: ignore[assignment] self.target_connections, max_workers=params.max_workers,