Explorer
/opt/struktur/worker-watcher/src/worker_watcher/normalizer.py
← Zurück ↓ Download
from datetime import datetime, timedelta
from typing import Any

from .enums import HealthStatus, ObservationStatus, OperationalState, Severity
from .models import (
    EvidenceItem,
    InstanceDescriptor,
    JobDescriptor,
    JobObservation,
    NormalizedStatus,
    RunDescriptor,
)
from .timeutil import ensure_aware_utc, parse_iso


def normalize_status(job: JobDescriptor, signals: dict[str, Any], now: datetime) -> NormalizedStatus:
    """Apply global status rules to collector evidence; collectors only provide signals."""
    now = ensure_aware_utc(now)
    if signals.get("permission_denied") or signals.get("unavailable"):
        return NormalizedStatus(OperationalState.UNKNOWN, HealthStatus.UNKNOWN, ObservationStatus.UNAVAILABLE, Severity.WARNING,
                                "collector could not read the configured source")
    if signals.get("paused"):
        return NormalizedStatus(OperationalState.PAUSED, HealthStatus.HEALTHY, ObservationStatus.FRESH, Severity.INFO,
                                "job is explicitly paused")
    exit_code = signals.get("exit_code")
    if exit_code is not None and int(exit_code) != 0:
        return NormalizedStatus(OperationalState.STOPPED, HealthStatus.FAILED, ObservationStatus.FRESH,
                                Severity(signals.get("failure_severity", Severity.CRITICAL.value)),
                                f"run failed with exit code {exit_code}")
    if signals.get("reported_state") == OperationalState.COMPLETED.value:
        return NormalizedStatus(OperationalState.COMPLETED, HealthStatus.HEALTHY, ObservationStatus.FRESH, Severity.INFO,
                                "run completed successfully")
    if signals.get("reported_state") in {state.value for state in OperationalState}:
        operational = OperationalState(signals["reported_state"])
    elif signals.get("process_exists") is True or signals.get("systemd_active_state") == "active":
        operational = OperationalState.RUNNING
    elif signals.get("expected_next_run_at") is not None:
        operational = OperationalState.SCHEDULED
    else:
        operational = OperationalState.UNKNOWN
    activity = signals.get("heartbeat_at") or signals.get("last_activity_at") or signals.get("progress_at")
    if isinstance(activity, str):
        activity = parse_iso(activity)
    stale_after = job.stale_after_seconds
    if (stale_after and activity is not None and now - ensure_aware_utc(activity) > timedelta(seconds=stale_after)
            and (operational == OperationalState.RUNNING or signals.get("run_active"))):
        return NormalizedStatus(operational, HealthStatus.STALE, ObservationStatus.FRESH,
                                Severity(signals.get("stale_severity", Severity.CRITICAL.value)),
                                f"no expected activity for more than {stale_after} seconds")
    expected = signals.get("expected_next_run_at")
    if isinstance(expected, str):
        expected = parse_iso(expected)
    if expected is not None and now > ensure_aware_utc(expected) + timedelta(seconds=job.timeout_seconds or 0):
        return NormalizedStatus(operational, HealthStatus.STALE, ObservationStatus.FRESH,
                                Severity(signals.get("stale_severity", Severity.WARNING.value)),
                                "expected run is overdue")
    reported_health = signals.get("reported_health")
    if reported_health in {status.value for status in HealthStatus}:
        return NormalizedStatus(operational, HealthStatus(reported_health), ObservationStatus.FRESH,
                                Severity(signals.get("severity", Severity.INFO.value)), "reported source status")
    if operational == OperationalState.UNKNOWN:
        return NormalizedStatus(operational, HealthStatus.UNKNOWN, ObservationStatus.FRESH, Severity.INFO,
                                "source was readable but did not provide enough evidence")
    return NormalizedStatus(operational, HealthStatus.HEALTHY, ObservationStatus.FRESH, Severity.INFO, "source is healthy")


def observation_from_signals(job: JobDescriptor, instance: InstanceDescriptor, run: RunDescriptor | None,
                             signals: dict[str, Any], observed_at: datetime) -> JobObservation:
    for key in ("expected_next_run_at", "heartbeat_at", "last_activity_at", "progress_at"):
        if isinstance(signals.get(key), str):
            signals[key] = parse_iso(signals[key])
    status = normalize_status(job, signals, observed_at)
    evidence: list[EvidenceItem] = []
    for key, value in signals.items():
        if key in {"observed_at", "expected_next_run_at", "heartbeat_at", "last_activity_at", "progress_at"}:
            continue
        evidence.append(EvidenceItem(key, value, str(signals.get("source", instance.runtime_type)), observed_at))
    return JobObservation(job, instance, run, status.operational_state, status.health_status, status.observation_status,
                           status.severity, status.reason, observed_at, signals.get("last_activity_at"),
                           signals.get("expected_next_run_at"), signals.get("heartbeat_at"), signals.get("exit_code"),
                           tuple(evidence), dict(signals))