Explorer
/opt/struktur/worker-watcher/src/worker_watcher/collectors/generic_worker.py
← Zurück ↓ Download
import json
from datetime import UTC, datetime
from pathlib import Path
from typing import Any

from ..models import CollectorError, InstanceDescriptor, JobDescriptor, RunDescriptor
from ..normalizer import observation_from_signals
from ..timeutil import parse_iso
from .base import CollectorContext, result


class GenericWorkerCollector:
    name = "generic-worker"
    version = "0.1"

    def collect(self, context: CollectorContext):  # type: ignore[no-untyped-def]
        started = context.now
        observations = []
        errors: list[CollectorError] = []
        for definition in context.config.definitions:
            job = JobDescriptor(str(definition["job_key"]), str(definition.get("display_name", definition["job_key"])), "generic_worker",
                                str(definition.get("job_type", "worker")), stale_after_seconds=definition.get("stale_after_seconds"),
                                timeout_seconds=context.config.timeout_seconds)
            source = str(definition.get("status_file", ""))
            source_reference = source or str(definition.get("heartbeat_file", definition.get("output_file", job.job_key)))
            instance = InstanceDescriptor(job, context.host, "generic_worker", source_reference,
                                          instance_name=definition.get("instance_name"), metadata={"signals": list(definition)})
            signals: dict[str, Any] = {"source": "generic_worker"}
            try:
                if source:
                    path = Path(source)
                    if path.stat().st_size > 1_000_000:
                        raise ValueError("status file exceeds configured 1 MiB limit")
                    payload = json.loads(path.read_text(encoding="utf-8"))
                    if not isinstance(payload, dict):
                        raise ValueError("status file must contain a JSON object")
                    signals.update(payload)
                elif definition.get("paused"):
                    signals["paused"] = True
                elif not any(definition.get(key) for key in ("heartbeat_file", "output_file", "log_file")):
                    signals["unavailable"] = True
                for signal_name, signal_key in (("heartbeat_file", "heartbeat_at"), ("output_file", "output_file_exists"),
                                                 ("log_file", "log_file_changed_at")):
                    configured_path = definition.get(signal_name)
                    if not configured_path:
                        continue
                    signal_path = Path(str(configured_path))
                    stat = signal_path.stat()
                    if signal_name == "heartbeat_file":
                        signals[signal_key] = datetime.fromtimestamp(stat.st_mtime, UTC)
                    elif signal_name == "output_file":
                        signals[signal_key] = True
                    else:
                        signals[signal_key] = datetime.fromtimestamp(stat.st_mtime, UTC)
            except PermissionError as exc:
                errors.append(CollectorError("permission_denied", str(exc), instance.instance_key))
                signals["permission_denied"] = True
            except (OSError, ValueError, json.JSONDecodeError) as exc:
                errors.append(CollectorError(type(exc).__name__, str(exc), instance.instance_key))
                signals["unavailable"] = True
            run = None
            if signals.get("run_id"):
                run = RunDescriptor(str(signals["run_id"]), parse_iso(signals.get("started_at")), parse_iso(signals.get("finished_at")),
                                    signals.get("exit_code"), signals.get("runtime_seconds"), signals.get("records_processed"), signals.get("error_message"))
            observations.append(observation_from_signals(job, instance, run, signals, started))
        return result(self.name, self.version, started, observations, errors, {"read_only": True, "declarative_signals": True})