Explorer
/opt/struktur/worker-watcher/src/worker_watcher/models.py
← Zurück ↓ Download
import re
from dataclasses import dataclass, field
from datetime import datetime
from typing import Any

from .enums import HealthStatus, ObservationStatus, OperationalState, Severity
from .timeutil import ensure_aware_utc


def validate_job_key(job_key: str) -> str:
    parts = job_key.split(".")
    if len(parts) < 3 or any(not part or not part.isalnum() or part != part.lower() for part in parts):
        raise ValueError("job_key must be lowercase dot-separated alphanumeric components")
    return job_key


def validate_agent_key(agent_key: str) -> str:
    if not re.fullmatch(r"[a-z0-9][a-z0-9._-]{1,127}", agent_key):
        raise ValueError("agent_key must be a stable lowercase identifier")
    return agent_key


@dataclass(frozen=True, slots=True)
class JobDescriptor:
    job_key: str
    name: str
    source_type: str
    job_type: str
    description: str | None = None
    expected_interval_seconds: int | None = None
    timeout_seconds: int | None = None
    stale_after_seconds: int | None = None
    enabled: bool = True
    metadata: dict[str, Any] = field(default_factory=dict)

    def __post_init__(self) -> None:
        validate_job_key(self.job_key)
        for value in (self.expected_interval_seconds, self.timeout_seconds, self.stale_after_seconds):
            if value is not None and value <= 0:
                raise ValueError("interval and timeout values must be positive")


@dataclass(frozen=True, slots=True)
class InstanceDescriptor:
    job: JobDescriptor
    host: str
    runtime_type: str
    source_reference: str
    instance_name: str | None = None
    service_name: str | None = None
    container_name: str | None = None
    process_name: str | None = None
    worker_key: str | None = None
    agent_key: str | None = None
    location_type: str = "local"
    location_name: str | None = None
    enabled: bool = True
    metadata: dict[str, Any] = field(default_factory=dict)

    @property
    def instance_key(self) -> str:
        readable = self.instance_name or self.source_reference
        return f"{self.job.job_key}|{self.host}|{self.runtime_type}|{readable}"


@dataclass(frozen=True, slots=True)
class RunDescriptor:
    external_run_id: str | None = None
    started_at: datetime | None = None
    finished_at: datetime | None = None
    exit_code: int | None = None
    runtime_seconds: float | None = None
    records_processed: int | None = None
    error_message: str | None = None

    def __post_init__(self) -> None:
        if self.started_at is not None:
            ensure_aware_utc(self.started_at)
        if self.finished_at is not None:
            ensure_aware_utc(self.finished_at)
        if self.started_at and self.finished_at and self.finished_at < self.started_at:
            raise ValueError("finished_at cannot be before started_at")
        if self.runtime_seconds is not None and self.runtime_seconds < 0:
            raise ValueError("runtime_seconds cannot be negative")


@dataclass(frozen=True, slots=True)
class EvidenceItem:
    key: str
    value: Any
    source: str
    observed_at: datetime
    description: str | None = None

    def __post_init__(self) -> None:
        ensure_aware_utc(self.observed_at)


@dataclass(frozen=True, slots=True)
class JobObservation:
    job: JobDescriptor
    instance: InstanceDescriptor
    run: RunDescriptor | None
    operational_state: OperationalState
    health_status: HealthStatus
    observation_status: ObservationStatus
    severity: Severity
    status_reason: str
    observed_at: datetime
    last_activity_at: datetime | None = None
    expected_next_run_at: datetime | None = None
    heartbeat_at: datetime | None = None
    exit_code: int | None = None
    evidence: tuple[EvidenceItem, ...] = ()
    raw_payload: dict[str, Any] = field(default_factory=dict)

    def __post_init__(self) -> None:
        for value in (self.observed_at, self.last_activity_at, self.expected_next_run_at, self.heartbeat_at):
            if value is not None:
                ensure_aware_utc(value)


@dataclass(frozen=True, slots=True)
class CollectorError:
    error_type: str
    message: str
    instance_key: str | None = None


@dataclass(frozen=True, slots=True)
class CollectorResult:
    collector_name: str
    collector_version: str
    started_at: datetime
    finished_at: datetime
    success: bool
    observations: tuple[JobObservation, ...] = ()
    errors: tuple[CollectorError, ...] = ()
    diagnostics: dict[str, Any] = field(default_factory=dict)


@dataclass(frozen=True, slots=True)
class NormalizedStatus:
    operational_state: OperationalState
    health_status: HealthStatus
    observation_status: ObservationStatus
    severity: Severity
    reason: str


@dataclass(frozen=True, slots=True)
class WatcherCycleResult:
    started_at: datetime
    finished_at: datetime
    collectors_total: int
    collectors_successful: int
    collectors_failed: int
    observations_created: int


@dataclass(frozen=True, slots=True)
class AgentRegistration:
    agent_key: str
    agent_name: str
    agent_version: str
    environment: str
    host: str
    runtime: str
    capabilities: tuple[str, ...] = ()
    metadata: dict[str, Any] = field(default_factory=dict)

    def __post_init__(self) -> None:
        validate_agent_key(self.agent_key)
        if not self.agent_name or not self.environment or not self.host:
            raise ValueError("agent_name, environment and host are required")


@dataclass(frozen=True, slots=True)
class RemoteEvent:
    event_id: str
    agent: AgentRegistration
    worker_key: str
    job_key: str
    event_type: str
    occurred_at: datetime
    payload: dict[str, Any] = field(default_factory=dict)
    sequence: int | None = None

    def __post_init__(self) -> None:
        if not self.event_id or len(self.event_id) > 128:
            raise ValueError("event_id must be present and at most 128 characters")
        if not self.worker_key or len(self.worker_key) > 256:
            raise ValueError("worker_key must be present and at most 256 characters")
        validate_job_key(self.job_key)
        if not self.event_type or len(self.event_type) > 64:
            raise ValueError("event_type must be present and at most 64 characters")
        ensure_aware_utc(self.occurred_at)