Files
box/tests/test_rate_limits.py
operator 0065d11e97 feat(supervision): add choice watcher daemon, HTTPS spec docs, and test suites
- bin/muse_choice_watcher.py + systemd/muse-choices-reconcile.*: automatic choice answering and timer reconciliation
- bin/digest.py: fleet log and health summarization
- docs/BOX-*-HTTPS.md: comprehensive HTTPS execution contracts and API documentation
- docs/MUSE-CHOICES-POLICY.md & docs/SUPERVISION-SPEC.md: autonomous execution specs
- tests/test_*.py: unit test suites for HTTPS API, choice watcher, fleet heal, and swarm pruning
2026-10-07 00:25:46 +00:00

216 lines
10 KiB
Python

#!/usr/bin/env python3
"""
test_rate_limits.py — Tests for rate-limiting, adaptive polling, backoff, and circuit breaker.
"""
import curses
import importlib.util
import json
import os
import sys
import time
import unittest
from pathlib import Path
from unittest.mock import MagicMock, patch
REPO_ROOT = Path(__file__).resolve().parent.parent
MUSE_TUI_PATH = REPO_ROOT / "bin" / "muse-tui.py"
spec = importlib.util.spec_from_file_location("muse_tui_rl", MUSE_TUI_PATH)
muse_tui_rl = importlib.util.module_from_spec(spec)
sys.modules["muse_tui_rl"] = muse_tui_rl
spec.loader.exec_module(muse_tui_rl)
FleetDataManager = muse_tui_rl.FleetDataManager
MuseTUI = muse_tui_rl.MuseTUI
run_command_isolated = muse_tui_rl.run_command_isolated
class TestRateLimitingAndBackoff(unittest.TestCase):
def setUp(self):
# In tests, start_poller is disabled by default
self.mgr = FleetDataManager(start_poller=False)
def test_poller_disabled_in_tests_by_default(self):
"""FleetDataManager does not leak background poller threads during unit test runs."""
self.assertIsNone(self.mgr.poller_thread)
def test_rate_limit_cooldown_marking(self):
"""Marking a node as rate-limited sets cooldown window and rate_limited flag."""
self.assertFalse(self.mgr.is_node_rate_limited("muse"))
self.assertEqual(self.mgr.get_node_cooldown_remaining("muse"), 0)
self.mgr.mark_node_rate_limited("muse", cooldown_seconds=30.0, reason="HTTP 429")
self.assertTrue(self.mgr.is_node_rate_limited("muse"))
self.assertGreater(self.mgr.get_node_cooldown_remaining("muse"), 20)
self.assertTrue(self.mgr.node_status["muse"]["rate_limited"])
def test_tiered_backoff_progression(self):
"""Consecutive rate-limit occurrences escalate (30s -> 60s -> 120s max), and clear resets."""
# Incident 1: 30s
self.mgr.mark_node_rate_limited("pip", reason="429 first")
cd1 = self.mgr.get_node_cooldown_remaining("pip")
self.assertTrue(25 <= cd1 <= 31)
# Incident 2: 60s
self.mgr.mark_node_rate_limited("pip", reason="429 second")
cd2 = self.mgr.get_node_cooldown_remaining("pip")
self.assertTrue(55 <= cd2 <= 61)
# Incident 3: 120s max
self.mgr.mark_node_rate_limited("pip", reason="429 third")
cd3 = self.mgr.get_node_cooldown_remaining("pip")
self.assertTrue(110 <= cd3 <= 121)
# Success resets backoff
self.mgr.clear_node_rate_limit("pip")
self.assertFalse(self.mgr.is_node_rate_limited("pip"))
self.assertEqual(self.mgr.get_node_cooldown_remaining("pip"), 0)
self.assertEqual(self.mgr.rate_limit_consecutive["pip"], 0)
def test_retry_after_header_parsing(self):
"""When output includes Retry-After, that explicit duration is used."""
with patch.object(muse_tui_rl, "run_command_isolated", return_value=(1, "", "Error 429: Too Many Requests. Retry-After: 85")):
self.mgr._fetch_history("646", "thread-xyz")
self.assertTrue(self.mgr.is_node_rate_limited("646"))
cd = self.mgr.get_node_cooldown_remaining("646")
self.assertTrue(80 <= cd <= 86)
def test_rate_limited_node_skips_fetches(self):
"""Rate-limited nodes are skipped by _fetch_history and _run_fetch_threads."""
self.mgr.mark_node_rate_limited("muse", cooldown_seconds=60.0)
with patch.object(muse_tui_rl, "run_command_isolated") as mock_cmd:
self.mgr._fetch_history("muse", "thread-123")
mock_cmd.assert_not_called()
self.mgr._run_fetch_threads("muse")
mock_cmd.assert_not_called()
def test_rate_limit_detection_from_command_output(self):
"""When CLI command outputs 429 or rate limit text, node is placed on cooldown."""
with patch.object(muse_tui_rl, "run_command_isolated", return_value=(1, "", "Error 429: Too Many Requests")):
self.mgr._fetch_history("pip", "thread-abc")
self.assertTrue(self.mgr.is_node_rate_limited("pip"))
self.assertGreater(self.mgr.get_node_cooldown_remaining("pip"), 0)
def test_record_activity_updates_timestamp(self):
"""record_activity updates last_user_activity timestamp."""
old_time = self.mgr.last_user_activity
time.sleep(0.01)
self.mgr.record_activity()
self.assertGreaterEqual(self.mgr.last_user_activity, old_time)
def test_run_command_isolated_handles_quick_command(self):
"""run_command_isolated runs a command and returns returncode, stdout, stderr."""
rc, stdout, stderr = run_command_isolated(["echo", "hello rate limit"], timeout=2.0)
self.assertEqual(rc, 0)
self.assertIn("hello rate limit", stdout)
def test_run_command_isolated_terminates_on_timeout(self):
"""run_command_isolated cleanly kills process group on timeout without zombies."""
rc, stdout, stderr = run_command_isolated(["sleep", "10"], timeout=0.1)
self.assertEqual(rc, -1)
self.assertIn("timed out", stderr)
class TestRateLimitUserInteraction(unittest.TestCase):
def setUp(self):
self.mock_stdscr = MagicMock()
self.mock_stdscr.getmaxyx.return_value = (30, 100)
self.mock_stdscr.getch.return_value = -1
self.tui = MuseTUI(self.mock_stdscr, initial_mode="muse", initial_node="muse")
self.tui.safe_addstr = MagicMock()
def test_soft_guardrail_send_during_cooldown(self):
"""Sending during cooldown warns on first Enter, preserves buffer, and forces on second Enter."""
self.tui.data.mark_node_rate_limited("muse", cooldown_seconds=30.0)
self.tui.editor_mode = "INSERT"
self.tui.input_buf = "Status report please"
self.tui.input_cursor = len(self.tui.input_buf)
# First Enter tap: warns and does not clear input_buf
with patch.object(self.tui, "_async_send_message") as mock_send:
self.tui._handle_insert_key(10)
mock_send.assert_not_called()
self.assertEqual(self.tui.input_buf, "Status report please")
self.assertIn("cooldown", self.tui.toast_msg)
self.assertIn("Press Enter again", self.tui.toast_msg)
# Second Enter tap within 2.5s: bypasses cooldown, loads into history, and delivers
with patch("threading.Thread") as mock_thread:
self.tui._handle_insert_key(10)
self.assertFalse(self.tui.data.is_node_rate_limited("muse"))
self.assertIn("Sending message to MUSE", self.tui.toast_msg)
self.assertEqual(self.tui.input_buf, "")
# Ensure message was immediately loaded into history cache!
msgs = self.tui.data.history_cache.get(("muse", self.tui.data.active_thread_id), [])
self.assertTrue(any(m.get("text") == "Status report please" and m.get("role") == "user" for m in msgs))
def test_two_tap_manual_sync_override(self):
"""Pressing 'r' during cooldown warns on first press and bypasses cooldown on double-tap."""
self.tui.data.mark_node_rate_limited("muse", cooldown_seconds=45.0)
# First 'r' tap: warns
self.tui._handle_normal_key(ord('r'))
self.assertTrue(self.tui.data.is_node_rate_limited("muse"))
self.assertIn("cooling down", self.tui.toast_msg)
self.assertIn("Press 'r' again", self.tui.toast_msg)
# Second 'r' tap within 2s: clears rate limit and forces sync
with patch.object(self.tui.data, "lazy_fetch_threads") as mock_fetch:
self.tui._handle_normal_key(ord('r'))
self.assertFalse(self.tui.data.is_node_rate_limited("muse"))
self.assertIn("Force-syncing", self.tui.toast_msg)
mock_fetch.assert_called_with("muse", force=True)
def test_agent_context_menu_sync_override(self):
"""Context menu 'sync' action also respects two-tap override during cooldown."""
self.tui.data.mark_node_rate_limited("pip", cooldown_seconds=30.0)
self.tui.context_agent = {"node": "pip"}
# First selection: warns
self.tui._execute_agent_action("sync")
self.assertTrue(self.tui.data.is_node_rate_limited("pip"))
self.assertIn("cooling down", self.tui.toast_msg)
# Second selection within 2s: forces sync
with patch.object(self.tui.data, "lazy_fetch_threads") as mock_fetch:
self.tui._execute_agent_action("sync")
self.assertFalse(self.tui.data.is_node_rate_limited("pip"))
self.assertIn("Force-syncing", self.tui.toast_msg)
mock_fetch.assert_called_with("pip", force=True)
def test_optimistic_message_persistence_across_fetch(self):
"""Optimistic user messages are preserved even if server history lags behind."""
self.tui.data.history_cache[("muse", "sess-test")] = []
self.tui.data.add_optimistic_message("muse", "sess-test", "New uncommitted instruction")
cached = self.tui.data.history_cache.get(("muse", "sess-test"), [])
self.assertEqual(len(cached), 1)
self.assertEqual(cached[0]["text"], "New uncommitted instruction")
# Simulate remote fetch returning older history that doesn't yet have the new message
old_server_msgs = [{"role": "assistant", "text": "Earlier reply"}]
with patch.object(muse_tui_rl, "run_command_isolated", return_value=(0, json.dumps(old_server_msgs), "")):
self.tui.data._fetch_history("muse", "sess-test")
updated = self.tui.data.history_cache.get(("muse", "sess-test"), [])
# Both the older server message AND the pending user message should exist!
self.assertEqual(len(updated), 2)
self.assertEqual(updated[0]["text"], "Earlier reply")
self.assertEqual(updated[1]["text"], "New uncommitted instruction")
def test_history_payload_mentioning_rate_limits_not_falsely_flagged(self):
"""Chat history discussing rate limits or HTTP 429 does NOT place node into cooldown when rc == 0."""
msgs_about_rate_limits = [
{"role": "user", "text": "Are we hitting any rate limits or 429 errors?"},
{"role": "assistant", "text": "No rate limit reached, all queues are unthrottled and healthy."}
]
with patch.object(muse_tui_rl, "run_command_isolated", return_value=(0, json.dumps(msgs_about_rate_limits), "")):
self.tui.data._fetch_history("pip", "thread-chat")
self.assertFalse(self.tui.data.is_node_rate_limited("pip"))
self.assertEqual(len(self.tui.data.history_cache.get(("pip", "thread-chat"), [])), 2)
if __name__ == "__main__":
import json
unittest.main()