Explorer
/opt/struktur/worker-watcher/src/worker_watcher/collectors/systemd.py
← Zurück ↓ Download
from typing import Any

from ..models import CollectorError, InstanceDescriptor, JobDescriptor
from ..normalizer import observation_from_signals
from .base import CollectorContext, CommandRunner, SubprocessCommandRunner, result


class SystemdCollector:
    name = "systemd"
    version = "0.1"

    def __init__(self, runner: CommandRunner | None = None) -> None:
        self.runner = runner or SubprocessCommandRunner()

    def collect(self, context: CollectorContext):  # type: ignore[no-untyped-def]
        started = context.now
        observations = []
        errors: list[CollectorError] = []
        properties = "LoadState,ActiveState,SubState,Result,ExecMainStatus,ActiveEnterTimestamp,InactiveEnterTimestamp"
        for definition in context.config.definitions:
            service = str(definition["service_name"])
            job = JobDescriptor(str(definition["job_key"]), str(definition.get("display_name", service)), "systemd",
                                str(definition.get("job_type", "service")), stale_after_seconds=definition.get("stale_after_seconds"),
                                timeout_seconds=context.config.timeout_seconds)
            instance = InstanceDescriptor(job, context.host, "systemd", service, instance_name=definition.get("instance_name"), service_name=service)
            signals: dict[str, Any] = {"source": "systemd", "service_name": service}
            try:
                code, stdout, stderr = self.runner.run(["systemctl", "show", "--no-pager", f"--property={properties}", service], context.config.timeout_seconds)
                if code != 0:
                    errors.append(CollectorError("systemctl_error", stderr.strip() or f"systemctl exited {code}", instance.instance_key))
                    signals["unavailable"] = True
                else:
                    values = dict(line.split("=", 1) for line in stdout.splitlines() if "=" in line)
                    signals.update({f"systemd_{key}": value for key, value in values.items()})
                    signals["reported_health"] = "healthy" if values.get("ActiveState") == "active" and values.get("Result", "success") == "success" else "degraded"
                    signals["reported_state"] = "running" if values.get("ActiveState") == "active" else "stopped"
                    if values.get("ExecMainStatus") and values["ExecMainStatus"] != "0":
                        signals["exit_code"] = int(values["ExecMainStatus"])
            except PermissionError as exc:
                errors.append(CollectorError("permission_denied", str(exc), instance.instance_key))
                signals["permission_denied"] = True
            except Exception as exc:  # noqa: BLE001 - adapter failures are isolated per source
                errors.append(CollectorError(type(exc).__name__, str(exc), instance.instance_key))
                signals["unavailable"] = True
            observations.append(observation_from_signals(job, instance, None, signals, started))
        return result(self.name, self.version, started, observations, errors, {"read_only": True, "commands": ["systemctl show"]})