"""Remote-agent event protocol and conversion into the common observation model."""
from datetime import UTC, datetime
from typing import Any
from .models import (
AgentRegistration,
InstanceDescriptor,
JobDescriptor,
JobObservation,
RemoteEvent,
RunDescriptor,
)
from .normalizer import observation_from_signals
from .timeutil import parse_iso
def _optional_positive(value: Any) -> int | None:
try:
number = int(value)
except (TypeError, ValueError):
return None
return number if number > 0 else None
def observation_from_remote_event(event: RemoteEvent) -> JobObservation:
payload = dict(event.payload)
job = JobDescriptor(
event.job_key,
str(payload.get("display_name", event.job_key)),
"remote_agent",
str(payload.get("job_type", "worker")),
description=payload.get("description"),
expected_interval_seconds=_optional_positive(payload.get("expected_interval_seconds")),
timeout_seconds=_optional_positive(payload.get("timeout_seconds")),
stale_after_seconds=_optional_positive(payload.get("stale_after_seconds", 300)),
metadata={"remote": True},
)
instance = InstanceDescriptor(
job,
event.agent.host,
event.agent.runtime,
event.worker_key,
instance_name=str(payload.get("instance_name", event.worker_key)),
worker_key=event.worker_key,
agent_key=event.agent.agent_key,
location_type=event.agent.environment,
location_name=event.agent.host,
metadata={"agent_key": event.agent.agent_key, "agent_version": event.agent.agent_version},
)
signals: dict[str, Any] = {"source": "remote_agent", "agent_key": event.agent.agent_key,
"event_id": event.event_id, "event_type": event.event_type}
signals.update(payload)
signals["heartbeat_at"] = event.occurred_at
if event.event_type in {"heartbeat", "run_started", "run_progress"}:
signals.setdefault("reported_state", "running")
signals.setdefault("process_exists", True)
if event.event_type == "run_finished":
signals["exit_code"] = payload.get("exit_code", 0)
signals["reported_state"] = "completed" if int(signals["exit_code"] or 0) == 0 else "stopped"
if event.event_type == "paused":
signals["paused"] = True
run: RunDescriptor | None = None
run_id = payload.get("run_id")
if run_id or event.event_type in {"run_started", "run_progress", "run_finished"}:
started_at = parse_iso(payload.get("started_at")) or (event.occurred_at if event.event_type == "run_started" else None)
finished_at = parse_iso(payload.get("finished_at")) or (event.occurred_at if event.event_type == "run_finished" else None)
run = RunDescriptor(str(run_id or f"{event.worker_key}:{event.event_id}"), started_at, finished_at,
signals.get("exit_code"), payload.get("runtime_seconds"), payload.get("records_processed"),
payload.get("error_message"))
return observation_from_signals(job, instance, run, signals, event.occurred_at)
def parse_agent_registration(payload: dict[str, Any]) -> AgentRegistration:
return AgentRegistration(
agent_key=str(payload["agent_key"]),
agent_name=str(payload.get("agent_name", payload["agent_key"])),
agent_version=str(payload.get("agent_version", "unknown")),
environment=str(payload.get("environment", "external")),
host=str(payload.get("host", "unknown")),
runtime=str(payload.get("runtime", "unknown")),
capabilities=tuple(str(value) for value in payload.get("capabilities", ())),
metadata=dict(payload.get("metadata", {})),
)
def parse_remote_event(payload: dict[str, Any], registration: AgentRegistration | None = None) -> RemoteEvent:
agent = registration or parse_agent_registration(dict(payload["agent"]))
return RemoteEvent(
event_id=str(payload["event_id"]),
agent=agent,
worker_key=str(payload["worker_key"]),
job_key=str(payload["job_key"]),
event_type=str(payload["event_type"]),
occurred_at=parse_iso(str(payload["occurred_at"])) or datetime.now(UTC),
payload=dict(payload.get("payload", {})),
sequence=payload.get("sequence"),
)