Implement ACTIVITY-WP-0026 ops_run claim queue (T01–T06).
All checks were successful
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / container-smoke (push) Successful in 2s
Build and Publish Container Image / build-and-push (push) Successful in 52s

Add durable claimable ops_runs table, emit dual-write on TaskSpec, REST
claim/lease/complete/fail API, ops status visibility, and consumer docs
aligned with ACT-ADR-005. T07 railiance rollout remains deploy-side.
This commit is contained in:
tegwick 2026-08-03 19:22:50 +02:00
parent a145cc4027
commit 15eb3a2066
15 changed files with 1294 additions and 27 deletions

View file

@ -26,6 +26,7 @@ from activity_core.db import make_engine
from activity_core.issue_sink import get_issue_sink
from activity_core.orm import ActivityDefinition as ActivityDefinitionRow
from activity_core.orm import ActivityRun, TaskInstance, TaskSpawnLog
from activity_core.ops_run_queue import create_ops_run_from_spec
from activity_core.llm_client import get_llm_client
from activity_core.models import InstructionDef
from activity_core.ops_evidence_sinks import persist_ops_inventory_evidence
@ -470,6 +471,25 @@ async def emit_tasks(payload: dict) -> list[str]:
activity_definition_id=activity_id,
)
try:
# ACTIVITY-WP-0026: claimable ops_run (primary for harness)
try:
ops_id = await create_ops_run_from_spec(
session,
spec,
approach_hint=spec_dict.get("approach_hint"),
)
if ops_id is not None:
activity.logger.info(
"emit_tasks: ops_run created id=%s key=%s:%s",
ops_id,
spec.source_id,
triggering_event_id,
)
except Exception as ops_exc:
activity.logger.warning(
"emit_tasks: ops_run insert failed — %s", ops_exc
)
ref = sink.emit(spec)
refs.append(ref.external_id)

View file

@ -41,6 +41,7 @@ from temporalio.client import Client
from activity_core.models import ActivityDefinition, CronTriggerConfig
from activity_core.ops_api import bind_ops_deps, router as ops_router
from activity_core.ops_runs_api import bind_ops_runs_deps, router as ops_runs_router
from activity_core.orm import ActivityDefinition as ActivityDefinitionRow, EventType as EventTypeRow
from activity_core.schedule_manager import delete_schedule, upsert_schedule
from activity_core.sync_service import run_sync
@ -73,6 +74,7 @@ async def lifespan(app: FastAPI): # type: ignore[type-arg]
_session_factory = async_sessionmaker(engine, expire_on_commit=False)
_temporal_client = await Client.connect(TEMPORAL_HOST, namespace=TEMPORAL_NAMESPACE)
bind_ops_deps(_get_db, _get_temporal)
bind_ops_runs_deps(_get_db)
yield
@ -82,6 +84,7 @@ async def lifespan(app: FastAPI): # type: ignore[type-arg]
app = FastAPI(title="activity-core API", lifespan=lifespan)
app.include_router(webhook_router)
app.include_router(ops_router)
app.include_router(ops_runs_router)
def _get_db() -> async_sessionmaker[AsyncSession]:

View file

@ -90,7 +90,7 @@ async def automations_status(
timezone_name: str = Query(default="Europe/Berlin", alias="timezone"),
activity_id: str | None = Query(default=None),
) -> dict[str, Any]:
return await ops_status(
report = await ops_status(
since=since,
until=until,
timezone_name=timezone_name,
@ -100,6 +100,37 @@ async def automations_status(
state_hub_url=os.environ.get("STATE_HUB_URL"),
activity_id=activity_id,
)
# ACTIVITY-WP-0026 T05: ops_run claim-queue visibility + SLA signals
try:
from datetime import timedelta
from sqlalchemy import and_, func, select
from activity_core.orm import OpsRun
from activity_core.ops_run_queue import ops_run_counts
sla_hours = float(os.environ.get("OPS_RUN_SLA_HOURS", "1") or "1")
now = datetime.now(timezone.utc)
sla_cutoff = now - timedelta(hours=max(0.1, sla_hours))
Session = _db()
async with Session() as session:
counts = await ops_run_counts(session)
stuck_stmt = select(func.count()).where(
and_(
OpsRun.state.in_(("open", "claimed")),
OpsRun.created_at < sla_cutoff,
)
)
stuck = int((await session.execute(stuck_stmt)).scalar_one() or 0)
report["ops_runs"] = {
"counts": counts,
"stuck_open_or_claimed": stuck,
"sla_hours": sla_hours,
"list_url": "/ops-runs?state=open",
}
except Exception as exc: # table may not exist pre-migration
report["ops_runs"] = {"error": str(exc), "counts": {}}
return report
@router.get("/automations/{definition_id}")

View file

@ -0,0 +1,299 @@
"""Ops run claim queue service (ACTIVITY-WP-0026 / ACT-ADR-005)."""
from __future__ import annotations
import os
import uuid
from datetime import datetime, timedelta, timezone
from typing import Any
from sqlalchemy import Select, and_, func, select, update
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession
from activity_core.orm import OpsRun
from activity_core.rules.models import TaskSpec
OPS_RUN_STATES = frozenset(
{"open", "claimed", "succeeded", "failed", "expired"}
)
def ops_run_queue_enabled() -> bool:
raw = (os.environ.get("OPS_RUN_QUEUE_ENABLED") or "true").strip().lower()
return raw in {"1", "true", "yes", "on", ""}
def default_lease_seconds() -> int:
try:
return max(30, int(os.environ.get("OPS_RUN_LEASE_SECONDS", "900")))
except ValueError:
return 900
def max_attempts() -> int:
try:
return max(1, int(os.environ.get("OPS_RUN_MAX_ATTEMPTS", "3")))
except ValueError:
return 3
def build_idempotency_key(spec: TaskSpec) -> str:
"""Stable key for one emit; prevents double open rows on redelivery."""
return (
f"{spec.activity_definition_id}:{spec.source_id}:{spec.triggering_event_id}"
)
def ops_run_to_dict(row: OpsRun) -> dict[str, Any]:
return {
"id": str(row.id),
"activity_definition_id": str(row.activity_definition_id),
"idempotency_key": row.idempotency_key,
"target_repo": row.target_repo,
"title": row.title,
"description": row.description,
"labels": list(row.labels or []),
"priority": row.priority,
"state": row.state,
"claim_owner": row.claim_owner,
"lease_until": row.lease_until.isoformat() if row.lease_until else None,
"attempt": row.attempt,
"source_type": row.source_type,
"source_id": row.source_id,
"triggering_event_id": row.triggering_event_id,
"approach_hint": row.approach_hint,
"result": dict(row.result or {}),
"created_at": row.created_at.isoformat() if row.created_at else None,
"updated_at": row.updated_at.isoformat() if row.updated_at else None,
}
async def create_ops_run_from_spec(
session: AsyncSession,
spec: TaskSpec,
*,
approach_hint: str | None = None,
) -> uuid.UUID | None:
"""Insert ops_run if queue enabled; return id or None if disabled/duplicate."""
if not ops_run_queue_enabled():
return None
if not spec.activity_definition_id:
return None
try:
def_id = uuid.UUID(str(spec.activity_definition_id))
except ValueError:
return None
key = build_idempotency_key(spec)
now = datetime.now(timezone.utc)
stmt = (
pg_insert(OpsRun)
.values(
id=uuid.uuid4(),
activity_definition_id=def_id,
idempotency_key=key,
target_repo=spec.target_repo,
title=spec.title or "(untitled)",
description=spec.description or "",
labels=list(spec.labels or []),
priority=spec.priority or "medium",
state="open",
attempt=0,
source_type=spec.source_type or "rule",
source_id=spec.source_id or "",
triggering_event_id=spec.triggering_event_id or "",
approach_hint=approach_hint,
result={},
created_at=now,
updated_at=now,
)
.on_conflict_do_nothing(index_elements=["idempotency_key"])
.returning(OpsRun.id)
)
result = await session.execute(stmt)
row_id = result.scalar_one_or_none()
return row_id
async def reopen_stale_claims(session: AsyncSession) -> int:
"""Return claimed rows with expired leases to open."""
now = datetime.now(timezone.utc)
stmt = (
update(OpsRun)
.where(
and_(
OpsRun.state == "claimed",
OpsRun.lease_until.is_not(None),
OpsRun.lease_until < now,
)
)
.values(
state="open",
claim_owner=None,
lease_until=None,
updated_at=now,
)
)
result = await session.execute(stmt)
return int(result.rowcount or 0)
async def claim_ops_runs(
session: AsyncSession,
*,
worker_id: str,
labels: list[str] | None = None,
labels_mode: str = "any",
limit: int = 1,
lease_seconds: int | None = None,
) -> list[OpsRun]:
"""Claim up to ``limit`` open ops_runs for worker_id."""
if not worker_id.strip():
raise ValueError("worker_id is required")
limit = max(1, min(limit, 20))
lease_seconds = lease_seconds or default_lease_seconds()
now = datetime.now(timezone.utc)
lease_until = now + timedelta(seconds=lease_seconds)
await reopen_stale_claims(session)
# Candidate ids with SKIP LOCKED
filters = [OpsRun.state == "open"]
# labels filter applied in Python after fetch for JSONB portability,
# but prefer SQL when possible — use jsonb containment for "all"
stmt: Select[tuple[uuid.UUID]] = (
select(OpsRun.id)
.where(*filters)
.order_by(OpsRun.created_at.asc())
.limit(limit * 5) # over-fetch if label filter drops rows
.with_for_update(skip_locked=True)
)
result = await session.execute(stmt)
candidate_ids = list(result.scalars().all())
if not candidate_ids:
return []
claimed: list[OpsRun] = []
for run_id in candidate_ids:
if len(claimed) >= limit:
break
row = await session.get(OpsRun, run_id)
if row is None or row.state != "open":
continue
if labels:
row_labels = {str(x) for x in (row.labels or [])}
want = {str(x) for x in labels}
if labels_mode == "all":
if not want.issubset(row_labels):
continue
else: # any
if not (want & row_labels):
continue
row.state = "claimed"
row.claim_owner = worker_id.strip()
row.lease_until = lease_until
row.attempt = int(row.attempt or 0) + 1
row.updated_at = now
claimed.append(row)
return claimed
async def heartbeat_ops_run(
session: AsyncSession,
run_id: uuid.UUID,
*,
worker_id: str,
lease_seconds: int | None = None,
) -> OpsRun | None:
row = await session.get(OpsRun, run_id)
if row is None:
return None
if row.state != "claimed" or row.claim_owner != worker_id:
return None
lease_seconds = lease_seconds or default_lease_seconds()
now = datetime.now(timezone.utc)
row.lease_until = now + timedelta(seconds=lease_seconds)
row.updated_at = now
return row
async def complete_ops_run(
session: AsyncSession,
run_id: uuid.UUID,
*,
worker_id: str,
result: dict[str, Any] | None = None,
) -> OpsRun | None:
row = await session.get(OpsRun, run_id)
if row is None:
return None
if row.state != "claimed" or row.claim_owner != worker_id:
return None
now = datetime.now(timezone.utc)
row.state = "succeeded"
row.lease_until = None
row.result = dict(result or {})
row.updated_at = now
return row
async def fail_ops_run(
session: AsyncSession,
run_id: uuid.UUID,
*,
worker_id: str,
error: str = "",
reopen: bool = False,
result: dict[str, Any] | None = None,
) -> OpsRun | None:
row = await session.get(OpsRun, run_id)
if row is None:
return None
if row.state != "claimed" or row.claim_owner != worker_id:
return None
now = datetime.now(timezone.utc)
payload = dict(result or {})
if error:
payload["error"] = error[:2000]
row.result = payload
row.updated_at = now
if reopen and int(row.attempt or 0) < max_attempts():
row.state = "open"
row.claim_owner = None
row.lease_until = None
else:
row.state = "failed"
row.lease_until = None
return row
async def list_ops_runs(
session: AsyncSession,
*,
state: str | None = None,
activity_definition_id: uuid.UUID | None = None,
since: datetime | None = None,
limit: int = 50,
) -> list[OpsRun]:
limit = max(1, min(limit, 200))
stmt = select(OpsRun).order_by(OpsRun.created_at.desc()).limit(limit)
if state:
stmt = stmt.where(OpsRun.state == state)
if activity_definition_id:
stmt = stmt.where(OpsRun.activity_definition_id == activity_definition_id)
if since:
stmt = stmt.where(OpsRun.created_at >= since)
result = await session.execute(stmt)
return list(result.scalars().all())
async def ops_run_counts(session: AsyncSession) -> dict[str, int]:
stmt = select(OpsRun.state, func.count()).group_by(OpsRun.state)
result = await session.execute(stmt)
counts = {s: 0 for s in OPS_RUN_STATES}
for state, n in result.all():
counts[str(state)] = int(n)
return counts

View file

@ -0,0 +1,339 @@
"""REST API for ops_run claim queue (ACTIVITY-WP-0026)."""
from __future__ import annotations
import os
import uuid
from datetime import datetime
from typing import Any, Callable
from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request
from pydantic import BaseModel, Field
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from activity_core.ops_auth import (
allow_unauth_mutations,
extract_operator_token,
extract_sso_principal,
operator_token_configured,
)
from activity_core.ops_run_queue import (
claim_ops_runs,
complete_ops_run,
default_lease_seconds,
fail_ops_run,
heartbeat_ops_run,
list_ops_runs,
ops_run_counts,
ops_run_to_dict,
reopen_stale_claims,
)
router = APIRouter(prefix="/ops-runs", tags=["ops-runs"])
_get_db: Callable[[], async_sessionmaker[AsyncSession]] | None = None
WORKER_TOKEN_ENV = "ACTIVITY_CORE_WORKER_TOKEN"
def bind_ops_runs_deps(
get_db: Callable[[], async_sessionmaker[AsyncSession]],
) -> None:
global _get_db
_get_db = get_db
def _db() -> async_sessionmaker[AsyncSession]:
if _get_db is None:
raise RuntimeError("ops_runs API not bound")
return _get_db()
def _worker_token_configured() -> bool:
return bool((os.environ.get(WORKER_TOKEN_ENV) or "").strip())
def _extract_worker_token(
*,
x_worker_token: str | None,
authorization: str | None,
) -> str | None:
if x_worker_token and x_worker_token.strip():
return x_worker_token.strip()
if authorization and authorization.lower().startswith("bearer "):
return authorization[7:].strip() or None
return None
def require_worker_or_operator(
request: Request,
*,
x_worker_token: str | None = None,
x_operator_token: str | None = None,
authorization: str | None = None,
) -> str:
"""Return principal string for worker or operator."""
worker_tok = _extract_worker_token(
x_worker_token=x_worker_token, authorization=authorization
)
expected_worker = (os.environ.get(WORKER_TOKEN_ENV) or "").strip()
if expected_worker and worker_tok and worker_tok == expected_worker:
return f"worker:{worker_tok[:8]}"
sso = extract_sso_principal(request)
if sso:
return f"sso:{sso}"
op_tok = extract_operator_token(
x_operator_token=x_operator_token, authorization=authorization
)
expected_op = (os.environ.get("ACTIVITY_CORE_OPERATOR_TOKEN") or "").strip()
if expected_op and op_tok and op_tok == expected_op:
return "operator:token"
if not _worker_token_configured() and not operator_token_configured():
if allow_unauth_mutations():
return "dev:unauth"
# Dev convenience when no tokens configured at all
return "dev:open"
raise HTTPException(status_code=401, detail="worker or operator auth required")
class ClaimBody(BaseModel):
worker_id: str = Field(..., min_length=1, max_length=256)
labels: list[str] | None = None
labels_mode: str = Field(default="any", pattern="^(any|all)$")
limit: int = Field(default=1, ge=1, le=20)
lease_seconds: int | None = Field(default=None, ge=30, le=86400)
class HeartbeatBody(BaseModel):
worker_id: str = Field(..., min_length=1, max_length=256)
lease_seconds: int | None = Field(default=None, ge=30, le=86400)
class CompleteBody(BaseModel):
worker_id: str = Field(..., min_length=1, max_length=256)
result: dict[str, Any] | None = None
class FailBody(BaseModel):
worker_id: str = Field(..., min_length=1, max_length=256)
error: str = ""
reopen: bool = False
result: dict[str, Any] | None = None
@router.get("")
async def get_ops_runs(
request: Request,
state: str | None = Query(default=None),
activity_definition_id: str | None = Query(default=None),
since: datetime | None = Query(default=None),
limit: int = Query(default=50, ge=1, le=200),
x_worker_token: str | None = Header(default=None, alias="X-Worker-Token"),
x_operator_token: str | None = Header(default=None, alias="X-Operator-Token"),
authorization: str | None = Header(default=None),
) -> dict[str, Any]:
require_worker_or_operator(
request,
x_worker_token=x_worker_token,
x_operator_token=x_operator_token,
authorization=authorization,
)
def_id = None
if activity_definition_id:
try:
def_id = uuid.UUID(activity_definition_id)
except ValueError as exc:
raise HTTPException(status_code=400, detail="invalid activity_definition_id") from exc
Session = _db()
async with Session() as session:
rows = await list_ops_runs(
session,
state=state,
activity_definition_id=def_id,
since=since,
limit=limit,
)
counts = await ops_run_counts(session)
await session.commit()
return {
"items": [ops_run_to_dict(r) for r in rows],
"counts": counts,
}
@router.get("/{run_id}")
async def get_ops_run(
run_id: uuid.UUID,
request: Request,
x_worker_token: str | None = Header(default=None, alias="X-Worker-Token"),
x_operator_token: str | None = Header(default=None, alias="X-Operator-Token"),
authorization: str | None = Header(default=None),
) -> dict[str, Any]:
require_worker_or_operator(
request,
x_worker_token=x_worker_token,
x_operator_token=x_operator_token,
authorization=authorization,
)
Session = _db()
async with Session() as session:
from activity_core.orm import OpsRun
row = await session.get(OpsRun, run_id)
if row is None:
raise HTTPException(status_code=404, detail="ops_run not found")
return ops_run_to_dict(row)
@router.post("/claim")
async def post_claim(
body: ClaimBody,
request: Request,
x_worker_token: str | None = Header(default=None, alias="X-Worker-Token"),
x_operator_token: str | None = Header(default=None, alias="X-Operator-Token"),
authorization: str | None = Header(default=None),
) -> dict[str, Any]:
require_worker_or_operator(
request,
x_worker_token=x_worker_token,
x_operator_token=x_operator_token,
authorization=authorization,
)
Session = _db()
async with Session() as session:
async with session.begin():
claimed = await claim_ops_runs(
session,
worker_id=body.worker_id,
labels=body.labels,
labels_mode=body.labels_mode,
limit=body.limit,
lease_seconds=body.lease_seconds,
)
return {
"items": [ops_run_to_dict(r) for r in claimed],
"lease_seconds": body.lease_seconds or default_lease_seconds(),
}
@router.post("/{run_id}/heartbeat")
async def post_heartbeat(
run_id: uuid.UUID,
body: HeartbeatBody,
request: Request,
x_worker_token: str | None = Header(default=None, alias="X-Worker-Token"),
x_operator_token: str | None = Header(default=None, alias="X-Operator-Token"),
authorization: str | None = Header(default=None),
) -> dict[str, Any]:
require_worker_or_operator(
request,
x_worker_token=x_worker_token,
x_operator_token=x_operator_token,
authorization=authorization,
)
Session = _db()
async with Session() as session:
async with session.begin():
row = await heartbeat_ops_run(
session,
run_id,
worker_id=body.worker_id,
lease_seconds=body.lease_seconds,
)
if row is None:
raise HTTPException(
status_code=409,
detail="not claimed by this worker or not found",
)
return ops_run_to_dict(row)
@router.post("/{run_id}/complete")
async def post_complete(
run_id: uuid.UUID,
body: CompleteBody,
request: Request,
x_worker_token: str | None = Header(default=None, alias="X-Worker-Token"),
x_operator_token: str | None = Header(default=None, alias="X-Operator-Token"),
authorization: str | None = Header(default=None),
) -> dict[str, Any]:
require_worker_or_operator(
request,
x_worker_token=x_worker_token,
x_operator_token=x_operator_token,
authorization=authorization,
)
Session = _db()
async with Session() as session:
async with session.begin():
row = await complete_ops_run(
session,
run_id,
worker_id=body.worker_id,
result=body.result,
)
if row is None:
raise HTTPException(
status_code=409,
detail="not claimed by this worker or not found",
)
return ops_run_to_dict(row)
@router.post("/{run_id}/fail")
async def post_fail(
run_id: uuid.UUID,
body: FailBody,
request: Request,
x_worker_token: str | None = Header(default=None, alias="X-Worker-Token"),
x_operator_token: str | None = Header(default=None, alias="X-Operator-Token"),
authorization: str | None = Header(default=None),
) -> dict[str, Any]:
require_worker_or_operator(
request,
x_worker_token=x_worker_token,
x_operator_token=x_operator_token,
authorization=authorization,
)
Session = _db()
async with Session() as session:
async with session.begin():
row = await fail_ops_run(
session,
run_id,
worker_id=body.worker_id,
error=body.error,
reopen=body.reopen,
result=body.result,
)
if row is None:
raise HTTPException(
status_code=409,
detail="not claimed by this worker or not found",
)
return ops_run_to_dict(row)
@router.post("/expire-leases")
async def post_expire_leases(
request: Request,
x_worker_token: str | None = Header(default=None, alias="X-Worker-Token"),
x_operator_token: str | None = Header(default=None, alias="X-Operator-Token"),
authorization: str | None = Header(default=None),
) -> dict[str, Any]:
require_worker_or_operator(
request,
x_worker_token=x_worker_token,
x_operator_token=x_operator_token,
authorization=authorization,
)
Session = _db()
async with Session() as session:
async with session.begin():
n = await reopen_stale_claims(session)
return {"reopened": n}

View file

@ -138,3 +138,45 @@ class TaskInstance(Base):
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now()
)
class OpsRun(Base):
"""Claimable automation run instance (ACT-ADR-005 / ACTIVITY-WP-0026)."""
__tablename__ = "ops_runs"
id: Mapped[uuid.UUID] = mapped_column(
UUID(as_uuid=True), primary_key=True, default=uuid.uuid4
)
activity_definition_id: Mapped[uuid.UUID] = mapped_column(
UUID(as_uuid=True),
ForeignKey("activity_definitions.id", ondelete="RESTRICT"),
nullable=False,
index=True,
)
idempotency_key: Mapped[str] = mapped_column(Text, nullable=False, unique=True)
target_repo: Mapped[str | None] = mapped_column(Text, nullable=True)
title: Mapped[str] = mapped_column(Text, nullable=False)
description: Mapped[str] = mapped_column(Text, nullable=False, default="")
labels: Mapped[list] = mapped_column(JSONB, nullable=False, default=list)
priority: Mapped[str] = mapped_column(Text, nullable=False, default="medium")
state: Mapped[str] = mapped_column(Text, nullable=False, default="open", index=True)
claim_owner: Mapped[str | None] = mapped_column(Text, nullable=True)
lease_until: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), nullable=True, index=True
)
attempt: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
source_type: Mapped[str] = mapped_column(Text, nullable=False)
source_id: Mapped[str] = mapped_column(Text, nullable=False)
triggering_event_id: Mapped[str] = mapped_column(Text, nullable=False, index=True)
approach_hint: Mapped[str | None] = mapped_column(Text, nullable=True)
result: Mapped[dict] = mapped_column(JSONB, nullable=False, default=dict)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now()
)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True),
nullable=False,
server_default=func.now(),
onupdate=func.now(),
)