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))