mirror of
https://github.com/rennf93/roboco.git
synced 2026-08-03 07:23:24 +02:00
* fix(rate-limit): real provider liveness probe instead of time-based stub
The rate-limit recovery sweeper cleared a provider and resumed parked agents
purely on elapsed time — _do_probe was a stub that always returned True once
the retry_after window passed, so it never confirmed the provider had actually
stopped rate-limiting us. Under a sustained limit that resumes agents straight
into another 429, re-parking them: avoidable churn.
Make the probe real. _do_probe now issues a free, unmetered liveness call —
Anthropic GET /v1/models or Ollama GET /api/tags — and treats any non-429
response as the limit having lifted. A 429 keeps the provider parked; a
network error keeps it parked too (retry next sweep). When the provider can't
be probed (no API key, or an unrecognized provider), it falls back to the
prior time-expiry optimism rather than stranding agents. _probe_target keeps
URL/header resolution separate and testable, and _do_probe stays a
monkeypatchable boundary so the existing sweep tests are unaffected.
Also drop two acceptance-criteria-number labels from comments in this file.
* chore(rate-limit): clear merged gate debt in rate-limit tests + deps lint
The rate-limit PR landed with ruff violations the full gate flags but the
authors' runs missed: test_rate_limit_sweep.py was unformatted, and
test_rate_limit_tracker.py had unsorted/unused imports and magic-value
comparisons. Format the sweep test, drop the dead imports, and bind the
magic comparison values to locals. Also strip acceptance-criteria-number
labels from comments/docstrings across the three rate-limit test files
(leaving genuine acceptance_criteria=[...] test data untouched), and add
api/deps.py to the PLC0415 per-file-ignore — it is the DI wiring hub and
defers a couple of service imports to call time to avoid import cycles,
the same rationale already applied to api/routes, runtime, and services.
* fix(rate-limit): resolve redis type errors in RateLimitStateTracker
A cold mypy run (the gate's true state — prior passes were warm-cache only)
flagged four redis-typing errors in rate_limit_tracker.py that the merge
missed: three unused type:ignore[type-arg] on redis.Redis, and an
aclose() the bundled redis type stub doesn't expose.
Drop the now-unused ignores, and close the scan client via
'async with redis.from_url(...) as r:' instead of a finally-block
aclose(). The context manager closes the client on exit using the modern
redis.asyncio API — no deprecated close(), no stub-missing aclose(), no
suppression. Extend the test's redis mock to model the async
context-manager protocol so it returns itself on enter.
* test(prompter): pass route='main_pm' in the product main-PM routing test
Pre-existing master failure, unrelated to the rate-limit work. The test is
named ...product_routes_to_main_pm and asserts team=MAIN_PM, but called
confirm_live_draft without a route, so it got the 'board' default — which
assigns the Product Owner and yields team=BOARD by design (the board-review
path keeps the root at team=board until the CEO approves). The Main-PM path
is selected with route='main_pm', exactly as the sibling
...main_pm_route_assigns_main_pm test does. Add the missing kwarg so the test
verifies the path it names; behaviour under test is unchanged.
* Updated uv.lock
* refactor(complexity): bring all rank-C blocks under the xenon B ceiling
The full quality gate's xenon step (--max-absolute B --max-modules A
--max-average A) failed on eight rank-C blocks plus the extraction module
average — debt the rate-limit and token-analytics merges deferred. Reduce
each by extracting cohesive helpers, behaviour unchanged:
- orchestrator._probe_one_provider: split into _too_early_to_probe,
_on_probe_success, _on_probe_failure, _parked_agents_for.
- rate_limit_tracker.list_rate_limited_providers: extract _read_rate_limited_entry
and a _decode helper.
- trigger_filter.decide_spawn: extract _stale_trigger_decision (drops the
PLR0911 suppression too).
- ollama_embedder (embed_query, _embed_batch_sync, aembed_query,
_embed_batch_async): share _rl_backoff / _map_embed_error / _log_429 /
_sleep_connect_retry / _asleep_connect_retry; remove a dead post-loop guard
in aembed_query.
- mentor._synthesize_answer: extract _select_system_prompt and
_answer_from_response.
- indexes/base.ask: extract the 429-retried LLM call into _ask_llm.
- extraction.__init__: extract _compile_patterns so the module average
lands at rank A.
xenon now exits 0; rate-limit, optimal_brain, extraction, and events suites
all green.
* chore(deps): drop obsolete types-redis stub; honor redis 8.0 inline types
types-redis 4.6 (typed for redis 4.x) shadowed redis 8.0's own inline types,
which both masked real annotation mismatches in stream_bus.py and forced
awkward workarounds elsewhere. The stale stub is why the mypy gate only ever
passed warm-cached: a cold run under the wrong stub disagreed with the code.
Remove types-redis (and its orphaned transitive stubs) so mypy uses redis's
shipped types. That surfaces that xreadgroup/xclaim return bytes-keyed records
while _handle_message is annotated str — the code already decodes bytes
defensively, so this is an annotation gap, not a runtime bug. Make the types
honest: cast each result to its concrete shape and decode the stream name and
message id to str at the dispatch boundary via a _to_str helper.
mypy roboco/ is now clean cold (247 files) against redis's real types; events
suite green.
* Updated uv.lock
* fix(workspace): install the dev extra so agents can run make quality
Agent workspaces were set up with plain `uv sync`, which installs only the
project's default dependency group (pytest) — not the `dev` *extra* where the
gate tools live (ruff, mypy, xenon, radon, vulture, bandit, deptry). So an
agent's .venv had pytest but no linters, and `make quality` died immediately
on `ruff: command not found`. Agents literally could not lint, type-check, or
complexity-check their own work, which is how format/mypy/xenon debt merged
unseen. Sync the `dev` extra (`uv sync --extra dev`) so the workspace gets the
full toolchain the setup's own docstring already promised.
* fix(panel): rate-limit endpoint shape + websocket path
Two panel-facing breakages from the rate-limit rework:
- GET /api/system/rate-limits returned a raw list, but the panel store reads
response.entries — so `r.entries is not iterable` crashed the banner sync on
page load. Return the panel's contract: a { entries: [...] } envelope whose
items are camelCase {provider, affectedAgents, hitAt, resumeAt,
retryAfterSeconds}, derived from the raw Redis state (resumeAt = hitAt +
retryAfter).
- The rate-limit websocket hook passed "/ws/system" while getWebSocketUrl()
already supplies the "/ws" base, producing the doubled "/ws/ws/system" URL.
Pass "/system" to match the agents/channels/notifications hooks.
Note: the backend /ws/system endpoint itself does not yet exist (the rework
shipped the panel hook only); the REST fix keeps the banner correct on load
and reconnect until that endpoint is built.
* test(workspace): assert uv sync installs the dev extra
Follow the workspace setup change: the dependency-install command is now
`uv sync --extra dev` so the agent workspace gets the lint/type/complexity
toolchain. Update the three assertions that pinned the old `uv sync`.
* feat(ws): add /ws/system stream and bridge rate-limit events to the panel
The rate-limit rework shipped the panel's websocket hook but no backend: there
was no /ws/system endpoint and nothing forwarded RATE_LIMIT_HIT/LIFTED to a
socket, so the banner got no live updates.
Build the missing half:
- ConnectionManager grows a system-wide connection set with connect_system /
broadcast_system, and disconnect() now clears it.
- A /ws/system websocket endpoint (operator stream, no per-agent keying) with
the same connected + ping/pong lifecycle as the other streams.
- websocket_bridge subscribes RATE_LIMIT_HIT/LIFTED and forwards each to
broadcast_system tagged with the type the panel switches on. Both events
ride the same StreamEventBus singleton, and the subscriptions register
before start_listening(), so the consumer reads their streams.
Pairs with the panel hook now passing '/system' (getWebSocketUrl supplies the
'/ws' base). Covered by handler, manager, and endpoint-lifecycle tests.
---------
Co-authored-by: Renn F <rennf93@users.noreply.github.com>
563 lines
20 KiB
Python
563 lines
20 KiB
Python
"""Unit tests for the rate-limit sweeper probe loop.
|
|
|
|
Tests cover:
|
|
- probe-success path: tracker.clear() + resolve_wait + RATE_LIMIT_LIFTED event
|
|
- probe-failure path: increment_probe_failures is called
|
|
- CEO notification fires at threshold 10 exactly once per episode
|
|
- ``_do_probe`` / ``_make_tracker`` are injectable boundaries for mocking
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import fnmatch
|
|
import json
|
|
from datetime import UTC, datetime
|
|
from typing import Any
|
|
from unittest.mock import AsyncMock, MagicMock, patch
|
|
from uuid import uuid4
|
|
|
|
from httpx import ASGITransport, AsyncClient
|
|
from roboco.api.app import create_app
|
|
from roboco.models.events import EventType
|
|
from roboco.models.runtime import WaitingRecord
|
|
from roboco.runtime.orchestrator import AgentOrchestrator
|
|
from roboco.services.gateway.rate_limit_tracker import RateLimitStateTracker
|
|
|
|
_HTTP_OK = 200
|
|
_HTTP_NOT_FOUND = 404
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Helpers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _make_redis_mock(initial_store: dict[str, Any] | None = None) -> AsyncMock:
|
|
"""Fake redis.asyncio.Redis backed by a plain dict."""
|
|
store: dict[str, Any] = initial_store if initial_store is not None else {}
|
|
|
|
async def _get(key: str) -> bytes | None:
|
|
val = store.get(key)
|
|
if val is None:
|
|
return None
|
|
return str(val).encode() if not isinstance(val, bytes) else val
|
|
|
|
async def _set(key: str, value: Any) -> None:
|
|
store[key] = value
|
|
|
|
async def _delete(key: str) -> int:
|
|
return 1 if store.pop(key, None) is not None else 0
|
|
|
|
async def _scan(
|
|
_cursor: int,
|
|
match: str = "*",
|
|
count: int = 100, # noqa: ARG001
|
|
) -> tuple[int, list[bytes]]:
|
|
# Simple in-memory scan: return all matching keys in one shot
|
|
matches = [k.encode() for k in store if fnmatch.fnmatch(k, match)]
|
|
return (0, matches)
|
|
|
|
async def _aclose() -> None:
|
|
pass
|
|
|
|
mock = AsyncMock()
|
|
mock.get = AsyncMock(side_effect=_get)
|
|
mock.set = AsyncMock(side_effect=_set)
|
|
mock.delete = AsyncMock(side_effect=_delete)
|
|
mock.scan = AsyncMock(side_effect=_scan)
|
|
mock.aclose = AsyncMock(side_effect=_aclose)
|
|
# Support `async with redis.from_url(...) as r:` — the client returns
|
|
# itself on enter so the configured side-effects are what the caller uses.
|
|
mock.__aenter__ = AsyncMock(return_value=mock)
|
|
mock.__aexit__ = AsyncMock(return_value=False)
|
|
mock._store = store
|
|
return mock
|
|
|
|
|
|
def _make_orchestrator() -> AgentOrchestrator:
|
|
"""Build a minimal orchestrator via __new__ (no __init__ side-effects)."""
|
|
orch = AgentOrchestrator.__new__(AgentOrchestrator)
|
|
orch._running = True
|
|
orch._waiting_records: dict[str, WaitingRecord] = {}
|
|
orch._instances: dict[str, Any] = {}
|
|
orch._rate_limit_ceo_notified: set[str] = set()
|
|
return orch
|
|
|
|
|
|
def _make_tracker_mock(failure_return: int = 1) -> AsyncMock:
|
|
"""Create an async mock RateLimitStateTracker instance."""
|
|
mock = AsyncMock()
|
|
mock.clear = AsyncMock()
|
|
mock.increment_probe_failures = AsyncMock(return_value=failure_return)
|
|
mock.reset_probe_failures = AsyncMock()
|
|
return mock
|
|
|
|
|
|
def _make_active_state(
|
|
_provider: str = "anthropic",
|
|
retry_after: float | None = None,
|
|
probe_failures: int = 0,
|
|
activated_at: datetime | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Build a tracker state dict."""
|
|
at = activated_at or datetime.now(UTC)
|
|
return {
|
|
"rate_limited": True,
|
|
"activated_at": at.isoformat(),
|
|
"retry_after": retry_after,
|
|
"affected_agents": ["be-dev-1"],
|
|
"probe_failures": probe_failures,
|
|
}
|
|
|
|
|
|
def _waiting_record(
|
|
agent_id: str,
|
|
provider: str = "anthropic",
|
|
task_id: str | None = None,
|
|
) -> WaitingRecord:
|
|
return WaitingRecord(
|
|
agent_id=agent_id,
|
|
task_id=task_id or str(uuid4()),
|
|
waiting_for="rate_limit_lifted",
|
|
waiting_since=datetime.now(UTC),
|
|
context={"provider": provider},
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tests: probe-success path
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestProbeSuccessPath:
|
|
"""When _do_probe returns True the rate limit should be cleared and
|
|
all parked agents resolved."""
|
|
|
|
async def test_tracker_clear_called_on_success(self) -> None:
|
|
"""tracker.clear() is invoked when the probe succeeds."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
tracker_mock = _make_tracker_mock()
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
orch.resolve_wait = AsyncMock(return_value=None)
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return True
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
with patch("roboco.events.get_event_bus") as mock_bus_fn:
|
|
bus_mock = AsyncMock()
|
|
bus_mock.publish = AsyncMock()
|
|
mock_bus_fn.return_value = bus_mock
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
tracker_mock.clear.assert_awaited_once()
|
|
|
|
async def test_resolve_wait_called_for_parked_agents(self) -> None:
|
|
"""resolve_wait is called for each agent waiting for rate_limit_lifted."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
agent1 = "be-dev-1"
|
|
agent2 = "be-dev-2"
|
|
orch._waiting_records = {
|
|
agent1: _waiting_record(agent1, provider),
|
|
agent2: _waiting_record(agent2, provider),
|
|
"be-qa-1": _waiting_record(
|
|
"be-qa-1", "other-provider"
|
|
), # different provider
|
|
}
|
|
|
|
orch.resolve_wait = AsyncMock(return_value=None)
|
|
|
|
tracker_mock = _make_tracker_mock()
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return True
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
with patch("roboco.events.get_event_bus") as mock_bus_fn:
|
|
bus_mock = AsyncMock()
|
|
bus_mock.publish = AsyncMock()
|
|
mock_bus_fn.return_value = bus_mock
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
# Only the two anthropic-parked agents should be resolved
|
|
assert orch.resolve_wait.await_count == 2 # noqa: PLR2004
|
|
resolved_ids = {call.args[0] for call in orch.resolve_wait.call_args_list}
|
|
assert agent1 in resolved_ids
|
|
assert agent2 in resolved_ids
|
|
assert "be-qa-1" not in resolved_ids
|
|
|
|
async def test_rate_limit_lifted_event_published(self) -> None:
|
|
"""RATE_LIMIT_LIFTED event is published to the bus on probe success."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
orch.resolve_wait = AsyncMock(return_value=None)
|
|
|
|
tracker_mock = _make_tracker_mock()
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
|
|
published_events: list[Any] = []
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return True
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
with patch("roboco.events.get_event_bus") as mock_bus_fn:
|
|
bus_mock = AsyncMock()
|
|
bus_mock.publish = AsyncMock(side_effect=published_events.append)
|
|
mock_bus_fn.return_value = bus_mock
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
assert len(published_events) == 1
|
|
event = published_events[0]
|
|
assert event.type == EventType.RATE_LIMIT_LIFTED
|
|
assert event.data["provider"] == provider
|
|
|
|
async def test_ceo_notified_flag_cleared_on_success(self) -> None:
|
|
"""_rate_limit_ceo_notified is cleared when probe succeeds."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
orch._rate_limit_ceo_notified.add(provider) # simulates prior episode
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
orch.resolve_wait = AsyncMock(return_value=None)
|
|
|
|
tracker_mock = _make_tracker_mock()
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return True
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
with patch("roboco.events.get_event_bus") as mock_bus_fn:
|
|
bus_mock = AsyncMock()
|
|
bus_mock.publish = AsyncMock()
|
|
mock_bus_fn.return_value = bus_mock
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
assert provider not in orch._rate_limit_ceo_notified
|
|
|
|
async def test_probe_skipped_before_estimated_lift_at(self) -> None:
|
|
"""When retry_after has not elapsed yet the probe is skipped entirely."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
# Set activated_at to now; retry_after = 300s → estimated lift in future
|
|
state = _make_active_state(
|
|
provider,
|
|
retry_after=300.0,
|
|
activated_at=datetime.now(UTC),
|
|
)
|
|
|
|
probe_called = []
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
probe_called.append(_p)
|
|
return True
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
assert probe_called == [] # probe was gated by time
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tests: probe-failure path
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestProbeFailurePath:
|
|
"""When _do_probe returns False the failure counter should be incremented."""
|
|
|
|
async def test_increment_probe_failures_called_on_failure(self) -> None:
|
|
"""increment_probe_failures is called when the probe fails."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
tracker_mock = _make_tracker_mock(failure_return=1)
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return False
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
orch._notify_rate_limit_ceo = AsyncMock()
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
tracker_mock.increment_probe_failures.assert_awaited_once()
|
|
|
|
async def test_clear_not_called_on_failure(self) -> None:
|
|
"""tracker.clear() must NOT be called when the probe fails."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
tracker_mock = _make_tracker_mock(failure_return=1)
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return False
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
orch._notify_rate_limit_ceo = AsyncMock()
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
tracker_mock.clear.assert_not_awaited()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tests: CEO notification threshold
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestCEONotificationThreshold:
|
|
"""CEO notification fires at count==10 exactly once per episode."""
|
|
|
|
async def test_notification_fires_at_exactly_10_failures(self) -> None:
|
|
"""_notify_rate_limit_ceo is called when failure count hits 10."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
# simulate already at 9 failures; next increment returns 10
|
|
tracker_mock = _make_tracker_mock(failure_return=10)
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
orch._notify_rate_limit_ceo = AsyncMock()
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return False
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
orch._notify_rate_limit_ceo.assert_awaited_once()
|
|
|
|
async def test_notification_not_fired_before_threshold(self) -> None:
|
|
"""No CEO notification below threshold 10."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
tracker_mock = _make_tracker_mock(failure_return=9)
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
orch._notify_rate_limit_ceo = AsyncMock()
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return False
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
orch._notify_rate_limit_ceo.assert_not_awaited()
|
|
|
|
async def test_notification_sent_only_once_per_episode(self) -> None:
|
|
"""Even if failures keep accumulating, the CEO is notified only once."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
state = _make_active_state(provider, retry_after=None)
|
|
|
|
# Mark this episode as already notified
|
|
orch._rate_limit_ceo_notified.add(provider)
|
|
|
|
tracker_mock = _make_tracker_mock(failure_return=15)
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
orch._notify_rate_limit_ceo = AsyncMock()
|
|
|
|
async def fake_do_probe(_p: str) -> bool:
|
|
return False
|
|
|
|
orch._do_probe = fake_do_probe # type: ignore[method-assign]
|
|
|
|
await orch._probe_one_provider(provider, state)
|
|
|
|
orch._notify_rate_limit_ceo.assert_not_awaited()
|
|
|
|
async def test_new_episode_allows_new_notification(self) -> None:
|
|
"""After a rate-limit clears (success) a new episode starts fresh."""
|
|
orch = _make_orchestrator()
|
|
provider = "anthropic"
|
|
# Episode 1: had a notification
|
|
orch._rate_limit_ceo_notified.add(provider)
|
|
|
|
success_state = _make_active_state(provider, retry_after=None)
|
|
orch.resolve_wait = AsyncMock(return_value=None)
|
|
|
|
tracker_mock = _make_tracker_mock(failure_return=10)
|
|
orch._make_tracker = MagicMock(return_value=tracker_mock) # type: ignore[method-assign]
|
|
notify_mock = AsyncMock()
|
|
orch._notify_rate_limit_ceo = notify_mock
|
|
|
|
async def fake_do_probe_success(_p: str) -> bool:
|
|
return True
|
|
|
|
orch._do_probe = fake_do_probe_success # type: ignore[method-assign]
|
|
|
|
with patch("roboco.events.get_event_bus") as mock_bus_fn:
|
|
bus_mock = AsyncMock()
|
|
bus_mock.publish = AsyncMock()
|
|
mock_bus_fn.return_value = bus_mock
|
|
|
|
# Success clears the episode flag
|
|
await orch._probe_one_provider(provider, success_state)
|
|
|
|
assert provider not in orch._rate_limit_ceo_notified
|
|
|
|
# Episode 2: simulate a new failure reaching threshold 10
|
|
async def fake_do_probe_fail(_p: str) -> bool:
|
|
return False
|
|
|
|
orch._do_probe = fake_do_probe_fail # type: ignore[method-assign]
|
|
|
|
failure_state = _make_active_state(provider, retry_after=None)
|
|
await orch._probe_one_provider(provider, failure_state)
|
|
|
|
# Notification SHOULD fire for the new episode
|
|
notify_mock.assert_awaited_once()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tests: list_rate_limited_providers
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestListRateLimitedProviders:
|
|
"""list_rate_limited_providers scans Redis for active rate-limit keys."""
|
|
|
|
async def test_returns_empty_when_no_keys(self) -> None:
|
|
redis_mock = _make_redis_mock()
|
|
with patch("redis.asyncio.from_url", return_value=redis_mock):
|
|
result = await RateLimitStateTracker.list_rate_limited_providers()
|
|
assert result == []
|
|
|
|
async def test_returns_active_provider(self) -> None:
|
|
state = {
|
|
"rate_limited": True,
|
|
"activated_at": datetime.now(UTC).isoformat(),
|
|
"retry_after": 60.0,
|
|
"affected_agents": ["be-dev-1"],
|
|
"probe_failures": 0,
|
|
}
|
|
store = {"roboco:rate_limit:anthropic:state": json.dumps(state).encode()}
|
|
redis_mock = _make_redis_mock(store)
|
|
|
|
with patch("redis.asyncio.from_url", return_value=redis_mock):
|
|
result = await RateLimitStateTracker.list_rate_limited_providers()
|
|
|
|
assert len(result) == 1
|
|
provider, returned_state = result[0]
|
|
assert provider == "anthropic"
|
|
assert returned_state["rate_limited"] is True
|
|
|
|
async def test_ignores_cleared_providers(self) -> None:
|
|
state = {
|
|
"rate_limited": False,
|
|
"activated_at": datetime.now(UTC).isoformat(),
|
|
"retry_after": 60.0,
|
|
"affected_agents": [],
|
|
"probe_failures": 2,
|
|
}
|
|
store = {"roboco:rate_limit:anthropic:state": json.dumps(state).encode()}
|
|
redis_mock = _make_redis_mock(store)
|
|
|
|
with patch("redis.asyncio.from_url", return_value=redis_mock):
|
|
result = await RateLimitStateTracker.list_rate_limited_providers()
|
|
|
|
assert result == []
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tests: GET /api/system/rate-limits endpoint schema
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestRateLimitsEndpoint:
|
|
"""GET /api/system/rate-limits returns correct schema."""
|
|
|
|
async def test_returns_empty_list_when_no_rate_limits(self) -> None:
|
|
app = create_app()
|
|
|
|
with patch(
|
|
"roboco.api.routes.system.RateLimitStateTracker"
|
|
".list_rate_limited_providers",
|
|
new_callable=AsyncMock,
|
|
return_value=[],
|
|
):
|
|
async with AsyncClient(
|
|
transport=ASGITransport(app=app), base_url="http://test"
|
|
) as client:
|
|
resp = await client.get("/api/system/rate-limits")
|
|
|
|
assert resp.status_code == _HTTP_OK
|
|
assert resp.json() == {"entries": []}
|
|
|
|
async def test_returns_provider_state_when_rate_limited(self) -> None:
|
|
app = create_app()
|
|
|
|
retry_after = 60.0
|
|
state = {
|
|
"rate_limited": True,
|
|
"activated_at": "2026-06-11T00:00:00+00:00",
|
|
"retry_after": retry_after,
|
|
"affected_agents": ["be-dev-1"],
|
|
"probe_failures": 3,
|
|
}
|
|
|
|
with patch(
|
|
"roboco.api.routes.system.RateLimitStateTracker"
|
|
".list_rate_limited_providers",
|
|
new_callable=AsyncMock,
|
|
return_value=[("anthropic", state)],
|
|
):
|
|
async with AsyncClient(
|
|
transport=ASGITransport(app=app), base_url="http://test"
|
|
) as client:
|
|
resp = await client.get("/api/system/rate-limits")
|
|
|
|
assert resp.status_code == _HTTP_OK
|
|
entries = resp.json()["entries"]
|
|
assert len(entries) == 1
|
|
entry = entries[0]
|
|
# Panel-shaped, camelCase fields (not the raw Redis state).
|
|
assert entry["provider"] == "anthropic"
|
|
assert entry["affectedAgents"] == ["be-dev-1"]
|
|
assert entry["hitAt"] == "2026-06-11T00:00:00+00:00"
|
|
assert entry["retryAfterSeconds"] == retry_after
|
|
assert entry["resumeAt"] == "2026-06-11T00:01:00+00:00"
|
|
|
|
async def test_endpoint_not_404(self) -> None:
|
|
"""The endpoint must be registered in app.py — no 404."""
|
|
app = create_app()
|
|
|
|
with patch(
|
|
"roboco.api.routes.system.RateLimitStateTracker"
|
|
".list_rate_limited_providers",
|
|
new_callable=AsyncMock,
|
|
return_value=[],
|
|
):
|
|
async with AsyncClient(
|
|
transport=ASGITransport(app=app), base_url="http://test"
|
|
) as client:
|
|
resp = await client.get("/api/system/rate-limits")
|
|
|
|
assert resp.status_code != _HTTP_NOT_FOUND
|