Akamai WAF SIEM

Use the Akamai WAF SIEM integration to retrieve security events from Akamai Web Application Firewall (WAF) service.

Analytics & SIEM · Akamai WAF SIEM

Details

IDAkamai WAF SIEM
ProviderAkamai Technologies
CategoryAnalytics & SIEM
From Version5.0.0
Docker Imagedemisto/auth-utils:1.0.0.10133006
Supported ModulesAgentix XSIAM

README

Get security event from Akamai Web Application Firewall (WAF) service

This integration was integrated and tested with API version 1.0 of Akamai WAF SIEM.

Use Cases

  • Get security events from Akamai WAF.
  • Analyze security events generated on the Akamai platform and correlate them with security events generated from other sources in Cortex XSOAR.

Detailed Description

A WAF (web application firewall) is a filter that protects against HTTP application attacks. It inspects HTTP traffic before it reaches your application and protects your server by filtering out threats that could damage your site functionality or compromise data.

API keys generating steps

  1. Go to WEB & DATA CENTER SECURITY > Security Configuration > choose your configuration > Advanced settings > Enable SIEM integration.
  2. Open Control panel and login with admin account.
  3. Open identity and access management menu.
  4. Create a user with assigned roles Manage SIEM or make sure the admin has rights to manage SIEM.
  5. Log in to the new account you created in the last step.
  6. Open identity and access management menu.
  7. Create new api client for me.
  8. Assign an API key to the relevant user group, and on the next page assign Read/Write access for SIEM.
  9. Save configuration and go to the API detail you created.
  10. Press new credentials and download or copy it.
  11. Now use the credentials to configure Akamai WAF in Cortex XSOAR.

Configure Akamai WAF SIEM on Cortex XSOAR

  1. Navigate to Settings > Integrations > Servers & Services.
  2. Search for Akamai WAF SIEM.
  3. Click Add instance to create and configure a new integration instance.

    Parameter Required Description
    Server URL (e.g., https://example.net) True  
    Client token False  
    Access token False  
    Client secret False  
    Config ids to fetch True Your Akamai security configuration ID(s). Multiple IDs can be separated by semicolons (e.g., 12345;2345;3456). Config IDs are unique to your account.
    Incident type False  
    First fetch timestamp False  
    Incident fetch limit False The maximum total number of incidents to retrieve per fetch. The maximum is 2000.
    Events fetch limit False The maximum total number of events to retrieve per fetch. The maximum is 80k.
    Skip events decoding False Use this parameter to avoid decoding the http message and attack data fields and speed up the ingestion rate.
    Long running instance False This is a beta feature for high performance fetch events. Use this param only if advised by CS. Make sure this feature is not used with fetch events configured in the integration params and that there’s no config ID used for 2 different instances / features.
    Page Size - high performance mode False The number of events to fetch per request to akamai Default is 200k, maximum is 600k as per Akamai documentation. Use this only when using the long running beta feature.
    Max allowed concurrent tasks False The number of tasks that can run concurrently - the higher the number, the bigger the gap between the ingested events and the events pulled from akamai can be. Maximum is 10k. Use this only when using the long running beta feature.
    Trust any certificate (not secure) False  
    Use system proxy settings False  
  4. Click Test to validate the new instance.

Commands

You can execute these commands from the CLI, as part of a script, or in a playbook.

Fetch Incidents

[
    {
        "name": "Akamai SIEM: 50170",
        "occurred": "2019-12-10T18:28:27Z",
        "rawJSON": {
            "type": "akamai_siem",
            "format": "json",
            "version": "1.0",
            "attackData": {
                "configId": "50170",
                ...
            }
        }
    },
    {
        "name": "Akamai SIEM: 50170",
        "occurred": "2019-12-10T18:28:26Z",
        "rawJSON": {
            "type": "akamai_siem",
            "format": "json",
            "version": "1.0",
            "attackData": {
                "configId": "50170",
                ...
            }
        }
    }
]

akamai-siem-reset-offset


Reset the last offset in case the offset is invalid.

Base Command

akamai-siem-reset-offset

Input

There are no input arguments for this command.

Context Output

There is no context output for this command.

akamai-siem-get-events


Get security events from Akamai WAF.

Base Command

akamai-siem-get-events

Input

Argument Name Description Required
config_ids Unique identifier for each security configuration. To report on more than one configuration, separate the integer identifiers with semicolons (;), for example: 12892;29182;82912. Required
offset This token denotes the last message. If specified, this operation fetches only security events that have occurred from the offset. This is a required parameter for offset mode and you can’t use it in time-based requests Optional
limit Defines the maximum number of security events returned per fetch. Optional
from_epoch The start of a specified time range, expressed in Unix epoch seconds. Optional
to_epoch The end of a specified time range, expressed in Unix epoch seconds. Optional
time_stamp Timestamp of events (<number> <time unit>. For example, 12 hours, 7 days). Optional

Context Output

Path Type Description
Akamai.SIEM.AttackData.clientIP String IP address involved in the attack.
Akamai.SIEM.AttackData.ConfigID String Unique identifier of the security configuration involved.
Akamai.SIEM.AttackData.PolicyID String Unique identifier of the policy configuration involved.
Akamai.SIEM.AttackData.Geo.Asn String Geographic ASN location of the IP address involved in the attack.
Akamai.SIEM.AttackData.Geo.City String City of the IP address involved in the attack.
Akamai.SIEM.AttackData.Geo.Continent String Continent of the IP address involved in the attack.
Akamai.SIEM.AttackData.Geo.Country String Country of the IP address involved in the attack.
Akamai.SIEM.AttackData.Geo.RegionCode String Region code of the IP address involved in the attack.
Akamai.SIEM.AttackData.HttpMessage.Bytes Number HTTP messege size in bytes.
Akamai.SIEM.AttackData.HttpMessage.Host String HTTP messege host.
Akamai.SIEM.AttackData.HttpMessage.Method String HTTP messege method.
Akamai.SIEM.AttackData.HttpMessage.Path String HTTP messege path.
Akamai.SIEM.AttackData.HttpMessage.Port String HTTP messege port.
Akamai.SIEM.AttackData.HttpMessage.Protocol String HTTP messege protocol.
Akamai.SIEM.AttackData.HttpMessage.Query String HTTP messege query.
Akamai.SIEM.AttackData.HttpMessage.RequestHeaders String HTTP messege request headers.
Akamai.SIEM.AttackData.HttpMessage.RequestID String HTTP messege request ID.
Akamai.SIEM.AttackData.HttpMessage.ResponseHeaders String HTTP message response headers.
Akamai.SIEM.AttackData.HttpMessage.Start Date HTTP messege epoch start time.
Akamai.SIEM.AttackData.HttpMessage.Status Number HTTP messege status code.
IP.Address String IP address.
IP.ASN String The autonomous system name for the IP address, for example: “AS8948”.
IP.Geo.Country String The country in which the IP address is located.
Context Example
{
  "Akamai": {
    "SIEM": [
        {
            "AttackData": {
                "ClientIP": "8.8.8.8",
                "ConfigID": "50170",
                "PolicyID": "1234_89452",
                "RuleActions": [
                    "alert",
                    "deny"
                ],
                "RuleMessages": [
                    "Custom_RegEX_Rule",
                    "No Accept Header AND No User Agent Header"
                ],
                "RuleTags": [
                    "example",
                    "No-AH-UA"
                ],
                "Rules": [
                    "642118",
                    "642119"
                ]
            },
            "Geo": {
                "Asn": "16509",
                "City": "FRANKFURT",
                "Continent": "EU",
                "Country": "DE",
                "RegionCode": "HE"
            },
            "HttpMessage": {
                "Bytes": "296",
                "Host": "wordpress.panw.ninja",
                "Method": "POST",
                "Path": "/wp-cron.php",
                "Port": "80",
                "Protocol": "HTTP/1.1",
                "RequestHeaders": "Host",
                "RequestId": "87bb604",
                "ResponseHeaders": "Server",
                "Start": "1576746102",
                "Status": "403"
            }
        },
        {
            "AttackData": {
                "ClientIP": "8.8.8.8",
                "ConfigID": "50170",
                "PolicyID": "1234_89452",
                "RuleActions": [
                    "alert",
                    "deny"
                ],
                "RuleMessages": [
                    "Custom_RegEX_Rule",
                    "No Accept Header AND No User Agent Header"
                ],
                "RuleTags": [
                    "example",
                    "No-AH-UA"
                ],
                "Rules": [
                    "642118",
                    "642119"
                ]
            },
            "Geo": {
                "Asn": "16509",
                "City": "FRANKFURT",
                "Continent": "EU",
                "Country": "DE",
                "RegionCode": "HE"
            },
            "HttpMessage": {
                "Bytes": "296",
                "Host": "wordpress.panw.ninja",
                "Method": "POST",
                "Path": "/wp-cron.php",
                "Port": "80",
                "Protocol": "HTTP/1.1",
                "RequestHeaders": "Header",
                "RequestId": "32e63ee2",
                "ResponseHeaders": "Server",
                "Start": "1576746179",
                "Status": "403"
            }
        }
    ]
  },
  "IP": [
    {
      "ASN": "5650",
      "Address": "8.8.8.8",
      "Geo": {
        "Country": "US"
      }
    },
    {
      "ASN": "5650",
      "Address": "8.8.8.8",
      "Geo": {
        "Country": "US"
      }
    }
  ]
}

Troubleshooting

receiving 416 error code / aggregated delay when fetching events

This may be due to not querying for enough events per interval / request.
The proposed solution in that case is to increase the Events fetch limit parameter.
Events fetch limit is the number of total events we want to retrieve each fetch interval. Note that the maximum allowed value is 80k.
Note that in cases where the ingestion rate from the Akamai API is higher, the integration will detect it and trigger the next fetch immediately.

A single fetch interval may execute multiple requests, each retrieving up to 20k events per request.

If after readjusting the limit you keep encountering errors, please contact support.

Known limitations

The config ID can only be configured on one instance

Due to limitations from Akamai, the config ID can only be configured on one instance on the same machine or on different machines (i.e. the same config ID can’t be configured both on dev and prod tenants or twice on the same tenant).
Configuring on multiple machines may lead to duplications or missing events.

<~PLATFORM>

License Requirements

The following configuration parameters require the Cortex XSIAM license:

  • Fetch Events
  • Skip events decoding
  • Long running instance

</~PLATFORM>

Configuration parameters

  • host — Server URL (e.g., https://akaa-xxxxxxxxxxxxxxxx-xxxxxxxxxxxxxxxx.luna.akamaiapis.net) (required)
  • clientToken — Client token
  • clienttoken_creds
  • accessToken — Access token
  • accesstoken_creds
  • clientSecret — Client secret
  • clientsecret_creds
  • configIds — Config IDs to fetch (required)
  • incidentType — Incident type
  • fetchTime — First fetch timestamp (<number> <time unit>, e.g., 12 hours, 7 days)
  • fetchLimit — Incident fetch limit
  • eventsFetchLimit — Events fetch limit
  • isFetch — Fetch incidents
  • insecure — Trust any certificate (not secure)
  • proxy — Use system proxy settings
  • incidentFetchInterval — Incidents Fetch Interval
  • eventFetchInterval — Events Fetch Interval
  • isFetchEvents — Fetch Events
  • should_skip_decode_events — Skip events decoding
  • longRunning — Long running instance
  • beta_page_size — Page Size - high performance mode
  • max_concurrent_tasks — Max allowed concurrent tasks

Commands (2)

  • akamai-siem-get-events

    Get security events from Akamai WAF.

  • akamai-siem-reset-offset

    Reset the last offset to start fetching from the first fetch timestamp, use this command when the instance is disabled.

"""Imports"""

# STD packages
import asyncio
from datetime import datetime
import re
import time
import json

import demistomock as demisto

# 3-rd party packages
import pytest
from freezegun import freeze_time

# Local imports
import Akamai_SIEM
from CommonServerPython import urljoin, DemistoException


"""Helper functions and fixrtures"""
BASE_URL = urljoin("https://akab-hnanog6ge5or6biz-ukavvo4zvqliqhlw.cloudsecurity.akamaiapis.net", "/siem/v1/configs")
with open("./Akamai_SIEM_test/TestCommandsFunctions/sec_events_empty.txt") as sec_events_empty:
    SEC_EVENTS_EMPTY_TXT = sec_events_empty.read()
with open("./Akamai_SIEM_test/TestCommandsFunctions/sec_events.txt") as sec_events:
    SEC_EVENTS_TXT = sec_events.read()
with open("./Akamai_SIEM_test/TestCommandsFunctions/sec_events_six_results.txt") as sec_events_six_results:
    SEC_EVENTS_SIX_RESULTS_TXT = sec_events_six_results.read()
with open("./Akamai_SIEM_test/TestCommandsFunctions/sec_events_two_results.txt") as sec_events_two_results:
    SEC_EVENTS_TWO_RESULTS_TXT = sec_events_two_results.read()


def load_params_from_json(json_path, type=""):
    with open(json_path) as f:
        file = json.load(f)
        if type == "incidents":
            for incident in file:
                incident["rawJSON"] = json.dumps(incident.get("rawJSON", {}))
    return file


@pytest.fixture(scope="module")
def client():
    from Akamai_SIEM import Client

    return Client(base_url=BASE_URL)


"""Tests"""


@pytest.mark.commands
@freeze_time(time.ctime(1576009202))
class TestCommandsFunctions:
    @pytest.mark.fetch
    def test_fetch_incidents_command_1(self, client, datadir, requests_mock):
        """Test - No last time exsits and event available"""
        from Akamai_SIEM import fetch_incidents_command

        requests_mock.get(f"{BASE_URL}/50170?limit=5&from=1575966002", text=SEC_EVENTS_TXT)
        tested_incidents, tested_last_run = fetch_incidents_command(
            client=client, fetch_time="12 hours", fetch_limit=5, config_ids="50170", last_run={}
        )
        expected_incidents = load_params_from_json(datadir["expected_fetch.json"], type="incidents")
        expected_last_run = {"lastRun": "1576002507"}
        assert expected_incidents == tested_incidents, "Incidents - No last time exsits and event available"
        assert tested_last_run == expected_last_run, "Last run - No last time exsits and event available"

    @pytest.mark.fetch
    def test_fetch_incidents_command_2(self, client, datadir, requests_mock):
        """Test - Last time exsits and events available"""
        from Akamai_SIEM import fetch_incidents_command

        requests_mock.get(f"{BASE_URL}/50170?from=1575966002&limit=5", text=SEC_EVENTS_TXT)
        tested_incidents, tested_last_run = fetch_incidents_command(
            client=client, fetch_time="12 hours", fetch_limit="5", config_ids="50170", last_run="1575966002"
        )
        expected_incidents = load_params_from_json(datadir["expected_fetch.json"], type="incidents")
        expected_last_run = {"lastRun": "1576002507"}
        assert expected_incidents == tested_incidents, "Incidents - Last time exsits and events available"
        assert tested_last_run == expected_last_run, "Last run - No last time exsits and event available"

    @pytest.mark.fetch
    def test_fetch_incidents_command_3(self, client, datadir, requests_mock):
        """Test - Last time exsits and no available data"""
        from Akamai_SIEM import fetch_incidents_command

        requests_mock.get(f"{BASE_URL}/50170?from=1575966002&limit=5", text=SEC_EVENTS_EMPTY_TXT)
        tested_incidents, tested_last_run = fetch_incidents_command(
            client=client, fetch_time="12 hours", fetch_limit=5, config_ids="50170", last_run="1575966002"
        )
        expected_last_run = {"lastRun": "1575966002"}
        expected_incidents = []
        assert expected_incidents == tested_incidents, "Incidents - Last time exsits and no available data"
        assert tested_last_run == expected_last_run, "Last run - No last time exsits and event available"

    @pytest.mark.fetch
    def test_fetch_incidents_command_4(self, client, datadir, requests_mock):
        """Test - No last time exsits and no available data"""
        from Akamai_SIEM import fetch_incidents_command

        requests_mock.get(f"{BASE_URL}/50170?from=1575966002&limit=5", text=SEC_EVENTS_EMPTY_TXT)
        tested_incidents, tested_last_run = fetch_incidents_command(
            client=client, fetch_time="12 hours", fetch_limit=5, config_ids="50170", last_run={}
        )
        expected_last_run = {"lastRun": "1575966002"}
        expected_incidents = []
        assert expected_incidents == tested_incidents, "Incidents - No last time exsits and no available data"
        assert tested_last_run == expected_last_run, "Last run - No last time exsits and no available data"

    @pytest.mark.get_events
    def test_get_events_command_1(self, client, datadir, requests_mock):
        """Test query response without security events - check only enrty context"""
        from Akamai_SIEM import get_events_command

        requests_mock.get(f"{BASE_URL}/50170?from=1575966002&limit=5", text=SEC_EVENTS_EMPTY_TXT)
        # About the drop some mean regex right now disable-secrets-detection-start
        human_readable, entry_context_tested, raw_response = get_events_command(
            client=client, config_ids="50170", from_epoch="1575966002", limit="5"
        )
        # Drops the mic disable-secrets-detection-end

        assert entry_context_tested == {}, "Test query response without security events - check only enrty context"

    @pytest.mark.get_events
    def test_get_events_command_2(self, client, datadir, requests_mock):
        """Test query response with security events - check only entry context"""
        from Akamai_SIEM import get_events_command

        # About the drop some mean regex right now disable-secrets-detection-start
        requests_mock.get(f"{BASE_URL}/50170?from=1575966002&limit=5", text=SEC_EVENTS_TXT)
        human_readable, entry_context_tested, raw_response = get_events_command(
            client=client, config_ids="50170", from_epoch="1575966002", limit="5"
        )
        # Drops the mic disable-secrets-detection-end
        expected_ec = load_params_from_json(json_path=datadir["get_events_expected_ec_2.json"])

        assert entry_context_tested == expected_ec, "Test query response with security events - check only entry context"

    def test_fetch_events_command_with_break_before_timeout(self, client, mocker):
        """
        Given:
        - A client object
        - 2 mock responses each one with one has 50 events (total 100).
        When:
        - Calling fetch_events_command() and getting is_interval_doesnt_have_enough_time_to_run in the second execution.
        Then:
        - Ensure there are only 50 total events received and auto_trigger_next_run = True.
        """
        page_size = 50
        events = [
            (
                [{"id": i + 1, "httpMessage": {"start": i + 1}} for i in range(page_size * j, page_size * (j + 1))],
                f"offset_{page_size * (j + 1)}",
            )
            for j in range(2)
        ]
        mocker.patch.object(Akamai_SIEM.Client, "get_events_with_offset", side_effect=events)
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        total_events_count = 0
        for _events, _, total_events_count, auto_trigger_next_run in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "3 days",
            220,
            "",
            {},
            50,
            False,
        ):
            mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(True, 1))
        assert total_events_count == 50
        assert auto_trigger_next_run

    def test_fetch_events_command_with_break_for_page_too_small(self, client, mocker):
        """
        Given:
        - A client object
        - 2 mock responses each one with one has 50 events (total 100).
        When:
        - Calling fetch_events_command() with page_size > amount of events obtained in first execution.
        Then:
        - Ensure there are only 50 total events received.
        """
        page_size = 50
        events = [
            (
                [{"id": i + 1, "httpMessage": {"start": i + 1}} for i in range(page_size * j, page_size * (j + 1))],
                f"offset_{page_size * (j + 1)}",
            )
            for j in range(2)
        ]
        mocker.patch.object(Akamai_SIEM.Client, "get_events_with_offset", side_effect=events)
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        total_events_count = 0
        for events, _, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "3 days",
            220,
            "",
            {},
            page_size=60,
            should_skip_decode_events=False,
        ):
            pass
        assert total_events_count == 50

    def test_fetch_events_command_sanity(self, client, mocker):
        """
        Given:
        - A client object
        - 500 events to pull in the 3rd party
        - A fetch_limit of 260
        When:
        - Calling fetch_events_command()
        Then:
        - Ensure offset is updated in each iteration by checking its value
        - Ensure 250 events are pulled (fetch_limit, rounded up to the nearest multiple of page_size=50)
        """
        num_of_results = 500
        page_size = 50
        limit = 250
        num_of_pages = num_of_results // page_size
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        mocker.patch.object(
            Akamai_SIEM.Client,
            "get_events_with_offset",
            side_effect=[
                (
                    [{"id": i + 1, "httpMessage": {"start": i + 1}} for i in range(page_size * j, page_size * (j + 1))],
                    f"offset_{page_size * (j + 1)}",
                )
                for j in range(num_of_pages)
            ],
        )
        total_events_count = 0

        for events, offset, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "3 days",
            limit,
            "",
            {},
            page_size,
            False,
        ):
            assert offset == f"offset_{events[-1]['id']}" if events else True
        assert total_events_count == 250

    def test_fetch_events_command_no_results(self, mocker, client, requests_mock):
        """
        Given:
        - A client object
        - no events to pull from the 3rd party
        - offset is 11111
        When:
        - Calling fetch_events_command()
        Then:
        - Ensure no events are returned and the offset is the same
        """
        # Historically this used the removed FETCH_EVENTS_MAX_PAGE_SIZE constant purely as a page-size
        # value. page_size is now passed explicitly to fetch_events_command, so use its former value.
        size = 20000

        total_events_count = 0
        last_offset = "11111"
        requests_mock.get(f"{BASE_URL}/50170?limit={size}&offset={last_offset}", text=SEC_EVENTS_EMPTY_TXT)
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))

        for _, offset, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "12 hours",
            size,
            "50170",
            {"offset": last_offset},
            size,
            False,
        ):
            last_offset = offset
        assert total_events_count == 0
        assert last_offset == "318d8"

    def test_attach_offset_param_to_url(self, client):
        """Test that the execute_get_events_request function doesn't encode the offset's ; into %3B."""
        params = {"offset": "das2;test", "limit": 100}
        try:
            client.execute_get_events_request(params, 50170)
        except Exception as e:
            error_str = str(e.args[1])
        assert "%3B" not in error_str

    def test_fetch_events_command_limit_is_smaller_than_page_size(self, client, requests_mock, mocker):
        """
        Given:
        - A client object
        - 8 events to pull from the 3rd party
        - page size is 6
        - limit is 6
        When:
        - Calling fetch_events_command()
        Then:
        - Ensure 6 events are returned
        """
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        total_events_count = 0
        last_offset = None
        requests_mock.get(f"{BASE_URL}/50170?limit=6&from=1575966002", text=SEC_EVENTS_SIX_RESULTS_TXT)
        requests_mock.get(f"{BASE_URL}/50170?limit=6&from=1575966002&offset=218d9", text=SEC_EVENTS_TXT)
        requests_mock.get(f"{BASE_URL}/50170?limit=6&from=1575966002&offset=318d8", text=SEC_EVENTS_EMPTY_TXT)

        for _, offset, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "12 hours",
            6,
            "50170",
            {},
            6,
            False,
        ):
            last_offset = offset
        assert total_events_count == 6
        assert last_offset == "218d9"

    def test_fetch_events_command_limit_is_higher_than_page_size(self, client, requests_mock, mocker):
        """
        Given:
        - A client object
        - 8 events to pull from the 3rd party
        - page size is 6
        - limit is 20
        When:
        - Calling fetch_events_command()
        Then:
        - Ensure 8 events are returned
        """
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        total_events_count = 0
        last_offset = None
        requests_mock.get(f"{BASE_URL}/50170?limit=6&from=1575966002", text=SEC_EVENTS_SIX_RESULTS_TXT)
        requests_mock.get(f"{BASE_URL}/50170?limit=6&offset=218d9", text=SEC_EVENTS_TXT)
        requests_mock.get(f"{BASE_URL}/50170?limit=6&offset=318d8", text=SEC_EVENTS_EMPTY_TXT)

        for _, offset, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "12 hours",
            20,
            "50170",
            {},
            6,
            False,
        ):
            last_offset = offset
        assert total_events_count == 8
        assert last_offset == "318d8"

    def test_fetch_events_command_limit_reached(self, client, requests_mock, mocker):
        """
        Given:
        - A client object
        - 4 events to pull from the 3rd party
        - page size is 2
        - limit is 2
        When:
        - Calling fetch_events_command()
        Then:
        - Ensure 2 events are returned
        - Ensure last_offset is the one returned from the last page we pulled events from (the 1st one)
        """
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        total_events_count = 0
        last_offset = None
        requests_mock.get(f"{BASE_URL}/50170?limit=2&from=1575966002", text=SEC_EVENTS_TWO_RESULTS_TXT)
        requests_mock.get(f"{BASE_URL}/50170?limit=2&offset=117d9", text=SEC_EVENTS_TXT)

        for _, offset, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "12 hours",
            2,
            "50170",
            {},
            2,
            False,
        ):
            last_offset = offset
        assert total_events_count == 2
        assert last_offset == "117d9"

    def test_fetch_events_command_with_page_truncated(self, mocker, client, requests_mock):
        """
        Given:
        - A client object
        - page_size = 2, fetch_limit = 3, and two requests_mock.
        When:
        - Calling fetch_events_command()
        Then:
        - The request was called correctly in the first fetch_events_command execution with limit = 2 and from_time.
        - The request was called correctly in the second fetch_events_command execution with limit = 1 and offset.
        - A total of 3 events received with offset = the offset from the second response.
        """
        page_size = 2
        fetch_limit = 3
        first_response_mock = '{"id": 1, "httpMessage": {"start": 1}}\n{"id": 2, "httpMessage": {"start": 2}}\n{"offset": "a"}'
        second_response_mock = '{"id": 3, "httpMessage": {"start": 3}}\n{"offset": "b"}'
        mocker.patch("CommonServerPython.parse_date_range", return_value="1575966002")
        requests_mock.get(f"{BASE_URL}/50170?limit=2&from=1575750002", text=first_response_mock)
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        total_events_count = 0
        for _, offset, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            fetch_time="3 days",
            fetch_limit=fetch_limit,
            config_ids="50170",
            ctx={},
            page_size=page_size,
            should_skip_decode_events=False,
        ):
            requests_mock.get(f"{BASE_URL}/50170?limit=1&offset={offset}", text=second_response_mock)
        assert total_events_count == fetch_limit
        assert offset == "b"

    def test_response_error_non_416_raises(self, mocker, client):
        """
        Given:
        - A client object that raises a non-416 error (e.g. 403 Unauthorized) from get_events_with_offset.
        When:
        - Calling fetch_events_command.
        Then:
        - Ensure the error is re-raised (non-recoverable errors are not swallowed).
        """
        error_entry = {
            "clientIp": "192.0.2.85",
            "detail": "The specified user is unauthorized to access the requested data",
            "instance": "https://test.akamaiapis.net/siem/v1/configs=12345?offset=123",
            "method": "GET",
            "requestId": "9cf2274",
            "requestTime": "2023-06-20T15:01:11Z",
            "serverIp": "1.1.1.1",
            "title": "Unauthorized",
        }
        err_msg = f"Error in API call [403] - Unauthorized\n{json.dumps(error_entry)}"
        mocker.patch.object(Akamai_SIEM.Client, "get_events_with_offset", side_effect=DemistoException(err_msg, res={}))
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=False)
        mocker.patch.object(demisto, "error")
        mocker.patch.object(demisto, "debug")
        with pytest.raises(DemistoException) as e:
            for _, _, _, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
                client,
                "3 days",
                220,
                "",
                {},
                5000,
                False,
            ):
                pass
        assert "Unauthorized" in str(e)

    def test_response_error_416_recovers_in_run(self, mocker, client):
        """
        Given:
        - A client whose first offset-based request fails with a 416 (expired/out-of-range offset),
          and whose second (recovery) request succeeds and returns events.
        When:
        - Calling fetch_events_command with a stored offset.
        Then:
        - Ensure the 416 does NOT raise; instead the offset is reset, the fetch restarts from the safe
          lookback window, and the events from the recovery request are yielded.
        """
        err_msg = "Error in API call [416] - Requested Range Not Satisfiable"
        recovery_events = ['{"event": "1"}', '{"offset": "recovered_offset"}']
        mocker.patch.object(
            Akamai_SIEM.Client,
            "get_events_with_offset",
            side_effect=[
                DemistoException(err_msg, res={}),  # first call: stale offset -> 416
                (recovery_events[:-1], "recovered_offset"),  # recovery call: time-based fetch succeeds
                ([], "recovered_offset"),  # loop termination
            ],
        )
        mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
        reset_offset_mock = mocker.patch.object(Akamai_SIEM, "reset_offset_command")
        mocker.patch.object(demisto, "error")
        mocker.patch.object(demisto, "debug")
        mocker.patch.object(demisto, "info")

        collected = []
        for events, _, _, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "3 days",
            220,
            "",
            {"offset": "stale_offset"},
            5000,
            True,  # should_skip_decode_events, so events pass through unchanged
        ):
            if events:
                collected.extend(events)

        # The stale offset was reset as part of recovery.
        reset_offset_mock.assert_called_once()
        # The recovery request's events were yielded (no exception raised).
        assert collected == recovery_events[:-1]


@pytest.mark.parametrize(
    "error, expected",
    [
        (DemistoException("Error in API call [416] - Requested Range Not Satisfiable"), True),
        (DemistoException("Requested Range Not Satisfiable"), True),
        (DemistoException("Error in API call [403] - Unauthorized"), False),
        (DemistoException("Some other error"), False),
    ],
)
def test_is_offset_out_of_range_error(error, expected):
    """
    Given:
    - Various Akamai errors (416 offset-expired and non-416).
    When:
    - Calling is_offset_out_of_range_error.
    Then:
    - Only the '416 Requested Range Not Satisfiable' errors are identified as offset-out-of-range.
    """
    assert Akamai_SIEM.is_offset_out_of_range_error(error) is expected


@pytest.mark.parametrize("reset_context_offset", [True, False])
def test_handle_offset_out_of_range(mocker, reset_context_offset):
    """
    Given:
    - A 416 offset-out-of-range error.
    When:
    - Calling handle_offset_out_of_range with reset_context_offset True (fetch-events) or False (long-running).
    Then:
    - The context offset is reset only when reset_context_offset is True.
    - The returned from_epoch corresponds to AKAMAI_MAX_LOOKBACK_MINUTES ago (within Akamai's 12h limit).
    """
    reset_offset_mock = mocker.patch.object(Akamai_SIEM, "reset_offset_command")
    mocker.patch.object(demisto, "error")
    mocker.patch.object(demisto, "info")
    error = DemistoException("Error in API call [416] - Requested Range Not Satisfiable")

    from_epoch = Akamai_SIEM.handle_offset_out_of_range(error, reset_context_offset=reset_context_offset)

    if reset_context_offset:
        reset_offset_mock.assert_called_once()
    else:
        reset_offset_mock.assert_not_called()

    # from_epoch is a Unix-epoch seconds string within the safe lookback window (just under 12h).
    now_epoch = int(time.time())
    expected_epoch = now_epoch - Akamai_SIEM.AKAMAI_MAX_LOOKBACK_MINUTES * 60
    assert isinstance(from_epoch, str)
    # Allow a small delta for execution time between computing expected and actual.
    assert abs(int(from_epoch) - expected_epoch) <= 5
    # The recovery window must stay under Akamai's hard 12-hour limit.
    assert (now_epoch - int(from_epoch)) < 12 * 60 * 60


@pytest.mark.parametrize("msg", ["", None])
def test_decode_message_empty_string(msg):
    """
    Given:
    - An empty / falsy message.
    When:
    - Calling decode_message.
    Then:
    - An empty list is returned (guard clause), without attempting to base64-decode.
    """
    assert Akamai_SIEM.decode_message(msg) == []


@pytest.mark.parametrize("headers", ["", None])
def test_decode_url_empty_string(headers):
    """
    Given:
    - Empty / falsy headers.
    When:
    - Calling decode_url.
    Then:
    - An empty dict is returned (guard clause), without attempting to parse.
    """
    assert Akamai_SIEM.decode_url(headers) == {}


def test_decode_event_decodes_attack_data_and_http_headers():
    """
    Given:
    - A raw JSON event string containing an attackData section (with base64 rule fields) and an
      httpMessage section (with URL-encoded request/response headers).
    When:
    - Calling decode_event on the raw string.
    Then:
    - The attackData rule fields are base64/URL decoded into lists, and the httpMessage headers are
      decoded into dicts. This characterizes the shared decode logic previously duplicated in
      fetch_events_command and process_and_send_events_to_xsiam.
    """
    request_headers = "Content-Type%3A%20application/json%3Bcharset%3DUTF-8%0Auser%3A%20test%40test.com"
    raw_event = json.dumps(
        {
            "attackData": {"rules": "cnVsZTE=", "ruleMessages": "bXNn"},
            "httpMessage": {"requestHeaders": request_headers, "responseHeaders": ""},
        }
    )

    decoded = Akamai_SIEM.decode_event(raw_event)

    assert decoded["attackData"]["rules"] == ["rule1"]
    assert decoded["attackData"]["ruleMessages"] == ["msg"]
    assert decoded["httpMessage"]["requestHeaders"] == {
        "Content_Type": "application/json;charset=UTF-8",
        "user": "test@test.com",
    }
    assert decoded["httpMessage"]["responseHeaders"] == {}


def test_decode_event_returns_raw_string_on_malformed_json():
    """
    Given:
    - A malformed JSON event string that cannot be parsed.
    When:
    - Calling decode_event on it.
    Then:
    - The original raw string is returned unchanged (no exception), matching the existing
      "leave malformed event in place" behavior.
    """
    malformed = "{not-valid-json"

    assert Akamai_SIEM.decode_event(malformed) == malformed


def test_decode_event_without_attack_or_http_sections_is_unchanged():
    """
    Given:
    - A valid JSON event with neither attackData nor httpMessage sections.
    When:
    - Calling decode_event on it.
    Then:
    - The event is parsed to a dict but its contents are otherwise unchanged.
    """
    raw_event = json.dumps({"id": 42, "geo": {"country": "US"}})

    decoded = Akamai_SIEM.decode_event(raw_event)

    assert decoded == {"id": 42, "geo": {"country": "US"}}


def test_events_to_ec_builds_expected_entry_context():
    """
    Given:
    - A single raw event containing attackData (with base64 rule fields), httpMessage, and geo sections.
    When:
    - Calling events_to_ec.
    Then:
    - The returned entry-context, IP-context, and human-readable structures match the expected shapes.
      This pins the exact output so the internal dict-lookup hoisting refactor is provably byte-identical.
    """
    raw_event = {
        "attackData": {
            "configId": "50170",
            "policyId": "pol1",
            "clientIP": "1.2.3.4",
            "rules": "cnVsZTE=",
            "ruleMessages": "bXNn",
            "ruleActions": "ZGVueQ==",
        },
        "httpMessage": {
            "requestId": "req1",
            "start": "1488816442",
            "method": "GET",
            "host": "example.com",
            "status": "403",
        },
        "geo": {"continent": "NA", "country": "US", "city": "SF", "asn": "123"},
    }

    events_ec, ip_ec, human_readable = Akamai_SIEM.events_to_ec([raw_event])

    assert events_ec == [
        {
            "AttackData": {
                "ConfigID": "50170",
                "PolicyID": "pol1",
                "ClientIP": "1.2.3.4",
                "Rules": ["rule1"],
                "RuleMessages": ["msg"],
                "RuleActions": ["deny"],
            },
            "HttpMessage": {
                "RequestId": "req1",
                "Start": "1488816442",
                "Method": "GET",
                "Host": "example.com",
                "Status": "403",
            },
            "Geo": {"Continent": "NA", "Country": "US", "City": "SF", "Asn": "123"},
        }
    ]
    assert ip_ec == [{"Address": "1.2.3.4", "ASN": "123", "Geo": {"Country": "US"}}]
    assert human_readable == [
        {
            "Attacking IP": "1.2.3.4",
            "Config ID": "50170",
            "Policy ID": "pol1",
            "Rules": ["rule1"],
            "Rule messages": ["msg"],
            "Rule actions": ["deny"],
            "Date occured": Akamai_SIEM.date_format_converter(from_format="epoch", date_before="1488816442"),
            "Location": {"Country": "US", "City": "SF"},
        }
    ]


def test_fetch_events_command_decodes_in_place_and_counts_last_page_size(client, mocker):
    """
    Given:
    - A single page of raw JSON events (with attackData and httpMessage) followed by an empty page.
    When:
    - Calling fetch_events_command with should_skip_decode_events=False.
    Then:
    - The events are decoded in place (each list item becomes a dict), and the yielded total_events_count
      equals the number of events in the page (last_page_size accounting).
    """
    raw_page = [
        json.dumps({"attackData": {"rules": "cnVsZTE="}, "httpMessage": {"requestHeaders": "", "responseHeaders": ""}}),
        json.dumps({"attackData": {"rules": "cnVsZTI="}, "httpMessage": {"requestHeaders": "", "responseHeaders": ""}}),
    ]
    mocker.patch.object(
        Akamai_SIEM.Client,
        "get_events_with_offset",
        side_effect=[(list(raw_page), "off1"), ([], "off1")],
    )
    mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))

    collected = []
    last_total = 0
    for events, _, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
        client, "3 days", 220, "50170", {}, 5000, False
    ):
        if events:
            collected = events
            last_total = total_events_count

    assert last_total == len(raw_page)
    # Every event was decoded from a raw JSON string into a dict in place.
    assert all(isinstance(event, dict) for event in collected)


def test_fetch_events_command_skip_decode_keeps_raw_strings(client, mocker):
    """
    Given:
    - A page of raw JSON event strings.
    When:
    - Calling fetch_events_command with should_skip_decode_events=True.
    Then:
    - The events are yielded as-is (raw strings), without being decoded to dicts.
    """
    raw_page = ['{"attackData": {"rules": "cnVsZTE="}}', '{"attackData": {"rules": "cnVsZTI="}}']
    mocker.patch.object(
        Akamai_SIEM.Client,
        "get_events_with_offset",
        side_effect=[(list(raw_page), "off1"), ([], "off1")],
    )
    mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))

    collected = []
    for events, _, _, _ in Akamai_SIEM.fetch_events_command(client, "3 days", 220, "50170", {}, 5000, True):  # noqa: B007
        if events:
            collected = events

    assert collected == raw_page
    assert all(isinstance(event, str) for event in collected)


def test_fetch_events_command_malformed_json_left_in_place(client, mocker):
    """
    Given:
    - A page containing a malformed JSON event that cannot be decoded.
    When:
    - Calling fetch_events_command with should_skip_decode_events=False.
    Then:
    - The malformed event is left in place (as its original raw string) instead of crashing the fetch,
      and it is still counted in total_events_count.
    """
    malformed = "{not-valid-json"
    valid = json.dumps({"attackData": {"rules": "cnVsZTE="}, "httpMessage": {"requestHeaders": "", "responseHeaders": ""}})
    raw_page = [valid, malformed]
    mocker.patch.object(
        Akamai_SIEM.Client,
        "get_events_with_offset",
        side_effect=[(list(raw_page), "off1"), ([], "off1")],
    )
    mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))

    collected = []
    last_total = 0
    for events, _, total_events_count, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
        client, "3 days", 220, "50170", {}, 5000, False
    ):
        if events:
            collected = events
            last_total = total_events_count

    assert last_total == len(raw_page)
    # Valid event decoded to dict; malformed event kept as its original raw string.
    assert isinstance(collected[0], dict)
    assert collected[1] == malformed


@pytest.mark.parametrize(
    "header",
    [
        (
            "Content-Type%3A%20application/json%3Bcharset%3DUTF-8%0D%0Auser%3A%20test%40test.com%0D%0Aclient%3A%"
            "20test_client%0D%0AX-Kong-Upstream-Latency%3A%2066%0D%0AX-Kong-Proxy-Latency%3A%202%0D%0AX-Kong-Request-Id%3A%20X"
            "_request_id%0D%0AEPM-Request-ID%3A%20EPM_request_id%0D%0AContent-Length%3A%20157%0D%0ADate%3A%20Mon%2C%2025%20Mar"
            "%202024%2013%3A52%3A11%20GMT%0D%0AConnection%3A%20keep-alive%0D%0AServer-Timing%3A%20cdn-cache%3B%20desc%3DMISS%0"
            "D%0AServer-Timing%3A%20edge%3B%20dur%3D23%0D%0AServer-Timing%3A%20origin%3B%20dur%3D72%0D%0AServer-Timing%3A%20int"
            "id%3Bdesc%3Ddd%0D%0AStrict-Transport-Security%3A%20max-age%3D31536000%20%3B%20includeSubDomains%20%3B%20preload%0D"
            "%0A"
        ),
        (
            "Content-Type%3A%20application/json%3Bcharset%3DUTF-8%0Auser%3A%20test%40test.com%0Aclient%3A%20"
            "test_client%0AX-Kong-Upstream-Latency%3A%2066%0AX-Kong-Proxy-Latency%3A%202%0AX-Kong-Request-Id%3A%20X_request_id%"
            "0AEPM-Request-ID%3A%20EPM_request_id%0AContent-Length%3A%20157%0ADate%3A%20Mon%2C%2025%20Mar%202024%2013%3A52%3A11"
            "%20GMT%0AConnection%3A%20keep-alive%0AServer-Timing%3A%20cdn-cache%3B%20desc%3DMISS%0AServer-Timing%3A%20edge%3B"
            "%20dur%3D23%0AServer-Timing%3A%20origin%3B%20dur%3D72%0AServer-Timing%3A%20intid%3Bdesc%3Ddd%0A"
            "Strict-Transport-Security%3A%20max-age%3D31536000%20%3B%20includeSubDomains%20%3B%20preload%0A"
        ),
    ],
)
def test_decode_url(header):
    """
    Given: A url decoded string.
        - Case 1: Each key separated by '\r\n'.
        - Case 2: Each key is separated by '\n'.
    When: Calling Akamai_SIEM.decode_url.
    Then: Ensure that the dict was decoded correctly and the same dict was extracted in both cases.
    """
    expected_decoded_dict = {
        "Content_Type": "application/json;charset=UTF-8",
        "user": "test@test.com",
        "client": "test_client",
        "X_Kong_Upstream_Latency": "66",
        "X_Kong_Proxy_Latency": "2",
        "X_Kong_Request_Id": "X_request_id",
        "EPM_Request_ID": "EPM_request_id",
        "Content_Length": "157",
        "Date": "Mon, 25 Mar 2024 13:52:11 GMT",
        "Connection": "keep-alive",
        "Server_Timing": "intid;desc=dd",
        "Strict_Transport_Security": "max-age=31536000 ; includeSubDomains ; preload",
    }
    assert Akamai_SIEM.decode_url(header) == expected_decoded_dict


@pytest.mark.parametrize(
    "freeze_mock, min_allowed_delta, worst_case_time, expected_time, expected_should_break",
    [
        (datetime(2024, 4, 10, 10, 4, 10), 30, 0, 250, True),
        (datetime(2024, 4, 10, 10, 4, 10), 310, 50, 50, True),
        (datetime(2024, 4, 10, 10, 1, 10), 30, 50, 50, False),
        (datetime(2024, 4, 10, 10, 1, 10), 30, 0, 70, False),
    ],
)
def test_is_interval_doesnt_have_enough_time_to_run(
    mocker, freeze_mock, min_allowed_delta, worst_case_time, expected_time, expected_should_break
):
    """
    Given: min_allowed_delta
        - Case 1: min_allowed_delta = 30, no worst_case_time set yet, and 50 seconds to timeout.
        - Case 2: min_allowed_delta = 310, worst_case_time = 50, and 50 seconds to timeout.
        - Case 3: min_allowed_delta = 30, worst_case_time = 50, and 230 seconds to timeout.
        - Case 4: min_allowed_delta = 30, no worst_case_time set yet, and 230 seconds to timeout.
    When: Running is_interval_doesnt_have_enough_time_to_run
    Then: Ensure that the right results and worst_case_time are returned.
        - Case 1: should return True (meaning we should break) and worst_case_time = 250.
        - Case 2: should return True (meaning we should break) and worst_case_time = 50.
        - Case 3: should return False (meaning we shouldn't break yet) and worst_case_time = 50.
        - Case 4: Should return False (meaning we shouldn't break yet) and worst_case_time = 70.
    """
    import demistomock as demisto

    mocker.patch.object(demisto, "callingContext", {"context": {"TimeoutDuration": 300000000000}})
    Akamai_SIEM.EXECUTION_START_TIME = datetime(2024, 4, 10, 10, 0, 0)
    with freeze_time(freeze_mock):
        should_break, worst_case_time = Akamai_SIEM.is_interval_doesnt_have_enough_time_to_run(min_allowed_delta, worst_case_time)
        assert expected_time == worst_case_time
        assert should_break == expected_should_break


@pytest.mark.parametrize(
    "num_events_from_previous_request, page_size, expected_results",
    [
        (300, 400, True),
        (380, 400, False),
        (400, 400, False),
    ],
)
def test_is_last_request_smaller_than_page_size(num_events_from_previous_request, page_size, expected_results):
    """
    Given: num_events_from_previous_request, and page_size
        - Case 1: num_events_from_previous_request = 300, page_size = 400.
        - Case 2: num_events_from_previous_request = 380, page_size = 400.
        - Case 3: num_events_from_previous_request = 400, page_size = 400.
    When: Running is_last_request_smaller_than_page_size with ALLOWED_PAGE_SIZE_DELTA_RATIO = 0.95
    Then: Ensure that the right results and worst_case_time are returned.
        - Case 1: should return True (meaning we should break).
        - Case 2: should return False (meaning we shouldn't break yet).
        - Case 3: Should return False (meaning we shouldn't break yet).
    """
    assert Akamai_SIEM.is_last_request_smaller_than_page_size(num_events_from_previous_request, page_size) is expected_results


@pytest.mark.asyncio
@pytest.mark.parametrize(
    "offset, request_query",
    [(None, "50170?limit=100&from=1691303422"), ("test_offset", "50170?limit=100&offset=test_offset")],
)
async def test_get_events_concurrently_success(client, offset, request_query, requests_mock):
    """
    Given:
    - A successful mock http call response with 2 events and offset context.
    When:
    - Calling get_events_concurrently().
    Then:
    - Ensure the right log type and text returned. And the right length of events, and offset returned.
    """
    response_mock = '{"id": 1, "httpMessage": {"start": 1}}\n{"id": 2, "httpMessage": {"start": 2}}\n{"offset": "a"}'
    requests_mock.get(f"{BASE_URL}/{request_query}", text=response_mock)
    events, response_offset = await client.get_events_concurrently("50170", offset, 100, "1691303422", 2)
    assert len(events) == 2
    assert response_offset == "a"


@pytest.mark.asyncio
async def test_get_events_concurrently_failure(client, requests_mock):
    """
    Given:
    - An error mock http response.
    When:
    - Calling get_events_concurrently().
    Then:
    - Ensure the right log type and text returned.
    """
    requests_mock.get(f"{BASE_URL}/50170?limit=100&offset=offset", status_code=416, text="Requested Range Not Satisfiable")
    with pytest.raises(DemistoException) as e:
        await client.get_events_concurrently("50170", "offset", 100, "1691303422", 2)
    assert "Error in API call [416]" in str(e.value)
    assert "Requested Range Not Satisfiable" in str(e.value)


@pytest.mark.asyncio
async def test_get_events_from_akamai_success(mocker, client, requests_mock):
    """
    Given:
    - A successful mock http request response with 2 events and offset context.
    When:
    - Calling get_events_from_akamai().
    Then:
    - Ensure the right log type and text returned. And the right length of events, offset, and counter returned.
    """
    response_mock = '{"id": 1, "httpMessage": {"start": 1}}\n{"id": 2, "httpMessage": {"start": 2}}\n{"offset": "a"}'
    demisto_debug = mocker.patch.object(demisto, "debug")
    requests_mock.get(f"{BASE_URL}/50170?limit=100&offset=test_offset", text=response_mock)
    async for events, counter, response_offset in Akamai_SIEM.get_events_from_akamai(
        client=client, config_ids="50170", from_time="1 day", page_size=100, offset="test_offset", max_concurrent_tasks=300
    ):
        assert counter == 1
        assert len(events) == 2
        assert response_offset == "a"
        demisto_debug.assert_any_call("Running in interval = 1. got 2 events and offset='a'.")
        break


@pytest.mark.asyncio
async def test_get_events_from_akamai_non_416_failure(mocker, client, requests_mock):
    """
    Given:
    - A non-416 error mock http response.
    When:
    - Calling get_events_from_akamai().
    Then:
    - Ensure the error is reported to module health with is_error=True (non-recoverable).
    """
    requests_mock.get(f"{BASE_URL}/50170?limit=100&offset=test_offset", status_code=403, text="Unauthorized")
    mocker.patch.object(demisto, "error")
    # cause an exception to break endless loop.
    update_module_health = mocker.patch.object(demisto, "updateModuleHealth", side_effect=Exception("Interrupted execution"))
    with pytest.raises(Exception) as e:
        async for _, _, _ in Akamai_SIEM.get_events_from_akamai(
            client=client, config_ids="50170", from_time="1 day", page_size=100, offset="test_offset", max_concurrent_tasks=300
        ):
            pass
    assert str(e.value) == "Interrupted execution"  # Ensure the exception indeed was the planned one.
    assert {"is_error": True} in update_module_health.mock_calls[0]


@pytest.mark.asyncio
async def test_get_events_from_akamai_416_recovers(mocker, client, requests_mock):
    """
    Given:
    - A 416 offset-out-of-range error on the offset-based request, followed by a successful
      time-based (recovery) request.
    When:
    - Calling get_events_from_akamai() with a stale offset.
    Then:
    - The 416 is treated as a self-healing situation: the offset is dropped, the fetch window is
      shortened to AKAMAI_MAX_LOOKBACK_MINUTES, module health is NOT flagged with is_error=True,
      and the recovery request's events are yielded.
    """
    # First request (offset-based) fails with 416; the recovery request is time-based (uses 'from').
    requests_mock.get(f"{BASE_URL}/50170?limit=100&offset=stale_offset", status_code=416, text="Requested Range Not Satisfiable")
    recovery_response = '{"id": 1, "httpMessage": {"start": 1}}\n{"offset": "recovered"}'
    requests_mock.get(re.compile(rf"{re.escape(BASE_URL)}/50170\?.*from="), text=recovery_response)
    mocker.patch.object(demisto, "error")
    mocker.patch.object(asyncio, "sleep")  # don't actually sleep between iterations
    update_module_health = mocker.patch.object(demisto, "updateModuleHealth")

    async for events, _, response_offset in Akamai_SIEM.get_events_from_akamai(
        client=client, config_ids="50170", from_time="3 days", page_size=100, offset="stale_offset", max_concurrent_tasks=300
    ):
        assert len(events) == 1
        assert response_offset == "recovered"
        break

    # The 416 recovery must NOT flag the run as an error (we keep running).
    for call in update_module_health.mock_calls:
        assert {"is_error": True} not in call


@pytest.mark.asyncio
async def test_get_events_from_akamai_no_events(mocker, client, requests_mock):
    """
    Given:
    - A successful mock http response with no events.
    When:
    - Calling get_events_from_akamai().
    Then:
    - Ensure the right log type and text returned.
    """
    response_mock = '{"offset": "a"}'
    requests_mock.get(f"{BASE_URL}/50170?limit=100&offset=test_offset", text=response_mock)
    demisto_debug = mocker.patch.object(demisto, "debug")
    # cause an exception to break endless loop.
    mocker.patch.object(asyncio, "sleep", side_effect=Exception("Interrupted execution"))
    with pytest.raises(Exception) as e:
        async for _, _, _ in Akamai_SIEM.get_events_from_akamai(
            client=client, config_ids="50170", from_time="1 day", page_size=100, offset="test_offset", max_concurrent_tasks=300
        ):
            pass
    assert str(e.value) == "Interrupted execution"  # Ensure the exception indeed was the planned one.
    demisto_debug.assert_called_with(
        "Running in interval = 1. No events were received from Akamai,going to sleep for 60 seconds."
    )


@pytest.mark.asyncio
async def test_process_and_send_events_to_xsiam_skip_events_decoding(mocker):
    """
    Given:
    - 2 non serialized events
    When:
    - Calling process_and_send_events_to_xsiam() with should_skip_decode_events=True.
    Then:
    - Ensure the events were not modified, and the right logs were printed.
    """
    requestHeaders = "Content-Type%3A%20application/json%3Bcharset%3DUTF-8%0Auser%3A%20test%40test.com%0Aclient%3A%20"
    "test_client%0AX-Kong-Upstream-Latency%3A%2066%0AX-Kong-Proxy-Latency%3A%202%0AX-Kong-Request-Id%3A%20X_request_id%"
    "0AEPM-Request-ID%3A%20EPM_request_id%0AContent-Length%3A%20157%0ADate%3A%20Mon%2C%2025%20Mar%202024%2013%3A52%3A11"
    "%20GMT%0AConnection%3A%20keep-alive%0AServer-Timing%3A%20cdn-cache%3B%20desc%3DMISS%0AServer-Timing%3A%20edge%3B"
    "%20dur%3D23%0AServer-Timing%3A%20origin%3B%20dur%3D72%0AServer-Timing%3A%20intid%3Bdesc%3Ddd%0A"
    "Strict-Transport-Security%3A%20max-age%3D31536000%20%3B%20includeSubDomains%20%3B%20preload%0A"
    events = [
        f'{{"id": 1, "httpMessage": {{"start": 1591303422, "requestHeaders": "{requestHeaders}"}}}}',
        f'{{"id": 2, "httpMessage": {{"start": 1591303422, "requestHeaders": "{requestHeaders}"}}}}',
    ]
    demisto_debug = mocker.patch.object(demisto, "debug")
    send_events_to_xsiam_akamai = mocker.patch(
        "Akamai_SIEM.send_events_to_xsiam_akamai", side_effect=Exception("Interrupted execution")
    )  # to break endless loop.
    with pytest.raises(Exception) as e:
        await Akamai_SIEM.process_and_send_events_to_xsiam(events, should_skip_decode_events=True, offset="test", counter=1)
    assert str(e.value) == "Interrupted execution"  # Ensure the exception indeed was the planned one.
    assert send_events_to_xsiam_akamai.call_args_list[0][0][0] == events
    assert isinstance(send_events_to_xsiam_akamai.call_args_list[0][0][0][0], str)
    demisto_debug.assert_has_calls(
        [
            mocker.call(f"Running in interval = 1. got {len(events)} events, moving to processing events data."),
            mocker.call("Running in interval = 1. Skipping decode events."),
            mocker.call(
                f"[Fetch] Running in interval = 1. Sending {len(events)} events to xsiam. "
                "latest event time is: 2020-06-04T20:43:42Z"
            ),
        ]
    )


@pytest.mark.asyncio
async def test_process_and_send_events_to_xsiam_with_events_decoding(mocker):
    """
    Given:
    - 2 non serialized events
    When:
    - Calling process_and_send_events_to_xsiam() with should_skip_decode_events=False.
    Then:
    - Ensure the events were loaded decoded correctly, and the right logs were printed.
    """
    requestHeaders = "Content-Type%3A%20application/json%3Bcharset%3DUTF-8%0Auser%3A%20test%40test.com%0Aclient%3A%20"
    "test_client%0AX-Kong-Upstream-Latency%3A%2066%0AX-Kong-Proxy-Latency%3A%202%0AX-Kong-Request-Id%3A%20X_request_id%"
    "0AEPM-Request-ID%3A%20EPM_request_id%0AContent-Length%3A%20157%0ADate%3A%20Mon%2C%2025%20Mar%202024%2013%3A52%3A11"
    "%20GMT%0AConnection%3A%20keep-alive%0AServer-Timing%3A%20cdn-cache%3B%20desc%3DMISS%0AServer-Timing%3A%20edge%3B"
    "%20dur%3D23%0AServer-Timing%3A%20origin%3B%20dur%3D72%0AServer-Timing%3A%20intid%3Bdesc%3Ddd%0A"
    "Strict-Transport-Security%3A%20max-age%3D31536000%20%3B%20includeSubDomains%20%3B%20preload%0A"
    events = [
        f'{{"id": 1, "httpMessage": {{"start": 1491303422, "requestHeaders": "{requestHeaders}"}}}}',
        f'{{"id": 2, "httpMessage": {{"start": 1591303422, "requestHeaders": "{requestHeaders}"}}}}',
    ]
    demisto_debug = mocker.patch.object(demisto, "debug")
    send_events_to_xsiam_akamai = mocker.patch(
        "Akamai_SIEM.send_events_to_xsiam_akamai", side_effect=Exception("Interrupted execution")
    )  # to break endless loop.
    with pytest.raises(Exception) as e:
        await Akamai_SIEM.process_and_send_events_to_xsiam(events, should_skip_decode_events=False, offset="test", counter=1)
    assert str(e.value) == "Interrupted execution"  # Ensure the exception indeed was the planned one.
    processed_events = [
        {
            "id": 1,
            "httpMessage": {
                "start": 1491303422,
                "requestHeaders": {"Content_Type": "application/json;charset=UTF-8", "user": "test@test.com", "client": ""},
                "responseHeaders": {},
            },
        },
        {
            "id": 2,
            "httpMessage": {
                "start": 1591303422,
                "requestHeaders": {"Content_Type": "application/json;charset=UTF-8", "user": "test@test.com", "client": ""},
                "responseHeaders": {},
            },
        },
    ]
    assert send_events_to_xsiam_akamai.call_args_list[0][0][0] == processed_events
    assert isinstance(send_events_to_xsiam_akamai.call_args_list[0][0][0][0], dict)
    demisto_debug.assert_has_calls(
        [
            mocker.call(f"Running in interval = 1. got {len(events)} events, moving to processing events data."),
            mocker.call("Running in interval = 1. decoding events."),
            mocker.call(
                f"[Fetch] Running in interval = 1. Sending {len(events)} events to xsiam. "
                "latest event time is: 2020-06-04T20:43:42Z"
            ),
        ]
    )


def test_test_fetch_events_long_running_command_flow(mocker, client, caplog):
    async def test_fetch_events_long_running_command_flow(mocker, client):
        """
        Given:
        - 2 mock response for 2 consecutive requests, one with 2 events, and one with 1 event.
        When:
        - Calling fetch_events_long_running_command() with page_size=2.
        Then:
        - Ensure the events were fetched and moved to send_events_to_xsiam as expected.
        The send_events_to_xsiam_akamai function was called once and the 2 requests were sent.
        The execution went into sleep after the second request as there were less the limit events to fetch.
        The right log was printed.
        """

        from Akamai_SIEM import Client

        requestHeaders = "Content-Type%3A%20application/json%3Bcharset%3DUTF-8%0Auser%3A%20test%40test.com%0Aclient%3A%20"
        "test_client%0AX-Kong-Upstream-Latency%3A%2066%0AX-Kong-Proxy-Latency%3A%202%0AX-Kong-Request-Id%3A%20X_request_id%"
        "0AEPM-Request-ID%3A%20EPM_request_id%0AContent-Length%3A%20157%0ADate%3A%20Mon%2C%2025%20Mar%202024%2013%3A52%3A11"
        "%20GMT%0AConnection%3A%20keep-alive%0AServer-Timing%3A%20cdn-cache%3B%20desc%3DMISS%0AServer-Timing%3A%20edge%3B"
        "%20dur%3D23%0AServer-Timing%3A%20origin%3B%20dur%3D72%0AServer-Timing%3A%20intid%3Bdesc%3Ddd%0A"
        "Strict-Transport-Security%3A%20max-age%3D31536000%20%3B%20includeSubDomains%20%3B%20preload%0A"
        event_1 = f'{{"id": 1, "httpMessage": {{"start": 1, "requestHeaders": "{requestHeaders}"}}}}'
        event_2 = f'{{"id": 2, "httpMessage": {{"start": 2, "requestHeaders": "{requestHeaders}"}}}}'
        event_3 = f'{{"id": 3, "httpMessage": {{"start": 3, "requestHeaders": "{requestHeaders}"}}}}'
        response_mock_1 = f'{event_1}\n{event_2}\n{{"offset": "a"}}'
        response_mock_2 = f'{event_3}\n{{"offset": "b"}}'
        execution_count = 0

        def count_execution(*args, **kwargs):
            nonlocal execution_count
            execution_count += 1
            if execution_count == 1:
                return response_mock_1
            else:
                return response_mock_2

        mocker.patch.object(Client, "_http_request", side_effect=count_execution)

        async def mock_func():
            return 2

        demisto_debug = mocker.patch.object(demisto, "debug")
        send_events_to_xsiam_akamai = mocker.patch(
            "Akamai_SIEM.send_events_to_xsiam_akamai", side_effect=asyncio.create_task(mock_func())
        )  # to break endless loop.
        mocker.patch.object(asyncio, "sleep", side_effect=Exception("Interrupted execution"))  # to break endless loop.
        with pytest.raises(Exception) as e:
            await Akamai_SIEM.fetch_events_long_running_command(client, "5 minutes", 2, "50170", {}, True, 5)
        assert str(e.value) == "Interrupted execution"  # Ensure the exception indeed was the planned one.
        assert send_events_to_xsiam_akamai.call_count == 1
        assert execution_count == 2
        assert send_events_to_xsiam_akamai.call_args_list[0][0][0] == [event_1, event_2]
        demisto_debug.assert_called_with(
            "Running in interval = 2. got 1 events which is less than 0.95 % of the page_size=2, going to sleep for 60 seconds."
        )

    asyncio.run(test_fetch_events_long_running_command_flow(mocker, client))
    caplog.clear()


def test_fetch_events_command_double_416_raises_no_spin(client, mocker):
    """
    Given:
    - A stored offset whose request fails with a 416 (expired offset), and whose recovery request
      ALSO fails with a 416.
    When:
    - Calling fetch_events_command with a bounded fetch_limit.
    Then:
    - Ensure the loop does NOT spin forever: after a single recovery, a second consecutive 416 is
      surfaced as a DemistoException (guard against an infinite retry loop). The offset is reset
      exactly once.
    """
    err_msg = "Error in API call [416] - Requested Range Not Satisfiable"
    mocker.patch.object(
        Akamai_SIEM.Client,
        "get_events_with_offset",
        side_effect=[
            DemistoException(err_msg, res={}),  # first: stale offset -> 416 (recovered)
            DemistoException(err_msg, res={}),  # recovery ALSO -> 416 -> must raise, not spin
        ],
    )
    mocker.patch.object(Akamai_SIEM, "is_interval_doesnt_have_enough_time_to_run", return_value=(False, 1))
    reset_offset_mock = mocker.patch.object(Akamai_SIEM, "reset_offset_command")
    mocker.patch.object(demisto, "error")
    mocker.patch.object(demisto, "debug")
    mocker.patch.object(demisto, "info")

    with pytest.raises(DemistoException) as e:
        for _events, _, _, _ in Akamai_SIEM.fetch_events_command(  # noqa: B007
            client,
            "3 days",
            220,
            "",
            {"offset": "stale_offset"},
            5000,
            True,
        ):
            pass

    assert "Requested Range Not Satisfiable" in str(e.value)
    # Only one recovery was attempted before aborting - no infinite reset spin.
    reset_offset_mock.assert_called_once()


""" main() tests """


def _run_main_fetch_events(mocker, params, ctx, fetch_events_pages, send_side_effect=None):
    """Helper to drive main() for the 'fetch-events' command with everything mocked.

    Returns a dict with the mocks so individual tests can assert on them.
    """
    mocker.patch.object(demisto, "command", return_value="fetch-events")
    mocker.patch.object(demisto, "params", return_value=params)
    mocker.patch.object(demisto, "debug")
    mocker.patch.object(demisto, "info")
    mocker.patch.object(demisto, "error")
    mocker.patch.object(Akamai_SIEM, "EdgeGridAuth", return_value=None)
    mocker.patch.object(Akamai_SIEM, "get_integration_context", return_value=ctx)
    set_context_mock = mocker.patch.object(Akamai_SIEM, "set_integration_context")
    set_last_run_mock = mocker.patch.object(demisto, "setLastRun")
    update_health_mock = mocker.patch.object(demisto, "updateModuleHealth")
    send_events_mock = mocker.patch.object(Akamai_SIEM, "send_events_to_xsiam", side_effect=send_side_effect)
    return_error_mock = mocker.patch.object(Akamai_SIEM, "return_error")

    captured_kwargs = {}

    def fake_fetch_events_command(*args, **kwargs):
        captured_kwargs.update(kwargs)
        yield from fetch_events_pages

    mocker.patch.object(Akamai_SIEM, "fetch_events_command", side_effect=fake_fetch_events_command)

    Akamai_SIEM.main()

    return {
        "set_context": set_context_mock,
        "set_last_run": set_last_run_mock,
        "update_health": update_health_mock,
        "send_events": send_events_mock,
        "return_error": return_error_mock,
        "captured_kwargs": captured_kwargs,
    }


BASE_MAIN_PARAMS = {
    "configIds": "50170",
    "host": "https://example.com",
    "isFetch": False,
}


@pytest.mark.parametrize(
    "events_fetch_limit_param, expected_page_size, expected_limit",
    [
        (80000, Akamai_SIEM.DEFAULT_PAGE_SIZE, 80000),
        (Akamai_SIEM.MAX_ALLOWED_FETCH_LIMIT + 1, Akamai_SIEM.DEFAULT_PAGE_SIZE, Akamai_SIEM.MAX_ALLOWED_FETCH_LIMIT),
        (3000, 3000, 3000),
        (40000, Akamai_SIEM.DEFAULT_PAGE_SIZE, 40000),
    ],
)
def test_main_fetch_events_param_clamps(mocker, events_fetch_limit_param, expected_page_size, expected_limit):
    """
    Given:
    - fetch-events params with various eventsFetchLimit values (out-of-bounds and in-bounds).
    When:
    - Running main() for the fetch-events command.
    Then:
    - Ensure page_size defaults to DEFAULT_PAGE_SIZE (lowered to the limit when limit is smaller) and
      fetch_limit is clamped correctly before being passed to fetch_events_command.
    """
    params = {**BASE_MAIN_PARAMS, "eventsFetchLimit": events_fetch_limit_param}
    result = _run_main_fetch_events(mocker, params, ctx={"offset": "ctx_offset"}, fetch_events_pages=[([], None, 0, False)])
    captured = result["captured_kwargs"]
    assert captured["page_size"] == expected_page_size
    assert captured["fetch_limit"] == expected_limit


def test_main_fetch_events_uses_events_default_when_no_limit_set(mocker):
    """
    Given:
    - fetch-events params without eventsFetchLimit or fetchLimit.
    When:
    - Running main() for the fetch-events command.
    Then:
    - Ensure fetch_limit defaults to DEFAULT_EVENTS_FETCH_LIMIT (60000).
    """
    result = _run_main_fetch_events(
        mocker, {**BASE_MAIN_PARAMS}, ctx={"offset": "ctx_offset"}, fetch_events_pages=[([], None, 0, False)]
    )
    assert result["captured_kwargs"]["fetch_limit"] == Akamai_SIEM.DEFAULT_EVENTS_FETCH_LIMIT


def test_main_fetch_events_ignores_incident_fetch_limit(mocker):
    """
    Given:
    - fetch-events params with only the incident fetchLimit set (no eventsFetchLimit).
    When:
    - Running main() for the fetch-events command.
    Then:
    - Ensure fetchLimit (which controls incident fetching only) is ignored for event collection,
      and the events limit falls back to DEFAULT_EVENTS_FETCH_LIMIT.
    """
    params = {**BASE_MAIN_PARAMS, "fetchLimit": 45000}
    result = _run_main_fetch_events(mocker, params, ctx={"offset": "ctx_offset"}, fetch_events_pages=[([], None, 0, False)])
    assert result["captured_kwargs"]["fetch_limit"] == Akamai_SIEM.DEFAULT_EVENTS_FETCH_LIMIT


def test_main_fetch_events_uses_events_limit_regardless_of_incident_limit(mocker):
    """
    Given:
    - fetch-events params with both eventsFetchLimit and the incident fetchLimit set.
    When:
    - Running main() for the fetch-events command.
    Then:
    - Ensure eventsFetchLimit controls event collection and the incident fetchLimit has no effect.
    """
    params = {**BASE_MAIN_PARAMS, "eventsFetchLimit": 50000, "fetchLimit": 20}
    result = _run_main_fetch_events(mocker, params, ctx={"offset": "ctx_offset"}, fetch_events_pages=[([], None, 0, False)])
    assert result["captured_kwargs"]["fetch_limit"] == 50000


def test_main_fetch_events_offset_persisted_on_success(mocker):
    """
    Given:
    - A successful fetch-events run that yields a page of events and a new offset.
    When:
    - Running main() and send_events_to_xsiam succeeds.
    Then:
    - Ensure the new offset is persisted to integration context and eventsPulled is reported.
    """
    pages = [
        (['{"id": 1}', '{"id": 2}'], "new_offset", 2, False),
        ([], "new_offset", 2, False),
    ]
    result = _run_main_fetch_events(mocker, {**BASE_MAIN_PARAMS}, ctx={"offset": "old"}, fetch_events_pages=pages)
    result["set_context"].assert_called_once_with({"offset": "new_offset"})
    result["update_health"].assert_called_once_with({"eventsPulled": 2})
    result["return_error"].assert_not_called()


def test_main_fetch_events_offset_not_persisted_on_send_failure(mocker):
    """
    Given:
    - A fetch-events run that yields events, but send_events_to_xsiam raises.
    When:
    - Running main().
    Then:
    - Ensure should_fail is set, the offset is NOT persisted (data-integrity guard),
      and the error is surfaced via return_error.
    """
    pages = [(['{"id": 1}'], "new_offset", 1, False)]
    result = _run_main_fetch_events(
        mocker,
        {**BASE_MAIN_PARAMS},
        ctx={"offset": "old"},
        fetch_events_pages=pages,
        send_side_effect=Exception("send failed"),
    )
    result["set_context"].assert_not_called()
    result["return_error"].assert_called_once()


@pytest.mark.parametrize(
    "pages, expected_next_trigger",
    [
        ([([f'{{"id": {i}}}' for i in range(20)], "off", 20, False), ([], "off", 20, False)], "0"),
        ([(['{"id": 1}'], "off", 1, True), ([], "off", 1, True)], "0"),
        ([(['{"id": 1}'], "off", 1, False), ([], "off", 1, False)], None),
    ],
)
def test_main_fetch_events_next_trigger(mocker, pages, expected_next_trigger):
    """
    Given:
    - fetch-events runs with different (total_events_count, auto_trigger_next_run) outcomes.
    When:
    - Running main() with eventsFetchLimit=20.
    Then:
    - Ensure nextTrigger is set to "0" when the interval hit the limit or requested an auto-trigger,
      and is absent otherwise.
    """
    params = {**BASE_MAIN_PARAMS, "eventsFetchLimit": 20}
    result = _run_main_fetch_events(mocker, params, ctx={"offset": "old"}, fetch_events_pages=pages)
    next_run = result["set_last_run"].call_args[0][0]
    assert next_run.get("nextTrigger") == expected_next_trigger


def test_main_fetch_events_streaming_clears_events_on_success(mocker):
    """
    Given:
    - A fetch-events run that yields a page of events and send_events_to_xsiam succeeds.
    When:
    - Running main().
    Then:
    - Ensure the yielded page list is cleared after a successful streaming send (memory hygiene).
    """
    page_events = ['{"id": 1}', '{"id": 2}']
    pages = [(page_events, "off", 2, False), ([], "off", 2, False)]
    result = _run_main_fetch_events(mocker, {**BASE_MAIN_PARAMS}, ctx={"offset": "old"}, fetch_events_pages=pages)
    result["send_events"].assert_called_once()
    assert page_events == []


def test_main_fetch_events_streaming_raises_and_keeps_events_on_failure(mocker):
    """
    Given:
    - A fetch-events run that yields a page of events, but send_events_to_xsiam raises.
    When:
    - Running main().
    Then:
    - Ensure the streaming failure is surfaced via return_error, the offset is NOT persisted
      (data-integrity guard), and the page list is NOT cleared (so the events can be retried).
    """
    page_events = ['{"id": 1}', '{"id": 2}']
    pages = [(page_events, "off", 2, False)]
    result = _run_main_fetch_events(
        mocker,
        {**BASE_MAIN_PARAMS},
        ctx={"offset": "old"},
        fetch_events_pages=pages,
        send_side_effect=Exception("send failed"),
    )
    result["send_events"].assert_called_once()
    result["return_error"].assert_called_once()
    result["set_context"].assert_not_called()
    # events.clear() is only reached after a successful send, so on failure the page is preserved.
    assert page_events == ['{"id": 1}', '{"id": 2}']