import json
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
from ..models import CollectorError, InstanceDescriptor, JobDescriptor, RunDescriptor
from ..normalizer import observation_from_signals
from ..timeutil import parse_iso
from .base import CollectorContext, result
class GenericWorkerCollector:
name = "generic-worker"
version = "0.1"
def collect(self, context: CollectorContext): # type: ignore[no-untyped-def]
started = context.now
observations = []
errors: list[CollectorError] = []
for definition in context.config.definitions:
job = JobDescriptor(str(definition["job_key"]), str(definition.get("display_name", definition["job_key"])), "generic_worker",
str(definition.get("job_type", "worker")), stale_after_seconds=definition.get("stale_after_seconds"),
timeout_seconds=context.config.timeout_seconds)
source = str(definition.get("status_file", ""))
source_reference = source or str(definition.get("heartbeat_file", definition.get("output_file", job.job_key)))
instance = InstanceDescriptor(job, context.host, "generic_worker", source_reference,
instance_name=definition.get("instance_name"), metadata={"signals": list(definition)})
signals: dict[str, Any] = {"source": "generic_worker"}
try:
if source:
path = Path(source)
if path.stat().st_size > 1_000_000:
raise ValueError("status file exceeds configured 1 MiB limit")
payload = json.loads(path.read_text(encoding="utf-8"))
if not isinstance(payload, dict):
raise ValueError("status file must contain a JSON object")
signals.update(payload)
elif definition.get("paused"):
signals["paused"] = True
elif not any(definition.get(key) for key in ("heartbeat_file", "output_file", "log_file")):
signals["unavailable"] = True
for signal_name, signal_key in (("heartbeat_file", "heartbeat_at"), ("output_file", "output_file_exists"),
("log_file", "log_file_changed_at")):
configured_path = definition.get(signal_name)
if not configured_path:
continue
signal_path = Path(str(configured_path))
stat = signal_path.stat()
if signal_name == "heartbeat_file":
signals[signal_key] = datetime.fromtimestamp(stat.st_mtime, UTC)
elif signal_name == "output_file":
signals[signal_key] = True
else:
signals[signal_key] = datetime.fromtimestamp(stat.st_mtime, UTC)
except PermissionError as exc:
errors.append(CollectorError("permission_denied", str(exc), instance.instance_key))
signals["permission_denied"] = True
except (OSError, ValueError, json.JSONDecodeError) as exc:
errors.append(CollectorError(type(exc).__name__, str(exc), instance.instance_key))
signals["unavailable"] = True
run = None
if signals.get("run_id"):
run = RunDescriptor(str(signals["run_id"]), parse_iso(signals.get("started_at")), parse_iso(signals.get("finished_at")),
signals.get("exit_code"), signals.get("runtime_seconds"), signals.get("records_processed"), signals.get("error_message"))
observations.append(observation_from_signals(job, instance, run, signals, started))
return result(self.name, self.version, started, observations, errors, {"read_only": True, "declarative_signals": True})