diff --git a/aikido_zen/background_process/realtime/listen_for_config_updates.py b/aikido_zen/background_process/realtime/listen_for_config_updates.py index fab7d6aff..4179df1ac 100644 --- a/aikido_zen/background_process/realtime/listen_for_config_updates.py +++ b/aikido_zen/background_process/realtime/listen_for_config_updates.py @@ -3,12 +3,16 @@ """ import json +import time from aikido_zen.helpers.token import Token from aikido_zen.helpers.logging import logger import aikido_zen.background_process.realtime as realtime from .sse_client import connect_to_sse +# Zen Realtime API should not send config-updated events more than once every 10 seconds +CONFIG_REFRESH_THROTTLE_SECONDS = 9 + def listen_for_config_updates(connection_manager, event_scheduler): """ @@ -26,6 +30,20 @@ def listen_for_config_updates(connection_manager, event_scheduler): token = connection_manager.token last_updated_at = connection_manager.conf.last_updated_at + last_config_refresh_started_at = None + + def config_update_arrived_too_fast(): + nonlocal last_config_refresh_started_at + + now = time.monotonic() + if ( + last_config_refresh_started_at is not None + and now - last_config_refresh_started_at < CONFIG_REFRESH_THROTTLE_SECONDS + ): + return True + + last_config_refresh_started_at = now + return False def on_event(event): nonlocal last_updated_at @@ -43,6 +61,10 @@ def on_event(event): logger.debug("SSE config-updated event has invalid payload: %s", event.data) return + if config_update_arrived_too_fast(): + logger.debug("SSE config-updated event ignored by refresh throttle") + return + logger.debug("SSE config-updated event, fetching new config") try: diff --git a/aikido_zen/background_process/realtime/listen_for_config_updates_test.py b/aikido_zen/background_process/realtime/listen_for_config_updates_test.py index 5a3c3010e..0c25678ab 100644 --- a/aikido_zen/background_process/realtime/listen_for_config_updates_test.py +++ b/aikido_zen/background_process/realtime/listen_for_config_updates_test.py @@ -188,6 +188,22 @@ def test_updates_last_updated_at_so_a_second_stale_event_is_ignored( mock_get_config.assert_called_once() +def test_throttles_newer_events(connection_manager): + connection_manager.conf.last_updated_at = 100 + event_scheduler = make_inline_scheduler() + on_event = get_on_event(connection_manager, event_scheduler) + + new_config = {"endpoints": [], "configUpdatedAt": 200} + with patch( + "aikido_zen.background_process.realtime.get_config", return_value=new_config + ) as mock_get_config: + on_event(make_event(data={"configUpdatedAt": 200})) + on_event(make_event(data={"configUpdatedAt": 300})) + + mock_get_config.assert_called_once_with(connection_manager.token) + connection_manager.update_firewall_lists.assert_called_once() + + def test_handles_get_config_failure_gracefully(connection_manager, caplog): connection_manager.conf.last_updated_at = 100 on_event = get_on_event(connection_manager)