2026-08-03 19:22:50 +02:00
|
|
|
"""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,
|
2026-08-05 15:30:00 +02:00
|
|
|
emit_triggering_event_id,
|
2026-08-03 19:22:50 +02:00
|
|
|
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"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
2026-08-05 15:30:00 +02:00
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
2026-08-03 19:22:50 +02:00
|
|
|
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"
|
|
|
|
|
|
|
|
|
|
|
2026-08-22 23:11:52 +02:00
|
|
|
@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"
|
2026-08-23 13:01:46 +02:00
|
|
|
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
|
2026-08-22 23:11:52 +02:00
|
|
|
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,
|
|
|
|
|
},
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
2026-08-03 19:22:50 +02:00
|
|
|
@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"
|
2026-08-23 13:01:46 +02:00
|
|
|
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
|
2026-08-03 19:22:50 +02:00
|
|
|
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"
|
2026-08-23 13:01:46 +02:00
|
|
|
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
|
2026-08-03 19:22:50 +02:00
|
|
|
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"
|
|
|
|
|
|
|
|
|
|
|
2026-08-22 23:11:52 +02:00
|
|
|
@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"
|
2026-08-23 13:01:46 +02:00
|
|
|
row.lease_until = datetime.now(timezone.utc) + timedelta(minutes=1)
|
2026-08-22 23:11:52 +02:00
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
2026-08-23 13:01:46 +02:00
|
|
|
@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}
|
|
|
|
|
|
|
|
|
|
|
2026-08-03 19:22:50 +02:00
|
|
|
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
|