Classify heartbeat lease failures

Assistant: codex
Assistant-Model: gpt-5.6-sol
Assistant-Session: 01a02b6f-7db1-7222-918b-e813a6bda38d
This commit is contained in:
tegwick 2026-08-23 14:42:30 +02:00
parent ad4074c20a
commit 84b4089f8a
6 changed files with 84 additions and 14 deletions

View file

@ -80,13 +80,22 @@ class _Heartbeat:
self._client.heartbeat(self._run_id, lease_seconds=self._lease) self._client.heartbeat(self._run_id, lease_seconds=self._lease)
logger.info("heartbeat ok run_id=%s", self._run_id) logger.info("heartbeat ok run_id=%s", self._run_id)
except OpsRunError as exc: except OpsRunError as exc:
evidence = self.monitor.mark_lost(type(exc).__name__) if exc.lease_rejected:
logger.warning( evidence = self.monitor.mark_lost(
"heartbeat lost run_id=%s error_type=%s observed_at=%s", f"{type(exc).__name__}:{exc.status_code}"
self._run_id, )
evidence.error_type, logger.warning(
evidence.observed_at, "heartbeat lost run_id=%s error_type=%s observed_at=%s",
) self._run_id,
evidence.error_type,
evidence.observed_at,
)
else:
logger.warning(
"heartbeat transient failure run_id=%s error_type=%s",
self._run_id,
type(exc).__name__,
)
self._thread = threading.Thread( self._thread = threading.Thread(
target=_loop, name=f"ops-hb-{self._run_id[:8]}", daemon=True target=_loop, name=f"ops-hb-{self._run_id[:8]}", daemon=True

View file

@ -33,7 +33,22 @@ DEFAULT_REPO_ROOTS = ("~", "~/work")
class OpsRunError(RuntimeError): class OpsRunError(RuntimeError):
pass """Bounded Activity Core failure with optional HTTP classification."""
def __init__(
self,
message: str,
*,
action: str | None = None,
status_code: int | None = None,
) -> None:
super().__init__(message)
self.action = action
self.status_code = status_code
@property
def lease_rejected(self) -> bool:
return self.status_code in {401, 403, 404, 409, 410, 412, 423}
@dataclass @dataclass
@ -159,8 +174,14 @@ class ActivityCoreOpsClient:
timeout=self.config.timeout, timeout=self.config.timeout,
) )
resp.raise_for_status() resp.raise_for_status()
except httpx.HTTPStatusError as exc:
raise OpsRunError(
f"list ops-runs failed: HTTP {exc.response.status_code}",
action="list",
status_code=exc.response.status_code,
) from exc
except httpx.HTTPError as exc: except httpx.HTTPError as exc:
raise OpsRunError(f"list ops-runs failed: {exc}") from exc raise OpsRunError("list ops-runs failed: transport error", action="list") from exc
data = resp.json() data = resp.json()
items = data.get("items") if isinstance(data, dict) else data items = data.get("items") if isinstance(data, dict) else data
if not isinstance(items, list): if not isinstance(items, list):
@ -190,8 +211,14 @@ class ActivityCoreOpsClient:
timeout=self.config.timeout, timeout=self.config.timeout,
) )
resp.raise_for_status() resp.raise_for_status()
except httpx.HTTPStatusError as exc:
raise OpsRunError(
f"claim failed: HTTP {exc.response.status_code}",
action="claim",
status_code=exc.response.status_code,
) from exc
except httpx.HTTPError as exc: except httpx.HTTPError as exc:
raise OpsRunError(f"claim failed: {exc}") from exc raise OpsRunError("claim failed: transport error", action="claim") from exc
data = resp.json() data = resp.json()
items = data.get("items") if isinstance(data, dict) else [] items = data.get("items") if isinstance(data, dict) else []
return [OpsRun.from_api(item) for item in items if isinstance(item, dict)] return [OpsRun.from_api(item) for item in items if isinstance(item, dict)]
@ -247,8 +274,16 @@ class ActivityCoreOpsClient:
timeout=self.config.timeout, timeout=self.config.timeout,
) )
resp.raise_for_status() resp.raise_for_status()
except httpx.HTTPStatusError as exc:
raise OpsRunError(
f"{action} ops_run {run_id} failed: HTTP {exc.response.status_code}",
action=action,
status_code=exc.response.status_code,
) from exc
except httpx.HTTPError as exc: except httpx.HTTPError as exc:
raise OpsRunError(f"{action} ops_run {run_id} failed: {exc}") from exc raise OpsRunError(
f"{action} ops_run {run_id} failed: transport error", action=action
) from exc
return OpsRun.from_api(resp.json()) return OpsRun.from_api(resp.json())

View file

@ -110,7 +110,7 @@ def test_process_one_refuses_close_after_lease_loss() -> None:
client = MagicMock(spec=ActivityCoreOpsClient) client = MagicMock(spec=ActivityCoreOpsClient)
client.config = OpsRunConfig(worker_id="w", lease_seconds=90) client.config = OpsRunConfig(worker_id="w", lease_seconds=90)
client.claim.return_value = [_claimed_run()] client.claim.return_value = [_claimed_run()]
client.heartbeat.side_effect = OpsRunError("lease rejected") client.heartbeat.side_effect = OpsRunError("lease rejected", status_code=409)
ar = ApproachResult( ar = ApproachResult(
ok=True, ok=True,
approach=APPROACH_FI_RESEARCH_BRIEF, approach=APPROACH_FI_RESEARCH_BRIEF,
@ -159,7 +159,7 @@ def test_profiled_exception_after_lease_loss_skips_close() -> None:
run = _claimed_run() run = _claimed_run()
run.harness_profile_ref = "harness.agent-dev-local@1.0.0" run.harness_profile_ref = "harness.agent-dev-local@1.0.0"
client.claim.return_value = [run] client.claim.return_value = [run]
client.heartbeat.side_effect = OpsRunError("lease rejected") client.heartbeat.side_effect = OpsRunError("lease rejected", status_code=409)
def slow_profile(*_args, **_kwargs): def slow_profile(*_args, **_kwargs):
time.sleep(0.05) time.sleep(0.05)

View file

@ -50,7 +50,7 @@ def test_concurrent_loss_still_notifies_once() -> None:
def test_heartbeat_publishes_loss_without_raw_exception_text() -> None: def test_heartbeat_publishes_loss_without_raw_exception_text() -> None:
client = MagicMock() client = MagicMock()
client.heartbeat.side_effect = OpsRunError("secret provider response") client.heartbeat.side_effect = OpsRunError("secret provider response", status_code=409)
heartbeat = _Heartbeat(client, "run-1", lease_seconds=90) heartbeat = _Heartbeat(client, "run-1", lease_seconds=90)
with patch("rein_aharness.claim_loop._heartbeat_interval", return_value=0.01): with patch("rein_aharness.claim_loop._heartbeat_interval", return_value=0.01):
@ -63,3 +63,16 @@ def test_heartbeat_publishes_loss_without_raw_exception_text() -> None:
assert heartbeat.monitor.lost is True assert heartbeat.monitor.lost is True
assert heartbeat.monitor.evidence is not None assert heartbeat.monitor.evidence is not None
assert "secret" not in heartbeat.monitor.evidence.error_type assert "secret" not in heartbeat.monitor.evidence.error_type
def test_heartbeat_transport_error_does_not_immediately_mark_loss() -> None:
client = MagicMock()
client.heartbeat.side_effect = OpsRunError("connection reset")
heartbeat = _Heartbeat(client, "run-1", lease_seconds=90)
with patch("rein_aharness.claim_loop._heartbeat_interval", return_value=0.01):
heartbeat.start()
time.sleep(0.04)
heartbeat.stop()
assert heartbeat.monitor.lost is False

View file

@ -137,6 +137,16 @@ def test_claim_http_error() -> None:
client.claim() client.claim()
def test_ops_run_error_classifies_lease_rejection_without_message_details() -> None:
rejected = OpsRunError("provider response included credentials", status_code=409)
transient = OpsRunError("connection reset")
assert rejected.lease_rejected is True
assert rejected.status_code == 409
assert transient.lease_rejected is False
assert transient.status_code is None
def test_ops_run_to_taskspec(tmp_path: Path) -> None: def test_ops_run_to_taskspec(tmp_path: Path) -> None:
import subprocess import subprocess

View file

@ -207,6 +207,9 @@ evidence and skips Activity Core completion/failure calls, leaving reconciliatio
to the owner of the expired lease. Adapter exception paths now also stop before to the owner of the expired lease. Adapter exception paths now also stop before
close when lease loss raced the failure; otherwise they return only an exception close when lease loss raced the failure; otherwise they return only an exception
class marker rather than leaving an unbound result or retaining provider text. class marker rather than leaving an unbound result or retaining provider text.
`OpsRunError` now carries bounded action/status metadata, and heartbeats treat
explicit ownership-rejection statuses as lease loss while allowing transport
errors to retry without immediately abandoning a run.
## Verify accepted commits and reconcile metrics/reporting ## Verify accepted commits and reconcile metrics/reporting