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)