activity-core/src/activity_core/ops_console.py
tegwick 7711bf9c70
All checks were successful
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / container-smoke (push) Successful in 1s
Build and Publish Container Image / build-and-push (push) Successful in 20s
feat(ops): run artefacts in API/UI + llm-connect host access design
ACTIVITY-WP-0027: document durable ClusterIP host access for llm-connect;
join activity_runs to ops_runs; surface Forgejo artefact links on
/ops/ui automation detail and new run detail pages.
2026-08-05 17:23:02 +02:00

410 lines
13 KiB
Python

"""Operator console services (ACTIVITY-WP-0024) — inventory, status, runs, audit."""
from __future__ import annotations
import argparse
import logging
import os
import uuid
from collections import deque
from datetime import datetime, timezone
from typing import Any
import httpx
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from temporalio.client import Client
from activity_core.automation_status import (
DEFAULT_STATE_HUB_URL,
DEFAULT_TEMPORAL_NAMESPACE,
DEFAULT_TIMEZONE,
automation_schedule_id,
build_inventory_report,
build_report,
inventory_row,
load_temporal_visibility,
public_trigger_config,
)
from activity_core.models import ActivityDefinition, CronTriggerConfig, ScheduledTriggerConfig
from activity_core.orm import ActivityDefinition as ActivityDefinitionRow
from activity_core.orm import ActivityRun
from activity_core.runtime_paths import custodian_working_memory_dir
from activity_core.schedule_manager import pause_schedule, unpause_schedule, upsert_schedule
from activity_core.state_hub_write import idempotency_headers, parse_state_hub_write_response
logger = logging.getLogger(__name__)
_AUDIT_BUFFER: deque[dict[str, Any]] = deque(maxlen=100)
SIDE_EFFECT_MARKERS = (
"forgejo_package_prune",
"cnpg_option_a_backup",
"apply: true",
'"apply": true',
"'apply': true",
)
def _status_namespace(
*,
since: str,
until: str | None,
timezone_name: str,
db_url: str | None,
temporal_host: str | None,
temporal_namespace: str,
state_hub_url: str | None,
timeout_seconds: float = 5.0,
activity_id: list[str] | None = None,
activity_name: list[str] | None = None,
) -> argparse.Namespace:
return argparse.Namespace(
since=since,
until=until,
timezone=timezone_name,
activity_id=activity_id or [],
activity_name=activity_name or [],
db_url=db_url,
state_hub_url=state_hub_url or os.environ.get("STATE_HUB_URL", DEFAULT_STATE_HUB_URL),
working_memory_dir=os.environ.get(
"AUTOMATION_STATUS_WORKING_MEMORY_DIR",
str(custodian_working_memory_dir()),
),
temporal_host=temporal_host or os.environ.get("TEMPORAL_HOST"),
temporal_namespace=temporal_namespace
or os.environ.get("TEMPORAL_NAMESPACE", DEFAULT_TEMPORAL_NAMESPACE),
timeout_seconds=timeout_seconds,
progress_limit=int(os.environ.get("AUTOMATION_STATUS_PROGRESS_LIMIT", "100")),
progress_event_type=None,
format="json",
)
def _inventory_namespace(
*,
db_url: str | None,
temporal_host: str | None,
temporal_namespace: str,
enabled: str = "all",
trigger: list[str] | None = None,
activity_id: list[str] | None = None,
activity_name: list[str] | None = None,
timeout_seconds: float = 5.0,
) -> argparse.Namespace:
return argparse.Namespace(
db_url=db_url,
temporal_host=temporal_host or os.environ.get("TEMPORAL_HOST"),
temporal_namespace=temporal_namespace
or os.environ.get("TEMPORAL_NAMESPACE", DEFAULT_TEMPORAL_NAMESPACE),
timeout_seconds=timeout_seconds,
enabled=enabled,
trigger_type=trigger or [],
activity_id=activity_id or [],
activity_name=activity_name or [],
format="json",
)
async def ops_inventory(
*,
db_url: str | None,
temporal_host: str | None = None,
temporal_namespace: str = DEFAULT_TEMPORAL_NAMESPACE,
enabled: str = "all",
) -> dict[str, Any]:
args = _inventory_namespace(
db_url=db_url,
temporal_host=temporal_host,
temporal_namespace=temporal_namespace,
enabled=enabled,
)
report, _exit = await build_inventory_report(args)
return report
async def ops_status(
*,
since: str = "today",
until: str | None = None,
timezone_name: str = DEFAULT_TIMEZONE,
db_url: str | None = None,
temporal_host: str | None = None,
temporal_namespace: str = DEFAULT_TEMPORAL_NAMESPACE,
state_hub_url: str | None = None,
activity_id: str | None = None,
) -> dict[str, Any]:
ids = [activity_id] if activity_id else []
args = _status_namespace(
since=since,
until=until,
timezone_name=timezone_name,
db_url=db_url,
temporal_host=temporal_host,
temporal_namespace=temporal_namespace,
state_hub_url=state_hub_url,
activity_id=ids,
)
report, exit_code = await build_report(args)
report["exit_code"] = exit_code
return report
async def ops_definition_detail(
session_factory: async_sessionmaker[AsyncSession],
temporal: Client | None,
definition_id: uuid.UUID,
*,
db_url: str | None,
temporal_host: str | None = None,
temporal_namespace: str = DEFAULT_TEMPORAL_NAMESPACE,
) -> dict[str, Any]:
async with session_factory() as session:
row = await session.get(ActivityDefinitionRow, definition_id)
if row is None:
return {}
definition = {
"id": str(row.id),
"name": row.name,
"enabled": row.enabled,
"trigger_type": row.trigger_type,
"trigger_config": row.trigger_config or {},
"source": "database",
"version": row.version,
"created_at": row.created_at.isoformat() if row.created_at else None,
"updated_at": row.updated_at.isoformat() if row.updated_at else None,
"side_effect": is_side_effect_definition(row),
}
temporal_map: dict[str, dict[str, Any]] = {}
if temporal is not None or temporal_host:
host = temporal_host or os.environ.get("TEMPORAL_HOST")
temporal_map, _src = await load_temporal_visibility(
host,
temporal_namespace,
[definition],
timeout_seconds=5.0,
)
row_out = inventory_row(definition, temporal_map.get(definition["id"]))
row_out["version"] = definition.get("version")
row_out["created_at"] = definition.get("created_at")
row_out["updated_at"] = definition.get("updated_at")
row_out["side_effect"] = definition.get("side_effect")
row_out["trigger_config_public"] = public_trigger_config(definition.get("trigger_config") or {})
return row_out
def is_side_effect_definition(row: ActivityDefinitionRow) -> bool:
blob = str(row.context_sources or []) + str(row.task_templates or []) + str(row.trigger_config or {})
lower = blob.lower()
return any(marker in lower for marker in SIDE_EFFECT_MARKERS)
async def ops_runs(
session_factory: async_sessionmaker[AsyncSession],
definition_id: uuid.UUID,
*,
since: datetime | None = None,
limit: int = 50,
) -> dict[str, Any]:
"""Recent activity_runs enriched with ops_runs + artefact links (WP-0027)."""
from activity_core.run_artifacts import enrich_activity_runs
limit = max(1, min(limit, 200))
async with session_factory() as session:
stmt = (
select(ActivityRun)
.where(ActivityRun.activity_id == definition_id)
.order_by(ActivityRun.fired_at.desc())
.limit(limit)
)
if since is not None:
stmt = stmt.where(ActivityRun.fired_at >= since)
runs = list((await session.scalars(stmt)).all())
items = await enrich_activity_runs(
session, definition_id, runs, since=since
)
return {
"activity_id": str(definition_id),
"count": len(items),
"runs": items,
}
async def ops_run_detail(
session_factory: async_sessionmaker[AsyncSession],
definition_id: uuid.UUID,
run_id: uuid.UUID,
) -> dict[str, Any] | None:
from activity_core.run_artifacts import get_enriched_run
async with session_factory() as session:
return await get_enriched_run(session, definition_id, run_id)
async def record_ops_audit(
*,
action: str,
activity_id: str,
activity_name: str | None,
principal: str,
detail: dict[str, Any] | None = None,
) -> dict[str, Any]:
event = {
"event_type": "ops_console_audit",
"action": action,
"activity_id": activity_id,
"activity_name": activity_name,
"principal": principal,
"detail": detail or {},
"at": datetime.now(tz=timezone.utc).isoformat(),
"audit_id": str(uuid.uuid4()),
}
_AUDIT_BUFFER.appendleft(event)
hub = (os.environ.get("STATE_HUB_URL") or "").rstrip("/")
if hub:
try:
body = {
"event_type": "ops_console_audit",
"summary": f"ops {action} {activity_name or activity_id} by {principal}",
"author": f"activity-core-ops:{principal}",
"detail": {
"action": action,
"activity_id": activity_id,
"activity_name": activity_name,
"principal": principal,
**(detail or {}),
},
}
headers = {
"Content-Type": "application/json",
**idempotency_headers("ops_console_audit", action, activity_id, event["audit_id"]),
}
with httpx.Client(timeout=3.0) as client:
resp = client.post(f"{hub}/progress/", json=body, headers=headers)
parse_state_hub_write_response(resp)
event["hub"] = "ok"
except Exception as exc: # noqa: BLE001 — audit must not fail mutation
logger.warning("ops audit hub write failed: %s", exc)
event["hub"] = f"failed:{type(exc).__name__}"
else:
event["hub"] = "skipped"
return event
def recent_audits(limit: int = 20) -> list[dict[str, Any]]:
return list(_AUDIT_BUFFER)[: max(1, min(limit, 100))]
def clear_audit_buffer_for_tests() -> None:
_AUDIT_BUFFER.clear()
async def set_definition_enabled(
session_factory: async_sessionmaker[AsyncSession],
temporal: Client,
definition_id: uuid.UUID,
*,
enabled: bool,
principal: str,
) -> dict[str, Any]:
async with session_factory() as session:
row = await session.get(ActivityDefinitionRow, definition_id)
if row is None:
raise KeyError("not found")
row.enabled = enabled
async with session.begin():
session.add(row)
name = row.name
defn = _row_to_model(row)
schedule_result: dict[str, Any] = {}
try:
if isinstance(defn.trigger_config, (CronTriggerConfig, ScheduledTriggerConfig)):
await upsert_schedule(temporal, defn)
schedule_result = {"upserted": True, "paused": not enabled}
else:
schedule_result = {"upserted": False, "reason": "non-scheduled trigger"}
except Exception as exc: # noqa: BLE001
schedule_result = {"upserted": False, "error": str(exc)}
audit = await record_ops_audit(
action="enable" if enabled else "disable",
activity_id=str(definition_id),
activity_name=name,
principal=principal,
detail={"enabled": enabled, "schedule": schedule_result},
)
return {
"activity_id": str(definition_id),
"name": name,
"enabled": enabled,
"schedule": schedule_result,
"audit": audit,
}
async def set_schedule_paused(
session_factory: async_sessionmaker[AsyncSession],
temporal: Client,
definition_id: uuid.UUID,
*,
paused: bool,
principal: str,
) -> dict[str, Any]:
async with session_factory() as session:
row = await session.get(ActivityDefinitionRow, definition_id)
if row is None:
raise KeyError("not found")
name = row.name
trigger_type = row.trigger_type
onetime = trigger_type == "scheduled"
note = f"{'paused' if paused else 'unpaused'} by ops console ({principal})"
if paused:
schedule_result = await pause_schedule(temporal, definition_id, note=note, onetime=onetime)
else:
schedule_result = await unpause_schedule(temporal, definition_id, note=note, onetime=onetime)
audit = await record_ops_audit(
action="pause" if paused else "unpause",
activity_id=str(definition_id),
activity_name=name,
principal=principal,
detail={"schedule": schedule_result},
)
return {
"activity_id": str(definition_id),
"name": name,
"schedule": schedule_result,
"audit": audit,
}
def _row_to_model(row: ActivityDefinitionRow) -> ActivityDefinition:
return ActivityDefinition.model_validate(
{
"id": row.id,
"name": row.name,
"enabled": row.enabled,
"trigger_config": row.trigger_config,
"context_sources": row.context_sources or [],
"task_templates": row.task_templates or [],
"rules": row.rules_json or [],
"instructions": row.instructions_json or [],
"dedupe_key_strategy": row.dedupe_key_strategy,
"version": row.version,
}
)
def definition_has_side_effect(
session_row: ActivityDefinitionRow | None,
*,
context_sources: Any = None,
) -> bool:
if session_row is not None:
return is_side_effect_definition(session_row)
blob = str(context_sources or "").lower()
return any(m in blob for m in SIDE_EFFECT_MARKERS)