Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
timokoessler marked this conversation as resolved.


def listen_for_config_updates(connection_manager, event_scheduler):
"""
Expand All @@ -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
Expand All @@ -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
Comment thread
timokoessler marked this conversation as resolved.

logger.debug("SSE config-updated event, fetching new config")

try:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading