#!/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()