Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions src/powerapi/cli/common_cli_parsing_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -460,6 +460,12 @@ def _register_k8s_pre_processor_parser(self):
help_text='Kubernetes API host for manual API mode',
)

subparser_k8s_pre_processor.add_argument(
'l', 'labels',
help_text='Comma-separated list of Kubernetes pod labels added to reports as metadata',
argument_type=list
)

self.add_subgroup_parser('pre-processor', subparser_k8s_pre_processor)

def _register_openstack_pre_processor_parser(self):
Expand All @@ -475,4 +481,10 @@ def _register_openstack_pre_processor_parser(self):
default_value=10.0
)

subparser_openstack_pre_processor.add_argument(
'm', 'metadata',
help_text='Comma-separated list of OpenStack server metadata fields added to reports',
argument_type=list
)

self.add_subgroup_parser('pre-processor', subparser_openstack_pre_processor)
15 changes: 9 additions & 6 deletions src/powerapi/cli/generator.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
from powerapi.puller import PullerActor
from powerapi.pusher import PusherActor
from powerapi.report import HWPCReport, PowerReport, Report, FormulaReport
from powerapi.utils.metadata import build_metadata_mapping

COMPONENT_TYPE_KEY = 'type'
COMPONENT_MODEL_KEY = 'model'
Expand Down Expand Up @@ -435,30 +436,32 @@ def _k8s_pre_processor_factory(processor_config: dict) -> ProcessorActor:
:param processor_config: Pre-Processor configuration
:return: Configured Kubernetes pre-processor actor
"""
from powerapi.processor.pre.k8s.actor import K8sPreProcessorActor
from powerapi.processor.pre.k8s.monitor_agent import K8sMonitorConfig
from powerapi.processor.pre.k8s.actor import KubernetesPreProcessorActor
from powerapi.processor.pre.k8s.monitor_agent import KubernetesMonitorConfig

api_mode = processor_config[K8S_API_MODE_KEY]
api_host = processor_config.get(K8S_API_HOST_KEY, None)
api_key = processor_config.get(K8S_API_KEY_KEY, None)
monitor_config = K8sMonitorConfig(api_mode, api_host, api_key)
label_mapping = build_metadata_mapping(processor_config.get('labels', []), prefix='k8s_pod_label_')
monitor_config = KubernetesMonitorConfig(api_mode, api_host, api_key, label_mapping)

name = processor_config[ACTOR_NAME_KEY]
level_logger = logging.DEBUG if processor_config[GENERAL_CONF_VERBOSE_KEY] else logging.INFO
return K8sPreProcessorActor(name, monitor_config, level_logger)
return KubernetesPreProcessorActor(name, monitor_config, level_logger)

@staticmethod
def _openstack_pre_processor_factory(processor_config: dict) -> ProcessorActor:
"""
Openstack pre-processor actor factory.
OpenStack pre-processor actor factory.
:param processor_config: Pre-Processor configuration
:return: Configured OpenStack pre-processor actor
"""
from powerapi.processor.pre.openstack.actor import OpenStackPreProcessorActor
from powerapi.processor.pre.openstack.monitor_agent import OpenStackMonitorConfig

api_polling_interval = processor_config['polling-interval']
monitor_config = OpenStackMonitorConfig(api_polling_interval)
metadata_mapping = build_metadata_mapping(processor_config.get('metadata', []), prefix='openstack_metadata_')
monitor_config = OpenStackMonitorConfig(api_polling_interval, metadata_mapping)

name = processor_config[ACTOR_NAME_KEY]
level_logger = logging.DEBUG if processor_config[GENERAL_CONF_VERBOSE_KEY] else logging.INFO
Expand Down
37 changes: 20 additions & 17 deletions src/powerapi/processor/pre/k8s/actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,54 +31,57 @@
from multiprocessing import Manager

from powerapi.actor import Actor, State
from powerapi.actor.message import StartMessage, PoisonPillMessage
from powerapi.actor.message import PoisonPillMessage, StartMessage
from powerapi.processor.processor_actor import ProcessorActor
from powerapi.report import HWPCReport
from .handlers import K8sPreProcessorActorHWPCReportHandler
from .handlers import K8sPreProcessorActorStartMessageHandler, K8sPreProcessorActorPoisonPillMessageHandler
from .metadata_cache_manager import K8sMetadataCacheManager
from .monitor_agent import K8sMonitorAgent, K8sMonitorConfig

from .handlers import (
ActorPoisonPillMessageHandler,
ActorStartMessageHandler,
HWPCReportHandler,
)
from .metadata_registry import KubernetesMetadataRegistry
from .monitor_agent import KubernetesMonitorAgent, KubernetesMonitorConfig

class K8sProcessorState(State):

class KubernetesProcessorState(State):
"""
State of the Kubernetes processor actor.
"""

def __init__(self, actor: Actor, monitor_config: K8sMonitorConfig):
def __init__(self, actor: Actor, monitor_config: KubernetesMonitorConfig):
"""
Initializes a Kubernetes pre-processor state.
"""
super().__init__(actor)

self.manager = Manager()
self.metadata_cache_manager = K8sMetadataCacheManager(self.manager)
self.monitor_agent = K8sMonitorAgent(self.metadata_cache_manager, monitor_config)
self.metadata_registry = KubernetesMetadataRegistry(self.manager)
self.monitor_agent = KubernetesMonitorAgent(self.metadata_registry, monitor_config)


class K8sPreProcessorActor(ProcessorActor):
class KubernetesPreProcessorActor(ProcessorActor):
"""
Pre-Processor Actor that adds Kubernetes related metadata to reports.
"""

def __init__(self, name: str, monitor_config: K8sMonitorConfig, level_logger: int = logging.WARNING, timeout: int = 5000):
def __init__(self, name: str, monitor_config: KubernetesMonitorConfig, level_logger: int = logging.WARNING):
"""
Initializes a Kubernetes pre-processor actor.
:param name: The name of the actor
:param monitor_config: Configuration of the monitoring agent
:param level_logger: logging level of the actor
:param timeout: timeout in seconds
"""
super().__init__(name, level_logger, timeout)
super().__init__(name, level_logger, 5000)

self.monitor_config = monitor_config

def setup(self):
"""
Set up the Kubernetes pre-processor actor.
"""
self.state = K8sProcessorState(self, self.monitor_config)
self.state = KubernetesProcessorState(self, self.monitor_config)

self.add_handler(StartMessage, K8sPreProcessorActorStartMessageHandler(self.state))
self.add_handler(HWPCReport, K8sPreProcessorActorHWPCReportHandler(self.state))
self.add_handler(PoisonPillMessage, K8sPreProcessorActorPoisonPillMessageHandler(self.state))
self.add_handler(StartMessage, ActorStartMessageHandler(self.state))
self.add_handler(PoisonPillMessage, ActorPoisonPillMessageHandler(self.state))
self.add_handler(HWPCReport, HWPCReportHandler(self.state))
20 changes: 11 additions & 9 deletions src/powerapi/processor/pre/k8s/handlers.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,17 @@
# OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
# OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.

from powerapi.handler import StartHandler, PoisonPillMessageHandler
from powerapi.handler import PoisonPillMessageHandler, StartHandler
from powerapi.processor.handlers import ProcessorReportHandler
from powerapi.report import HWPCReport
from ._utils import extract_container_id_from_k8s_cgroups_path, is_target_a_valid_k8s_cgroups_path

from ._utils import (
extract_container_id_from_k8s_cgroups_path,
is_target_a_valid_k8s_cgroups_path,
)

class K8sPreProcessorActorStartMessageHandler(StartHandler):

class ActorStartMessageHandler(StartHandler):
"""
Start message handler for the Kubernetes processor actor.
"""
Expand All @@ -48,7 +52,7 @@ def initialization(self):
self.state.monitor_agent.start()


class K8sPreProcessorActorPoisonPillMessageHandler(PoisonPillMessageHandler):
class ActorPoisonPillMessageHandler(PoisonPillMessageHandler):
"""
Poison Pill message handler for the Kubernetes processor actor.
"""
Expand All @@ -66,7 +70,7 @@ def teardown(self, soft: bool = False):
actor.disconnect()


class K8sPreProcessorActorHWPCReportHandler(ProcessorReportHandler):
class HWPCReportHandler(ProcessorReportHandler):
"""
HWPCReport message handler for the Kubernetes processor actor.
"""
Expand All @@ -78,14 +82,12 @@ def handle(self, msg: HWPCReport):
"""
if is_target_a_valid_k8s_cgroups_path(msg.target):
container_id = extract_container_id_from_k8s_cgroups_path(msg.target)
container_metadata = self.state.metadata_cache_manager.get_container_metadata(container_id)

container_metadata = self.state.metadata_registry.get_metadata(container_id)
if container_metadata is None:
# Drop the report if the container metadata is not present in the cache.
# This is mainly to filter out the empty pause container present for every running POD.
return

msg.target = container_metadata.container_name
msg.metadata['k8s'] = vars(container_metadata)
msg.metadata.update(container_metadata)

self._send_report(msg)
Original file line number Diff line number Diff line change
Expand Up @@ -28,58 +28,32 @@
# OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.

from collections.abc import MutableMapping
from dataclasses import dataclass
from multiprocessing.managers import SyncManager

ADDED_EVENT = 'ADDED'
DELETED_EVENT = 'DELETED'
MODIFIED_EVENT = 'MODIFIED'


@dataclass
class K8sContainerMetadata:
"""
Represents a metadata cache entry for Kubernetes containers.
"""
container_id: str
container_name: str
namespace: str
pod_name: str
pod_labels: dict


class K8sMetadataCacheManager:
class KubernetesMetadataRegistry:
"""
Kubernetes container metadata cache manager.
Kubernetes metadata registry.
"""

def __init__(self, manager: SyncManager):
"""
:param manager: Manager of the shared metadata cache
:param manager: Manager of the shared metadata registry
"""
self.metadata_cache: MutableMapping[str, K8sContainerMetadata] = manager.dict()
self._container_metadata: MutableMapping[str, dict[str, str]] = manager.dict()

def update_container_metadata(self, event: str, container_metadata: K8sContainerMetadata):
def set_metadata(self, container_id: str, metadata: dict[str, str]) -> None:
"""
Updates the metadata cache according to an event.
:param event: Event of the metadata cache
:param container_metadata: Container metadata entry
Set the metadata for the given container ID.
:param container_id: Container ID
:param metadata: Metadata entry
"""
if event in {ADDED_EVENT, MODIFIED_EVENT}:
self.metadata_cache[container_metadata.container_id] = container_metadata
if event == DELETED_EVENT:
self.metadata_cache.pop(container_metadata.container_id, None)
self._container_metadata[container_id] = metadata

def get_container_metadata(self, container_id: str) -> K8sContainerMetadata | None:
def get_metadata(self, container_id: str) -> dict[str, str] | None:
"""
Get metadata for a specific container from the cache.
Get metadata for a specific container.
:param container_id: Container ID (hexadecimal string of 64 characters, short format is not supported)
:return: Container metadata entry
"""
return self.metadata_cache.get(container_id)

def clear_metadata_cache(self):
"""
Clears all container metadata entries from the cache.
:return: Metadata entry or None if not found
"""
self.metadata_cache.clear()
return self._container_metadata.get(container_id)
Loading