activity-core/tests/test_ops_run_queue.py

441 lines
14 KiB
Python
Raw Permalink Normal View History

"""Tests for ops_run claim queue (ACTIVITY-WP-0026)."""
from __future__ import annotations
import os
import uuid
from datetime import datetime, timedelta, timezone
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
from activity_core.ops_run_queue import (
build_idempotency_key,
emit_triggering_event_id,
max_attempts,
ops_run_queue_enabled,
ops_run_to_dict,
)
from activity_core.rules.models import TaskSpec
def test_ops_run_queue_enabled_default(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.delenv("OPS_RUN_QUEUE_ENABLED", raising=False)
assert ops_run_queue_enabled() is True
monkeypatch.setenv("OPS_RUN_QUEUE_ENABLED", "false")
assert ops_run_queue_enabled() is False
monkeypatch.setenv("OPS_RUN_QUEUE_ENABLED", "1")
assert ops_run_queue_enabled() is True
def test_build_idempotency_key() -> None:
spec = TaskSpec(
title="t",
activity_definition_id="aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee",
source_id="emit-fi",
triggering_event_id="wf-1",
)
assert build_idempotency_key(spec) == (
"aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee:emit-fi:wf-1"
)
def test_emit_triggering_event_id_scheduled_uses_run_id() -> None:
"""Bare 'scheduled' must not be the ops_run key for every weekday fire."""
day1 = emit_triggering_event_id("scheduled", "run-day-1")
day2 = emit_triggering_event_id("scheduled", "run-day-2")
assert day1 == "run-day-1"
assert day2 == "run-day-2"
assert day1 != day2
def test_emit_triggering_event_id_prefers_scheduled_for() -> None:
got = emit_triggering_event_id(
"scheduled",
"run-fallback",
scheduled_for="2026-08-05T05:30:00+00:00",
)
assert got == "scheduled:2026-08-05T05:30:00+00:00"
def test_emit_triggering_event_id_passes_through_events_and_manual() -> None:
assert emit_triggering_event_id("evt-abc", "run-x") == "evt-abc"
assert (
emit_triggering_event_id("manual-3da4cf06", "run-y") == "manual-3da4cf06"
)
def test_scheduled_fires_produce_distinct_ops_keys() -> None:
"""Regression: two cron days must not share one ops_run idempotency key."""
def_id = "3169ab1f-882b-59c7-9763-d014dc96f4fc"
source = "emit-fi-daily-brief-task"
key_aug4 = build_idempotency_key(
TaskSpec(
title="FI",
activity_definition_id=def_id,
source_id=source,
triggering_event_id=emit_triggering_event_id("scheduled", "run-aug4"),
)
)
key_aug5 = build_idempotency_key(
TaskSpec(
title="FI",
activity_definition_id=def_id,
source_id=source,
triggering_event_id=emit_triggering_event_id("scheduled", "run-aug5"),
)
)
# Old bug: both would be "...:scheduled"
legacy = f"{def_id}:{source}:scheduled"
assert key_aug4 != key_aug5
assert key_aug4 != legacy
assert key_aug5 != legacy
def test_max_attempts(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.delenv("OPS_RUN_MAX_ATTEMPTS", raising=False)
assert max_attempts() == 3
monkeypatch.setenv("OPS_RUN_MAX_ATTEMPTS", "5")
assert max_attempts() == 5
def test_ops_run_to_dict_shape() -> None:
now = datetime.now(timezone.utc)
row = MagicMock()
row.id = uuid.uuid4()
row.activity_definition_id = uuid.uuid4()
row.idempotency_key = "k"
row.target_repo = "freedom-intelligence"
row.title = "FI brief"
row.description = "d"
row.labels = ["automated", "research-brief"]
row.priority = "medium"
row.state = "open"
row.claim_owner = None
row.lease_until = None
row.attempt = 0
row.source_type = "rule"
row.source_id = "emit"
row.triggering_event_id = "e1"
row.approach_hint = None
row.result = {}
row.created_at = now
row.updated_at = now
d = ops_run_to_dict(row)
assert d["target_repo"] == "freedom-intelligence"
assert d["state"] == "open"
assert "research-brief" in d["labels"]
assert d["created_at"]
@pytest.mark.asyncio
async def test_create_ops_run_disabled(monkeypatch: pytest.MonkeyPatch) -> None:
from activity_core.ops_run_queue import create_ops_run_from_spec
monkeypatch.setenv("OPS_RUN_QUEUE_ENABLED", "false")
session = AsyncMock()
spec = TaskSpec(
title="t",
activity_definition_id=str(uuid.uuid4()),
source_id="s",
triggering_event_id="e",
)
assert await create_ops_run_from_spec(session, spec) is None
session.execute.assert_not_called()
@pytest.mark.asyncio
async def test_create_ops_run_inserts(monkeypatch: pytest.MonkeyPatch) -> None:
from activity_core.ops_run_queue import create_ops_run_from_spec
monkeypatch.setenv("OPS_RUN_QUEUE_ENABLED", "true")
run_id = uuid.uuid4()
result = MagicMock()
result.scalar_one_or_none.return_value = run_id
session = AsyncMock()
session.execute = AsyncMock(return_value=result)
spec = TaskSpec(
title="FI daily",
target_repo="freedom-intelligence",
labels=["automated", "research-brief"],
activity_definition_id=str(uuid.uuid4()),
source_id="emit-fi-daily-brief-task",
triggering_event_id="manual-1",
)
got = await create_ops_run_from_spec(session, spec)
assert got == run_id
session.execute.assert_awaited()
@pytest.mark.asyncio
async def test_claim_and_complete_roundtrip(monkeypatch: pytest.MonkeyPatch) -> None:
"""In-memory style: claim filters labels and complete transitions state."""
from activity_core import ops_run_queue as oq
now = datetime.now(timezone.utc)
open_run = MagicMock()
open_run.id = uuid.uuid4()
open_run.state = "open"
open_run.labels = ["automated", "research-brief"]
open_run.attempt = 0
open_run.claim_owner = None
open_run.lease_until = None
open_run.result = {}
session = AsyncMock()
reopen_result = MagicMock()
reopen_result.rowcount = 0
id_result = MagicMock()
id_result.scalars.return_value.all.return_value = [open_run.id]
async def execute_side_effect(stmt, *args, **kwargs):
# reopen_stale_claims uses Update; claim select uses Select (may contain FOR UPDATE)
name = type(stmt).__name__.lower()
if name == "update":
return reopen_result
return id_result
session.execute = AsyncMock(side_effect=execute_side_effect)
session.get = AsyncMock(return_value=open_run)
claimed = await oq.claim_ops_runs(
session,
worker_id="worker-1",
labels=["research-brief"],
limit=1,
lease_seconds=60,
)
assert len(claimed) == 1
assert open_run.state == "claimed"
assert open_run.claim_owner == "worker-1"
assert open_run.attempt == 1
assert open_run.lease_until is not None
assert open_run.lease_until > now
done = await oq.complete_ops_run(
session,
open_run.id,
worker_id="worker-1",
result={"path": "briefs/x.md"},
)
assert done is not None
assert open_run.state == "succeeded"
assert open_run.result["path"] == "briefs/x.md"
@pytest.mark.asyncio
async def test_complete_persists_only_normalized_glas_evidence() -> None:
from activity_core import ops_run_queue as oq
row = MagicMock()
row.state = "claimed"
row.claim_owner = "worker-1"
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
row.result = {}
session = AsyncMock()
session.get = AsyncMock(return_value=row)
done = await oq.complete_ops_run(
session,
uuid.uuid4(),
worker_id="worker-1",
result={
"ok": True,
"evidence": {
"request_id": "request-1",
"profile_ref": "harness.agent-dev@1.0.0",
"rein_id": "rein-aharness",
"outcome": "succeeded",
"duration_s": 0,
},
"tool_output": "must not persist",
},
)
assert done is row
assert row.result == {
"ok": True,
"execution_evidence": {
"request_id": "request-1",
"profile_ref": "harness.agent-dev@1.0.0",
"rein_id": "rein-aharness",
"outcome": "succeeded",
"duration_s": 0,
},
}
@pytest.mark.asyncio
async def test_fail_reopen_under_max_attempts(monkeypatch: pytest.MonkeyPatch) -> None:
from activity_core import ops_run_queue as oq
monkeypatch.setenv("OPS_RUN_MAX_ATTEMPTS", "3")
row = MagicMock()
row.id = uuid.uuid4()
row.state = "claimed"
row.claim_owner = "w1"
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
row.attempt = 1
row.result = {}
session = AsyncMock()
session.get = AsyncMock(return_value=row)
out = await oq.fail_ops_run(
session, row.id, worker_id="w1", error="timeout", reopen=True
)
assert out is not None
assert row.state == "open"
assert row.claim_owner is None
@pytest.mark.asyncio
async def test_fail_permanent_at_max_attempts(monkeypatch: pytest.MonkeyPatch) -> None:
from activity_core import ops_run_queue as oq
monkeypatch.setenv("OPS_RUN_MAX_ATTEMPTS", "3")
row = MagicMock()
row.id = uuid.uuid4()
row.state = "claimed"
row.claim_owner = "w1"
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
row.attempt = 3
row.result = {}
session = AsyncMock()
session.get = AsyncMock(return_value=row)
out = await oq.fail_ops_run(
session, row.id, worker_id="w1", error="timeout", reopen=True
)
assert out is not None
assert row.state == "failed"
@pytest.mark.asyncio
async def test_fail_persists_redacted_failure_evidence() -> None:
from activity_core import ops_run_queue as oq
row = MagicMock()
row.state = "claimed"
row.claim_owner = "worker-1"
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
row.attempt = 1
row.result = {}
session = AsyncMock()
session.get = AsyncMock(return_value=row)
failed = await oq.fail_ops_run(
session,
uuid.uuid4(),
worker_id="worker-1",
error="execution refused",
result={
"ok": False,
"evidence": {
"request_id": "request-failed",
"profile_ref": "harness.agent-dev@1.0.0",
"rein_id": "rein-aharness",
"outcome": "failed",
"failure_stage": "execution",
"error": "execution failed; inspect direct caller error",
"duration_s": 3.5,
},
"tool_error": "raw provider failure must not persist",
},
)
assert failed is row
assert row.result["error"] == "execution refused"
assert row.result["execution_evidence"]["failure_stage"] == "execution"
assert "tool_error" not in row.result
@pytest.mark.asyncio
@pytest.mark.parametrize("lease_offset", [None, timedelta(0), timedelta(seconds=-1)])
async def test_heartbeat_rejects_missing_or_expired_lease(
lease_offset: timedelta | None,
) -> None:
from activity_core import ops_run_queue as oq
now = datetime(2026, 8, 23, 10, 0, tzinfo=timezone.utc)
row = MagicMock()
row.state = "claimed"
row.claim_owner = "worker-1"
row.lease_until = now + lease_offset if lease_offset is not None else None
session = AsyncMock()
session.get = AsyncMock(return_value=row)
with patch.object(oq, "_utcnow", return_value=now):
heartbeat = await oq.heartbeat_ops_run(
session,
uuid.uuid4(),
worker_id="worker-1",
lease_seconds=60,
)
assert heartbeat is None
assert row.lease_until == (now + lease_offset if lease_offset is not None else None)
session.get.assert_awaited_once()
assert session.get.await_args.kwargs == {"with_for_update": True}
@pytest.mark.asyncio
async def test_complete_rejects_expired_lease_without_mutation() -> None:
from activity_core import ops_run_queue as oq
now = datetime(2026, 8, 23, 10, 0, tzinfo=timezone.utc)
row = MagicMock()
row.state = "claimed"
row.claim_owner = "worker-1"
row.lease_until = now - timedelta(microseconds=1)
row.result = {"before": True}
session = AsyncMock()
session.get = AsyncMock(return_value=row)
with patch.object(oq, "_utcnow", return_value=now):
completed = await oq.complete_ops_run(
session,
uuid.uuid4(),
worker_id="worker-1",
result={"ok": True},
)
assert completed is None
assert row.state == "claimed"
assert row.result == {"before": True}
@pytest.mark.asyncio
async def test_fail_rejects_wrong_owner_with_active_lease() -> None:
from activity_core import ops_run_queue as oq
now = datetime(2026, 8, 23, 10, 0, tzinfo=timezone.utc)
row = MagicMock()
row.state = "claimed"
row.claim_owner = "worker-1"
row.lease_until = now + timedelta(minutes=1)
row.result = {"before": True}
session = AsyncMock()
session.get = AsyncMock(return_value=row)
with patch.object(oq, "_utcnow", return_value=now):
failed = await oq.fail_ops_run(
session,
uuid.uuid4(),
worker_id="worker-2",
error="must not persist",
)
assert failed is None
assert row.state == "claimed"
assert row.result == {"before": True}
def test_label_filter_any_vs_all() -> None:
"""Document labels_mode semantics used by claim_ops_runs."""
row_labels = {"automated", "research-brief"}
want = {"research-brief", "missing"}
assert bool(want & row_labels) # any
assert not want.issubset(row_labels) # all