import json from unittest.mock import call import pytest from stream.exceptions import RateLimitReached from getstream_connector.throttling import DummyThrottler from tests import constants from tests.functional import data def test_add_activity(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 1 activity_data = json.loads(request_mock.last_request_data["data"]) assert activity_data["some"] == "data" @pytest.mark.parametrize("stream_client", ({"inject_trace_context": False},), indirect=True) def test_add_activity_inject_trace_context_disabled(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 1 activity_data = json.loads(request_mock.last_request_data["data"]) assert constants.TRACE_CONTEXT_KEY not in activity_data def test_add_activity_trace_context_injected(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 1 activity_data = json.loads(request_mock.last_request_data["data"]) assert constants.TRACE_CONTEXT_KEY in activity_data @pytest.mark.parametrize("stream_client", ({"inject_correlation_id": False},), indirect=True) def test_add_activity_inject_correlation_id_disabled(stream_client, stream_feed, request_mock): correlation_id = "corr_id_123" stream_client.set_correlation_id(correlation_id) request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 1 activity_data = json.loads(request_mock.last_request_data["data"]) assert constants.CORRELATION_ID_KEY not in activity_data def test_add_activity_correlation_id_not_injected(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 1 activity_data = json.loads(request_mock.last_request_data["data"]) assert constants.CORRELATION_ID_KEY not in activity_data def test_add_activity_correlation_id_injected_from_self(stream_client, stream_feed, request_mock): correlation_id = "corr_id_123" stream_client.set_correlation_id(correlation_id) request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 1 activity_data = json.loads(request_mock.last_request_data["data"]) assert activity_data[constants.CORRELATION_ID_KEY] == correlation_id def test_add_activity_correlation_id_injected_from_flask(stream_client, stream_feed, request_mock, mocker): correlation_id = "corr_id_123" request_mock.add_many(data.ONE_POST_SUCCESS) mocker.patch("getstream_connector.integrations.flask.g", correlation_id=correlation_id) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 1 activity_data = json.loads(request_mock.last_request_data["data"]) assert activity_data[constants.CORRELATION_ID_KEY] == correlation_id @pytest.mark.parametrize("stream_client", ({"throttler": DummyThrottler()},), indirect=True) def test_add_activity_retry_success(stream_client, stream_feed, frozen_time, sleep_mock, request_mock): request_mock.add_many(data.TWO_POST_RATELIMIT_REACHED_ONE_SUCCESS) stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 3 assert sleep_mock.call_args_list == [call(123), call(456)] @pytest.mark.parametrize("stream_client", ({"throttler": DummyThrottler()},), indirect=True) def test_add_activity_retry_fail(stream_client, stream_feed, frozen_time, sleep_mock, request_mock): request_mock.add_many(data.THREE_POST_RATELIMIT_REACHED) with pytest.raises(RateLimitReached): stream_feed.add_activity({"some": "data"}) assert request_mock.call_count == 3 assert sleep_mock.call_args_list == [call(123), call(456)] def test_add_activity_delay_request(stream_client, stream_feed, frozen_time, sleep_mock, request_mock): request_mock.add_many(data.THREE_POST_SUCCESS) stream_feed.add_activity({"some": "data_1"}) stream_feed.add_activity({"some": "data_2"}) stream_feed.add_activity({"some": "data_3"}) assert request_mock.call_count == 3 assert sleep_mock.call_args_list == [call(123), call(456)] def test_add_activities(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 1 request_data = json.loads(request_mock.last_request_data["data"]) activity_data = request_data["activities"] assert len(activity_data) == 2 assert activity_data[0]["some"] == "data_1" assert activity_data[1]["some"] == "data_2" @pytest.mark.parametrize("stream_client", ({"inject_trace_context": False},), indirect=True) def test_add_activities_inject_trace_context_disabled(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 1 request_data = json.loads(request_mock.last_request_data["data"]) activity_data = request_data["activities"] assert constants.TRACE_CONTEXT_KEY not in activity_data[0] assert constants.TRACE_CONTEXT_KEY not in activity_data[1] def test_add_activities_trace_context_injected(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 1 request_data = json.loads(request_mock.last_request_data["data"]) activity_data = request_data["activities"] assert constants.TRACE_CONTEXT_KEY in activity_data[0] assert constants.TRACE_CONTEXT_KEY in activity_data[1] def test_add_activities_correlation_id_not_injected(stream_client, stream_feed, request_mock): request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 1 request_data = json.loads(request_mock.last_request_data["data"]) activity_data = request_data["activities"] assert constants.CORRELATION_ID_KEY not in activity_data[0] assert constants.CORRELATION_ID_KEY not in activity_data[1] @pytest.mark.parametrize("stream_client", ({"inject_correlation_id": False},), indirect=True) def test_add_activities_inject_correlation_id_disabled(stream_client, stream_feed, request_mock): correlation_id = "corr_id_123" stream_client.set_correlation_id(correlation_id) request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 1 request_data = json.loads(request_mock.last_request_data["data"]) activity_data = request_data["activities"] assert constants.CORRELATION_ID_KEY not in activity_data[0] assert constants.CORRELATION_ID_KEY not in activity_data[1] def test_add_activities_correlation_id_injected_from_self(stream_client, stream_feed, request_mock): correlation_id = "corr_id_123" stream_client.set_correlation_id(correlation_id) request_mock.add_many(data.ONE_POST_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 1 request_data = json.loads(request_mock.last_request_data["data"]) activity_data = request_data["activities"] assert activity_data[0][constants.CORRELATION_ID_KEY] == correlation_id assert activity_data[1][constants.CORRELATION_ID_KEY] == correlation_id def test_add_activities_correlation_id_injected_from_flask(stream_client, stream_feed, request_mock, mocker): correlation_id = "corr_id_123" request_mock.add_many(data.ONE_POST_SUCCESS) mocker.patch("getstream_connector.integrations.flask.g", correlation_id=correlation_id) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 1 request_data = json.loads(request_mock.last_request_data["data"]) activity_data = request_data["activities"] assert activity_data[0][constants.CORRELATION_ID_KEY] == correlation_id assert activity_data[1][constants.CORRELATION_ID_KEY] == correlation_id @pytest.mark.parametrize("stream_client", ({"throttler": DummyThrottler()},), indirect=True) def test_add_activities_retry_success(stream_client, stream_feed, frozen_time, sleep_mock, request_mock): request_mock.add_many(data.TWO_POST_RATELIMIT_REACHED_ONE_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 3 assert sleep_mock.call_args_list == [call(123), call(456)] @pytest.mark.parametrize("stream_client", ({"throttler": DummyThrottler()},), indirect=True) def test_add_activities_retry_fail(stream_client, stream_feed, frozen_time, sleep_mock, request_mock): request_mock.add_many(data.THREE_POST_RATELIMIT_REACHED) with pytest.raises(RateLimitReached): stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) assert request_mock.call_count == 3 assert sleep_mock.call_args_list == [call(123), call(456)] def test_add_activities_delay_request(stream_client, stream_feed, frozen_time, sleep_mock, request_mock): request_mock.add_many(data.THREE_POST_SUCCESS) stream_feed.add_activities([{"some": "data_1"}, {"some": "data_2"}]) stream_feed.add_activities([{"some": "data_3"}, {"some": "data_4"}]) stream_feed.add_activities([{"some": "data_5"}, {"some": "data_6"}]) assert request_mock.call_count == 3 assert sleep_mock.call_args_list == [call(123), call(456)]