SplunkPy v2

Run queries on Splunk and fetch Splunk ES Findings and Investigations (Splunk ES 8.2+).

Analytics & SIEM · Splunk

Details

IDSplunkPy v2
ProviderCisco Systems
CategoryAnalytics & SIEM
From Version6.0.0
Docker Imagedemisto/splunksdk-py3:1.0.0.10133006
Supported ModulesAgentix XSIAM Cloud Posture Security EDR Cortex Cloud Cloud Runtime Security

README

Use the SplunkPy v2 integration to:

  • Fetch events (logs) from Splunk into Cortex
  • Push events from Cortex to Splunk
  • Fetch Splunk Enterprise Security events (Findings) into Cortex.

Important: This integration is designed specifically for Splunk Enterprise Security version 8.2 and above.
For Splunk ES versions prior to 8.2, please use the legacy SplunkPy integration.

This integration was integrated and tested with Splunk Enterprise v10.0.0 and Enterprise Security v8.2.3.

Use Cases


User Configuration Requirements

Option one

Assign the following roles to the user: admin, ess_admin (for working with Splunk Enterprise Security).

Option two

When assigning admin is not an option.
Define a custom role and include all necessary capabilities: (permissions)
image
Define the indexes configuration for the custom role:
image
At the end of the process, the Splunk user (not admin user) should receive the previously created role. Following is the list of capabilities that covers both the Splunk UI access and using the SplunkPy v2 integration including Enterprise security:

  • accelerate_search
  • admin_all_objects
  • can_own_notable_events
  • change_own_password
  • edit_analyticstories
  • edit_cam_queue
  • edit_correlationsearches
  • edit_lookups
  • edit_notable_events
  • edit_own_objects
  • edit_tcp
  • edit_tcp_stream
  • edit_upload_and_index
  • get_metadata
  • get_typeahead
  • input_file
  • list_accelerate_search
  • list_all_objects
  • list_inputs
  • list_introspection
  • list_metrics_catalog
  • list_search_head_clustering
  • manage_all_investigations
  • manage_behavioral_analytics
  • output_file
  • rest_access_server_endpoints
  • rest_apps_view
  • rest_properties_get
  • rest_properties_set
  • rtsearch
  • run_collect
  • run_mcollect
  • run_msearch
  • run_sendalert
  • schedule_rtsearch
  • schedule_search
  • search
  • search_process_config_refresh
  • upload_lookup_files

SplunkPy v2 command permissions by example

splunk-finding-event-edit

custom roles required: (at least one)
ess_analyst, ess_admin
Can replace the Splunk power role for ES users.

User-Level Permissions:
Read: Access to view resources (e.g., dashboards, reports).
Write: Ability to modify existing resources or create new ones.

Query Load Analysis

Mirroring

When mirroring in enabled, 2-3 simultaneous queries are expected.
Each query can have more that one API call. On mirror out - API call for updating each finding event that changed. This also includes User mapping queries and mirroring queries.

Enrichment

Fetching finding event - for each fetch iteration 2 queries.
For each finding event that fetched - we have 3 enrichments(max).
<Amount of fetched findings> * <amount of defined enrichments>
In case of more that one drilldown - number of drilldown queries.
Each query can have more that one API call.

Fetch

Configured by the instance configuration max_fetch (behind the scenes an query can made few API calls).

Configure SplunkPy v2 in Cortex

Parameter Description Required
Server URL The Splunk server URL. Port 8089 (Splunk’s default REST API port) is used automatically. Only include the port in the URL if using a non-default port. Examples: ‘https://splunk.example.com’ (uses default port 8089) or ‘https://splunk.example.com:8090’ (uses custom port 8090). True
Splunk Token   True
Fetch events query The Splunk search query by which to fetch events. The default query fetches ES finding events. You can edit this query to fetch other types of events. Note, that to fetch ES finding events, make sure to include the \`notable\` macro in your query. False
Fetch Limit (Max.- 200, Recommended less than 50)   False
Fetch incidents   False
Incident type   False
Parse Raw Part of Finding Events Whether to parse the raw part of the Findings, or not. False
Replace with Underscore in Incident Fields Whether to replace special characters to underscore when parsing the raw data of the Findings, or not. False
First fetch timestamp (<number> <time unit>, e.g., 12 hours, 7 days, 3 months, 1 year) The amount of time to go back when performing the first fetch, or when creating a mapping using the Select Schema option. False
Event types to fetch Select the Splunk event types to ingest. Default is `Finding`. False
First fetch timestamp (Investigations) The relative time interval to look back during the initial investigation fetch (for example, 12 hours, 7 days, 3 months). False
Investigations fetch query The SPL query used when “Investigation” is selected for “Event types to fetch”. The query must include the `FETCH_FILTER_PLACEHOLDER` token. Do not modify or remove this token. For more information on customizing the query (for example, adding the filter &status=New), see the integration documentation under “Fetching investigation events”. False
Maximum investigations per fetch The maximum number of investigations to fetch per cycle. Limited to 100 by the Splunk investigations endpoint. False
Extract Fields - CSV fields that will be parsed out of raw finding events   False
Incident Mirroring Direction Choose the direction to mirror the incident: Incoming (from Splunk to Cortex XSOAR), Outgoing (from Cortex XSOAR to Splunk), or Incoming and Outgoing (from/to Cortex XSOAR and Splunk). False
Close Mirrored Cortex XSOAR Incidents (Incoming Mirroring) When selected, closing the Splunk finding event with a “Closed” status will close the Cortex XSOAR incident. False
Additional Splunk status labels to close on mirror (Incoming Mirroring) A comma-separated list of Splunk status labels to mirror as closed Cortex XSOAR incident (Example: Resolved,False-Positive). False
Enable Splunk statuses marked as “End Status” to close on mirror (Incoming Mirroring) When selected, automatically close the Cortex XSOAR incident when the Splunk ES event (Finding or Investigation) is marked as ‘End Status’. False
Close Mirrored Splunk ES Events (Outgoing Mirroring) When selected, automatically close the corresponding Splunk ES event (Finding or Investigation) when the Cortex XSOAR incident is closed. False
Trust any certificate (not secure)   False
Use system proxy settings   False
The app context of the namespace   False
HEC Token (HTTP Event Collector)   False
HEC BASE URL (e.g: https://localhost:8088 or https://example.splunkcloud.com/).   False
Enrichment Types Enrichment types to enrich each fetched finding. If none are selected, the integration will fetch findings as usual (without enrichment).  
For more info about enrichment types see Enriching Finding Events. False  
Asset enrichment lookup tables CSV of the Splunk lookup tables from which to take the Asset enrichment data. False
Identity enrichment lookup tables CSV of the Splunk lookup tables from which to take the Identity enrichment data. False
Enrichment Timeout (Minutes) When the selected timeout was reached, Finding events that were not enriched will be saved without the enrichment. False
Number of Events Per Enrichment Type The limit of how many events to retrieve per each one of the enrichment types (Drilldown, Asset, and Identity). In a case of multiple drilldown enrichments the limit will apply for each drilldown search query. To retrieve all events, enter “0” (not recommended). False
Advanced: Extensive logging (for debugging purposes). Do not use this option unless advised otherwise.   False
Advanced: Time type to use when fetching events Defines which timestamp will be used to filter the events:
- creation time: Filters based on when the event actually occurred.
- index time (Beta): *Beta feature* – Filters based on when the event was ingested into Splunk.
This option is still in testing and may not behave as expected in all scenarios.
When using this mode, the parameter “Fetch backwards window for the events occurrence time (minutes)” should be set to `0``, as indexing time ensures there are no delay-based gaps.
The default is “creation time”.
 
Advanced: Fetch backwards window for the events occurrence time (minutes) The fetch time range will be at least the size specified here. This will support events that have a gap between their occurrence time and their index time in Splunk. To decide how long the backwards window should be, you need to determine the average time between them both in your Splunk environment. False
Advanced: Unique ID Fields A comma-separated list of additional fields to use when generating unique incident IDs for events that are not findings (i.e., queries without the `notable` macro). By default, the integration uses: _cd, index,_time, _indextime,_raw. If these fields do not provide unique values in your environment, specify additional fields here to ensure incident uniqueness. Example: source,host,unique_field False
Enable user mapping Whether to enable the user mapping between Cortex XSOAR and Splunk, or not. For more information see https://xsoar.pan.dev/docs/reference/integrations/splunk-py#configure-user-mapping-between-splunk-and-cortex-xsoar False
Users Lookup table name The name of the lookup table in Splunk, containing the username’s mapping data. False
XSOAR user key The name of the lookup column containing the Cortex XSOAR username. False
SPLUNK user key The name of the lookup table containing the Splunk username. False
Note tag from Splunk Add this tag to an entry to mirror it as a note from Splunk. False
Note tag to Splunk Add this tag to an entry to mirror it as a note to Splunk. False
Incidents Fetch Interval   False

Note: To use a Splunk Cloud instance, contact Splunk support to request API access. Use a non-SAML account to access the API.

Splunk Enterprise Security Users

Note: The following information is for Splunk Enterprise Security version 8.2+ users.
This integration requires Splunk ES 8.2 or higher due to the new Finding Events API and terminology changes (Notable Events → Finding Events, Comments → Notes).
For Splunk ES versions prior to 8.2, use the legacy SplunkPy integration.
For Splunk non-Enterprise Security Users, see Splunk non-Enterprise Security Users.

Fetching finding events

The integration allows for fetching Splunk Finding events using a default query. The query can be changed and modified to support different Splunk use cases.

Enriching finding events

This integration allows 3 types of enrichments for fetched findings: Drilldown, Asset, and Identity.

Enrichment types

  1. Drilldown search enrichment: Fetches the drilldown searches configured by the user in the rule name that triggered the finding event and performs this search. The results are stored in the context of the incident under the Drilldown field as follows: [{‘query_name’:, 'query_search': , 'query_results': [{result1}, {result2}, {result3}], 'enrichment_status': }].
  2. Asset search enrichment: Runs the following query:
    | inputlookup append=T asset_lookup_by_str where asset=$ASSETS_VALUE | inputlookup append=t asset_lookup_by_cidr where asset=$ASSETS_VALUE | rename _key as asset_id | stats values(*) as * by asset_id
    where the $ASSETS_VALUE is replaced with the src, dest, src_ip and dst_ip from the fetched finding. The results are stored in the context of the incident under the Asset field.
  3. Identity search enrichment: Runs the following query
    | inputlookup identity_lookup_expanded where identity=$IDENTITY_VALUE
    where the $IDENTITY_VALUE is replaced with the user and src_user from the fetched finding event. The results are stored in the context of the incident under the Identity field.

How to configure

  1. Configure the integration to fetch incidents.
  2. Enrichment Types: Select the enrichment types you want to enrich each fetched finding with. If none are selected, the integration will fetch findings as usual (without enrichment).
  3. Fetch events query: The query for fetching events. The default query is for fetching finding events. You can edit this query to fetch other types of events. Note that to fetch finding events, make sure the query uses the `notable` macro.
  4. Enrichment Timeout (Minutes): The timeout for each enrichment (default is 5min). When the selected timeout was reached, finding events that were not enriched will be saved without the enrichment.
  5. Number of Events Per Enrichment Type: The maximal amount of events to fetch per enrichment type (Drilldown, Asset, and Identity). In a case of multiple drilldown enrichments the limit will apply for each drilldown search query. (default to 20).

Troubleshooting enrichment status

Each enriched incident contains the following fields in the incident context:

  • successful_drilldown_enrichment: whether the drilldown enrichment was successful. In a case of multiple drilldown enrichments, the status is successful if at least one drilldown search enrichment was successful.
  • successful_asset_enrichment: whether the asset enrichment was successful.
  • successful_identity_enrichment: whether the identity enrichment was successful.

Resetting the enriching fetch mechanism

  • Run the Last Run button
  • Run the splunk-reset-enriching-fetch-mechanism command and the mechanism will be reset to the initial configuration.

Enrichment Limitations

  • As the enrichment process is asynchronous, fetching enriched incidents takes longer. The integration was tested with 20+ findings simultaneously that were fetched and enriched after approximately ~4min.
  • If you wish to configure a mapper, wait for the integration to perform the first fetch successfully. This is to make the fetch mechanism logic stable.
  • The drilldown search, does not support Splunk’s advanced syntax. For example: Splunk filters (** s**, ** h**, etc.)

Configure User Mapping between Splunk and Cortex XSOAR/XSIAM

When fetching incidents from Splunk to Cortex XSOAR and when mirroring incidents between Splunk and Cortex XSOAR, the Splunk Owner Name (user) associated with an incident needs to be mapped to the relevant Cortex XSOAR Owner Name (user).
You can use Splunk to define a user lookup table and then configure the SplunkPy v2 integration instance to enable the user mapping. Alternatively, you can map the users with a script or a transformer.

Note:

  • When mapping users, the specified Cortex XSOAR/XSIAM user must be a valid user in the system.
  • The Cortex XSOAR Owner incident field can only be used for mirroring changes out to Splunk, you cannot use it to update Cortex XSOAR incidents based on values from Splunk. To mirror changes in from Splunk, use the Assigned User incident field.

Configure User Mapping Using Splunk

  1. Define the lookup table in Splunk.
    1. Under App: Lookup Editor, select Lookup Editor.

      image

    2. Select Create a New Lookup > KV Store lookup.

      image

    3. Enter the Name for the table. For example, splunk_xsoar_users is the default lookup table name defined in the SplunkPy v2 integration settings.
    4. Under App, select Enterprise Security.
    5. Assign two Key-value collection schema fields, one for the Cortex XSOAR/XSIAM usernames and one for the corresponding Splunk usernames. For example, xsoar_user and splunk_user are the default field values defined in the SplunkPy v2 integration settings.
    6. Click Create Lookup.
      image

    7. Add values to the table to map Cortex XSOAR/XSIAM users to the Splunk users.
      image

    Note:
    If the user keys are defined already in another table, you can use that table name and relevant key names in the SplunkPy integration settings.

  2. Configure the SplunkPy v2 integration instance.
    1. Under Settings > Integrations, search for the SplunkPy v2 integration and create an instance.
    2. In the Integration Settings:
      1. Select Enable user mapping.
      2. Set Users Lookup table name to the name of the lookup table defined in Splunk. By default it is splunk_xsoar_users.
      3. Set the XSOAR user key to the field defined in the Splunk lookup table. By default it is xsoar_user.
      4. Set the SPLUNK user key to the field defined in the Splunk lookup table. By default it is splunk_user.

        image

Incident Mirroring

Important Notes

  • Mirroring-in is not supported when multiple Splunk integration instances are connected to the same Splunk server, meaning only one instance per Splunk server can be configured to perform mirroring-in.
  • This feature is available from Cortex XSOAR version 6.0.0.
  • This feature is supported by Splunk Enterprise Security only.
  • In order for the mirroring to work, the Incident Mirroring Direction parameter needs to be set before the incident is fetched.
  • In order to ensure the mirroring works as expected, mappers are required, both for incoming and outgoing, to map the expected fields in Cortex XSOAR and Splunk.
  • For mirroring the owner field, the usernames need to be transformed to the corresponding in Cortex XSOAR and Splunk.

Splunk Notes Mirroring

  • Splunk Notes Updates/Deletions - Editing or deleting existing Finding notes will NOT trigger mirroring to XSOAR on their own. Notes changes will only appear in the Splunk Notes field when a “real” change occurs in another finding field (such as status, owner, urgency, etc.), which triggers the mirror-in process. However, these changes will NOT update War Room notes. War Room notes from Splunk will NOT be deleted or updated.
  • Notes Display Behavior - The Splunk Notes field will display notes from the past week only. However, all notes that were mirrored via mirror-in will appear in the War Room notes.
  • Notes Time - Note timestamps will display the same time as shown in Splunk for the user who created the authentication token. The timezone offset is based on the timezone configured for that user in Splunk.

You can enable incident mirroring between Cortex XSOAR incidents and Splunk findings.

To set up mirroring:

  1. Navigate to Settings > Integrations > Servers & Services.
  2. Search for SplunkPy v2 and select your integration instance.
  3. Enable Fetches incidents.
  4. You can go to the Fetch events query parameter and select the query to fetch the findings from Splunk. Make sure to provide a query which uses the `notable` macro, See the default query as an example.
  5. In the Incident Mirroring Direction integration parameter, select in which direction the incidents should be mirrored:
    • Incoming - Any changes in Splunk findings (finding’s status, status_label, urgency, notes, and owner) will be reflected in Cortex XSOAR incidents.
    • Outgoing - Any changes in Cortex XSOAR incidents (finding’s status, urgency, notes, and owner) will be reflected in Splunk findings.
    • Incoming And Outgoing - Changes in Cortex XSOAR incidents and Splunk findings will be reflected in both directions.
    • None - Turns off incident mirroring.
  6. Optional: Check the Close Mirrored Cortex XSOAR Incidents (Incoming Mirroring) integration parameter to close the Cortex XSOAR incident when the corresponding finding is closed on the Splunk side.
    By default, only Findings closed with a “Closed” label will be mirrored. You can specify specific statuses (comma-separated) in the Additional Splunk status labels to close on mirror (Incoming Mirroring), and enable the Enable Splunk statuses marked as “End Status” to close on mirror (Incoming Mirroring) option to add statuses marked as “End Status” in Splunk, and to add additional statuses to the mirroring process.
  7. Optional: Check the Close Mirrored Splunk Finding Event integration parameter to close the Splunk finding when the corresponding Cortex XSOAR incident is closed.

Mapping fetched incidents using Select Schema

This integration supports the Select Schema feature of XSOAR 6.0 by providing the get-mapping-fields command.
When creating a new field mapping for fetched incidents, the Pull Instances option retrieves current incidents which can be clicked to visually map fields.
The Select Schema option retrieves possible objects, even if they are not the next objects to be fetched, or have not been triggered in the past 24 hours.
This enables you to map fields for an incident without having to generate a new alert or incident just for the sake of mapping.
The get-mapping-fields command can be executed in the Playground to test and review the list of sample objects that are returned under the current configuration.

To use this feature, you must set several integration instance parameters:

  • Fetch events query - The query used for fetching new incidents. Select Schema will run a modified version of this query to get the object samples, so it is important to have the correct query here.
  • First fetch timestamp - The time scope of objects to be pulled. You may choose to go back further in time to include samples for alert types that haven’t triggered recently - so long as your Splunk server can handle the more intensive Search Job involved.

Splunk non-Enterprise Security Users

Configure Splunk to Produce Alerts for SplunkPy v2 for non-ES Splunk Users

It is recommended that Splunk is configured to produce basic alerts that the SplunkPy v2 integration can ingest, by creating a summary index in which alerts are stored. The SplunkPy v2 integration can then query that index for incident ingestion. It is not recommended to use the Cortex XSOAR/XSIAM application with Splunk for routine event consumption because this method is not able to be monitored and is not scalable.

  1. Create a summary index in Splunk. For more information, click here.
  2. Build a query to return relevant alerts.
    image
  3. Identify the fields list from the Splunk query and save it to a local file.
    image
  4. Define a search macro to capture the fields list that you saved locally. For more information, click here.
    Use the following naming convention: (demisto_fields_{type}).
    image
    image
  5. Define a scheduled search, the results of which are stored in the summary index. For more information about scheduling searches, click here.
    image
  6. In the Summary indexing section, select the summary index, and enter the {key:value} pair for Cortex XSOAR/XSIAM classification.
    image
  7. Configure the incident type in Cortex XSOAR by navigating to Settings > Advanced > Incident Types. Note: In the example, Splunk Generic is a custom incident type.
    image
  8. Configure the classification. Make sure that your non ES incident fields are associated with your custom incident type.
    1. Navigate to Settings > Integrations > Classification & Mapping.
    2. Click your classifier.
    3. Select your instance.
    4. Click the fetched data.
    5. Drag the value to the appropriate incident type.
      image
  9. Configure the mapping. Make sure to map your non ES fields accordingly and make sure that these incident fields are associated with their custom incident type.
    1. Navigate to Settings > Integrations > Classification & Mapping.
    2. Click your mapper.
    3. Select your instance.
    4. Click the Choose data path link for the field you want to map.
    5. Click the data from the Splunk fields to map it to Cortex XSOAR.
      image
  10. (Optional) Create custom fields.
  11. Build a playbook and assign it as the default for this incident type.

Constraints

The following features are not supported in non-ES (Enterprise Security) Splunk.

  • Incident Mirroring
  • Enrichment.
  • Content in the Splunk content pack (such as mappers, layout, playbooks, incident fields, and the incident type). Therefore, you will need to create your own content. See the Cortex XSOAR Administrator’s Guide for information.

Create KV Store

KV Store stores your data as key-value pairs in collections. It provides a way to save and retrieve data within your Splunk apps. The following is an example for how to create a KV Store.

  1. In Cortex XSOAR/XSIAM, create a new KV Store.

    !splunk-kv-store-collection-create kv_store_name=”<kv_store_name>“

    For example:

    !splunk-kv-store-collection-create kv_store_name=”test_kvstore”

  2. Define the fields and their type in the KV Store.

    !splunk-kv-store-collection-config kv_store_collection_name=”<kv_store_name>” kv_store_fields=”field.<field-name>=<type>,index.<index-name>=<type>,field.<field-name-or-index>=<type>,…“

    For example:

    !splunk-kv-store-collection-config kv_store_collection_name=”test_kvstore” kv_stre_fields=”field.src=cidr,field.t=number,field.description=string”

    Note: To see the fields in Splunk, you must install the Splunk App for Lookup File Editing app in Splunk. For more information, see Define a KV Store lookup in Splunk.

  3. Make the KV Store usable in Splunk queries.

    !splunk-kv-store-collection-create-transform kv_store_collection_name=<kv-store-name> supported_fields=<field-name-or-index>,<field-name-or-index>,<field-name-or-index>,…

    For example:

    !splunk-kv-store-collection-create-transform kv_store_collection_name=<test_kvstore> supported_fields=src,t,description

    Note: If no value is specified, the KV Store collection configuration will be used.

Add data to the KV Store

To add data to the fields in the KV Store, run the following command:

!splunk-kv-store-collection-add-entries kv_store_data=”{"<field-name-or-index>": "<value>", "<field-name-or-index>": "<value>", "<field-name-or-index>": "<value>"…}”

For example:

!splunk-kv-store-collection-add-entries kv_store_data=”{"src": "88.88.88.88", "t": 9, "description": This is the description"}”

Commands

You can execute these commands from the CLI, as part of an automation, or in a playbook.
After you successfully execute a command, a DBot message appears in the War Room with the command details.

splunk-results


Returns the results of a previous Splunk search. This command can be used in conjunction with the splunk-job-create command.

Base Command

splunk-results

Input
Argument Name Description Required
sid The ID of the search for which to return results. Required
limit The maximum number of returned results per search. To retrieve all results, enter “0” (not recommended). Optional
Context Output

There is no context output for this command.

Command Example

!splunk-results sid="1566221331.1186" limit="200"

splunk-search


Searches Splunk for events. For human readable output, the table command is supported in the query argument. For example, query=" * | table field1 field2 field3" will generate a table with field1, field2, and field3 as headers.

Base Command

splunk-search

Input
Argument Name Description Required
query The Splunk search language string to execute. For example, “index=* | head 3”. Required
earliest_time Specifies the earliest time in the time range to search. The time string can be a UTC time (with fractional seconds), a relative time specifier (to now), or a formatted time string. The default is 1 week ago, in the format “-7d”. You can also specify time in the format: 2014-06-19T12:00:00.000-07:00. Optional
latest_time Specifies the latest time in the time range to search. The time string can be a UTC time (with fractional seconds), a relative time specifier (to now), or a formatted time string. For example: “2014-06-19T12:00:00.000-07:00” or “-3d” (for 3 days ago). Optional
event_limit The maximum number of events to return. The default is 100. If “0” is selected, all results are returned. Optional
app The string that contains the application namespace in which to restrict searches. Optional
batch_limit The maximum number of returned results to process at a time. For example, if 100 results are returned, and you specify a batch_limit of 10, the results will be processed 10 at a time over 10 iterations. This does not affect the search or the context and outputs returned. In some cases, specifying a batch_size enhances search performance. If you think that the search execution is suboptimal, it is recommended to try several batch_size values to determine which works best for your search. The default is 25,000. Optional
update_context Determines whether the results will be entered into the context. Optional
polling Use XSOAR built-in polling to retrieve the result when it’s ready. Optional
interval_in_seconds Interval in seconds between each poll. Optional
sid The job sid. Optional
fast_mode Determines whether to retrieve the results in fast mode Optional
Context Output
Path Type Description
Splunk.Result Unknown The results of the Splunk search. The results are a JSON array, in which each item is a Splunk event.
Splunk.JobStatus.SID String ID of the job.
Splunk.JobStatus.Status String Status of the job.
Splunk.JobStatus.TotalResults String The number of events that were returned by the job.
Command Example

!splunk-search query="* | head 3" earliest_time="-1000d"

Note: To display empty columns as well, the following should be added to the query: | fillnull value=

Human Readable Output

Splunk Search results for query: * | head 3

_bkt _cd _indextime _kv _raw _serial _si _sourcetype _time host index linecount source sourcetype splunk_server
main~445~66D21DF4-F4FD-4886-A986-82E72ADCBFE9 445:897774 1585462906 1 InsertedAt=”2020-03-29 06:21:43”; EventID=”837005”; EventType=”Application control”; Action=”None”; ComputerName=”ACME-code-007”; ComputerDomain=”DOMAIN”; ComputerIPAddress=”127.0.0.1”; EventTime=”2020-03-29 06:21:43”; EventTypeID=”5”; Name=”LogMeIn”; EventName=”LogMeIn”; UserName=””; ActionID=”6”; ScanTypeID=”200”; ScanType=”Unknown”; SubTypeID=”23”; SubType=”Remote management tool”; GroupName=””;\u003cbr\u003e 2 ip-172-31-44-193, main sophos:appcontrol 2020-03-28T23:21:43.000-07:00 127.0.0.1 main 2 eventgen sophos:appcontrol ip-172-31-44-193

splunk-submit-event


Creates a new event in Splunk.

Base Command

splunk-submit-event

Input
Argument Name Description Required
index The Splunk index to which to push the data. Run the splunk-get-indexes command to get all of the indexes. Required
data The new event data to push. Can be any string. Required
sourcetype The event source type. Required
host The event host. Can be “Local” or “120.0.0.1”. Required
Context Output

There is no context output for this command.

Command Example

!splunk-submit-event index="main" data="test" sourcetype="demisto-ci" host="localhost"

Human Readable Output

image

splunk-get-indexes


Prints all Splunk index names.

Base Command

splunk-get-indexes

Input

There are no input arguments for this command.

Context Output

There is no context output for this command.

Command Example

!splunk-get-indexes extend-context="indexes="

Human Readable Output

image

splunk-finding-event-edit


Update an existing finding event in Splunk ES.

Base Command

splunk-finding-event-edit

Input
Argument Name Description Required
event_ids A comma-separated list of event IDs of finding events. Required
owner The Splunk user to assign to the finding events. Optional
note The Note to add to the finding events. Optional
urgency The urgency of the finding events. Optional
status The status of the finding events. Can be one of the default options: Unassigned, Assigned, In Progress, Pending, Resolved, Closed Or you can specif another custom status. Optional
disposition The disposition of the finding events. Can be one of the default options: Unassigned, True Positive - Suspicious Activity, Benign Positive - Suspicious But Expected, False Positive - Incorrect Analytic Logic, False Positive - Inaccurate Data, Other, Undetermined. Or you can specify custom dispositions as disposition:# where # is the number of the custom configured disposition on Splunk. Optional
finding_time The time associated with the finding event (e.g., the _time field of the finding). Use this argument only when the command fails with error code MC_01202 or MC_0210, which indicate that the finding event time is required to complete the update. Optional
Context Output

There is no context output for this command.

Command Example

!splunk-finding-event-edit event_ids=${incident.eventid} note="note from `splunk-finding-event-edit` command"

Human Readable Output

image

splunk-update-investigation


Updates existing investigations in Splunk ES. Supports updating fields such as owner, status, urgency, disposition, name, and description, adding a note, and appending finding IDs to the investigation.

Base Command

splunk-update-investigation

Input

Argument Name Description Required
event_ids A comma-separated list of investigation IDs. Required
owner A Splunk user to assign to the investigations. Optional
note Note to add to the investigation. Optional
disposition Disposition of the investigation. If more options exist on the server, specifying the disposition as disposition:# will work in place of choosing one of the default values from the list. Possible values are: Unassigned, True Positive - Suspicious Activity, Benign Positive - Suspicious But Expected, False Positive - Incorrect Analytic Logic, False Positive - Inaccurate Data, Other, Undetermined. Optional
status Investigation status. Possible values are: New, Unassigned, In progress, Pending, Resolved, Closed. Optional
urgency Investigation urgency. Possible values are: critical, high, medium, low, informational. Optional
name Updated name for the investigation. Optional
description Updated description for the investigation. Optional
findings Comma-separated list of finding IDs to add (append) to the investigation. Only allowed when exactly one investigation ID is provided in event_ids. Optional
finding_times The list of times for findings added to the investigation. Value can be in relative, ISO, or epoch time. Ignored when findings is not provided. Only allowed when exactly one investigation ID is provided in event_ids. Optional

Context Output

There is no context output for this command.

Command example

!splunk-update-investigation event_ids="ES00019"

Human Readable Output

Splunk ES events updated successfully:
Successfully updated Splunk ES event ES-00019

Command example (update name, description and append findings)

!splunk-update-investigation event_ids="ES00019" name="New investigation name" description="Updated description" findings="FND-1,FND-2" finding_times="1700000001,1700000002"

Human Readable Output

Splunk ES events updated successfully:
Successfully updated Splunk ES event ES-00019

splunk-job-create


Creates a new search job in Splunk.

Base Command

splunk-job-create

Input
Argument Name Description Required
query The Splunk search language string to execute. For example, “index=* | head 3”. Required
app The string that contains the application namespace in which to restrict searches. Optional
Context Output
Path Type Description
Splunk.Job Unknown The SID of the created job.
Command Example

!splunk-job-create query="index=* | head 3"

Context Example
{
    "Splunk.Job": "1566221733.1628"
}
Human Readable Output

image

splunk-parse-raw


Parses the raw part of the event.

Base Command

splunk-parse-raw

Input
Argument Name Description Required
raw The raw data of the Splunk event (string). Optional
Context Output
Path Type Description
Splunk.Raw.Parsed unknown The raw event data (parsed).
Command Example

!splunk-parse-raw

splunk-submit-event-hec


Sends events Splunk. if batch_event_data or entry_id arguments are provided then all arguments related to a single event are ignored.

Base Command

splunk-submit-event-hec

Input
Argument Name Description Required
event The event payload key-value pair. An example string: “event”: “Access log test message.”. Optional
fields Fields for indexing that do not occur in the event payload itself. Accepts multiple, comma-separated, fields. Optional
index The index name. Optional
host The hostname. Optional
source_type The user-defined event source type. Optional
source The user-defined event source. Optional
time The epoch-formatted time. Optional
batch_event_data A batch of events to send to Splunk. For example, {"event": "something happened at 14/10/2024 12:29", "fields": {"severity": "INFO", "category": "test2, test2"}, "index": "index0","sourcetype": "sourcetype0","source": "/example/something" } {"event": "something happened at 14/10/2024 13:29", "index": "index1", "sourcetype": "sourcetype1","source": "/example/something", "fields":{ "fields" : "severity: INFO, category: test2, test2"}}. If provided, the arguments related to a single event and the entry_id argument are ignored. Optional
batch_event_data A batch of events to send to splunk. For example, {"event": "something happened at 14/10/2024 12:29", "fields": {"severity": "INFO", "category": "test2, test2"}, "index": "index0","sourcetype": "sourcetype0","source": "/example/something" } {"event": "something happened at 14/10/2024 13:29", "index": "index1", "sourcetype": "sourcetype1","source": "/exeample/something", "fields":{ "fields" : "severity: INFO, category: test2, test2"}}. If provided, the arguments related to a single event and the entry_id argument are ignored. Optional
entry_id The entry id in Cortex XSOAR of the file containing a batch of events. Content of the file should be valid batch event’s data, as it would be provided to the batch_event_data. If provided, the arguments related to a single event are ignored. Optional
Batched events description

This command allows sending events to Splunk, either as a single event or a batch of multiple events.
To send a single event: Use the event, fields, host, index, source, source_type, and time arguments.
To send a batch of events, there are two options, either use the batch_event_data argument or use the entry_id argument (for a file uploaded to Cortex XSOAR).
Batch format requirements: The batch must be a single string containing valid dictionaries, each representing an event. Events should not be separated by commas. Each dictionary should include all necessary fields for an event. For example: {"event": "event occurred at 14/10/2024 12:29", "fields": {"severity": "INFO", "category": "test1"}, "index": "index0", "sourcetype": "sourcetype0", "source": "/path/event1"} {"event": "event occurred at 14/10/2024 13:29", "index": "index1", "sourcetype": "sourcetype1", "source": "/path/event2", "fields": {"severity": "INFO", "category": "test2"}}.
This formatted string can be passed directly via batch_event_data, or, if saved in a file, the file can be uploaded to Cortex XSOAR, and the entry_id (e.g., ${File.[4].EntryID}) should be provided.

Context Output

There is no context output for this command.

Command Example

!splunk-submit-event-hec event="something happened" fields="severity: INFO, category: test, test1" source_type=access source="/var/log/access.log"

Human Readable Output

The event was sent successfully to Splunk.

splunk-job-status


Returns the status of a job.

Base Command

splunk-job-status

Input
Argument Name Description Required
sid Comma-separated list of job IDs for which to retrieve the statuses. Required
Context Output
Path Type Description
Splunk.JobStatus.SID Unknown The ID of the job.
Splunk.JobStatus.Status Unknown The status of the job.
Command Example

!splunk-job-status sid=1234.5667

Context Example
Splank.JobStatus = {
    'SID': 1234.5667,
    'Status': DONE
}
Human Readable Output

image

get-mapping-fields


Gets one sample alert per alert type. Used only for creating a mapping with Select Schema.

Base Command

get-mapping-fields

Input

There are no input arguments for this command.

Context Output

There is no context output for this command.

Command Example

!get-mapping-fields using="SplunkPy_v2_instance" raw-response="true"

Human Readable Output
{
    "Access - Brute Force Access Behavior Detected - Rule": {
        "_bkt": "notable~712~66D21DF4-F4FD-4886-A986-82E72ADCBFE9",
        "_cd": "712:21939",
        "_indextime": "1598464820",
        "_serial": "0",
        "_si": [
            "ip-1-1-1-1",
            "notable"
        ],
        "_sourcetype": "stash",
        "_time": "2020-08-26T11:00:20.000-07:00",
        "host": "ip-1-1-1-1",
        "host_risk_object_type": "system",
        "host_risk_score": "0",
        "index": "notable",
        "linecount": "1",
        "priority": "unknown",
        "risk_score": "460",
        "rule_description": "Access - Brute Force Access Behavior Detected - Rule",
        "rule_name": "Access - Brute Force Access Behavior Detected - Rule",
        "rule_title": "Access - Brute Force Access Behavior Detected - Rule",
        "security_domain": "Access - Brute Force Access Behavior Detected - Rule",
        "severity": "unknown",
        "source": "Access - Brute Force Access Behavior Detected - Rule",
        "sourcetype": "stash",
        "splunk_server": "ip-1-1-1-1",
        "src": "1.1.1.1",
        "src_risk_object_type": "system",
        "src_risk_score": "460",
        "urgency": "low"
    },
    "Access - Excessive Failed Logins - Rule": {
        "_bkt": "notable~712~66D21DF4-F4FD-4886-A986-82E72ADCBFE9",
        "_cd": "712:21515",
        "_indextime": "1598460945",
        "_serial": "22",
        "_si": [
            "ip-1-1-1-1",
            "notable"
        ],
        "_sourcetype": "stash",
        "_time": "2020-08-26T09:55:45.000-07:00",
        "host": "ip-1-1-1-1",
        "host_risk_object_type": "system",
        "host_risk_score": "0",
        "index": "notable",
        "linecount": "1",
        "priority": "unknown",
        "risk_score": "380",
        "rule_description": "Access - Excessive Failed Logins - Rule",
        "rule_name": "Access - Excessive Failed Logins - Rule",
        "rule_title": "Access - Excessive Failed Logins - Rule",
        "security_domain": "Access - Excessive Failed Logins - Rule",
        "severity": "unknown",
        "source": "Access - Excessive Failed Logins - Rule",
        "sourcetype": "stash",
        "splunk_server": "ip-1-1-1-1",
        "src": "1.1.1.1",
        "src_risk_object_type": "system",
        "src_risk_score": "380",
        "urgency": "low"
}

splunk-kv-store-collection-create


Creates a new KV store table.

Base Command

splunk-kv-store-collection-create

Input

Argument Name Description Required
kv_store_name The name of the KV store collection. Required
app_name The name of the Splunk application in which to create the KV store. The default is “search”. Required

Context Output

There is no context output for this command.

Command Example

!splunk-kv-store-collection-create app_name=search kv_store_name=demisto_store

Human Readable Output

KV store collection search created successfully

splunk-kv-store-collection-config


Configures the KV store fields.

Base Command

splunk-kv-store-collection-config

Input

Argument Name Description Required
kv_store_collection_name The name of the KV store collection. Required
kv_store_fields The list of names and value types to define the KV store collection scheme, e.g., id=number, name=string, address=string.
Required
app_name The name of the Splunk application that contains the KV store collection. The default is “search”. Required

Context Output

There is no context output for this command.

Command Example

!splunk-kv-store-collection-config app_name=search kv_store_collection_name=demisto_store kv_store_fields=addr=string

Human Readable Output

KV store collection search configured successfully

splunk-kv-store-collection-add-entries


Adds objects to a KV store utilizing the batch-save API.

Base Command

splunk-kv-store-collection-add-entries

Input

Argument Name Description Required
kv_store_data The data to add to the KV store collection, according to the collection JSON format, e.g., {“name”: “Splunk HQ”, “id”: 123, “address”: { “street”: “250 Brannan Street”, “city”: “San Francisco”, “state”: “CA”, “zip”: “94107”}} Required
kv_store_collection_name The name of the KV store collection. Required
indicator_path The path to the indicator value in kv_store_data. Optional
app_name The name of the Splunk application that contains the KV store collection. The default is “search”. Required

Context Output

There is no context output for this command.

Command Example

!splunk-kv-store-collection-add-entries app_name=search kv_store_collection_name=demisto_store kv_store_data="{\"addr\": \"0.0.0.0\"}" indicator_path=addr

Human Readable Output

Data added to demisto_store

splunk-kv-store-collections-list


Lists all collections for the specified application.

Base Command

splunk-kv-store-collections-list

Input

Argument Name Description Required
app_name The name of the Splunk application in which to create the KV store. The default is “search”. Required

Context Output

Path Type Description
Splunk.CollectionList String List of collections.

Command Example

!splunk-kv-store-collections-list app_name=search

Context Example

{
    "Splunk": {
        "CollectionList": [
            "autofocus_tags",
            "files"
        ]
    }
}

Human Readable Output

list of collection names search

name
autofocus_tags
files

splunk-kv-store-collection-data-list


Lists all data within a specific KV store collection or collections.

Base Command

splunk-kv-store-collection-data-list

Input

Argument Name Description Required
app_name The name of the Splunk application that contains the KV store collection. Default is search. Required
kv_store_collection_name A comma-separated list of KV store collections. Required
limit Maximum number of records to return. The default is 50. Optional

Context Output

Path Type Description
Splunk.KVstoreData Unknown An array of collection names. Each collection name will have an array of values, e.g., Splunk.KVstoreData.<collection_name> is a list of the data in the collection.

Command Example

!splunk-kv-store-collection-data-list app_name=search limit=3 kv_store_collection_name=demisto_store

Context Example

{
    "Splunk": {
        "KVstoreData": {
            "demisto_store": [
                {
                    "_key": "5f4e2e9c097d9e6749453536",
                    "_user": "nobody",
                    "addr": "0.0.0.0"
                }
            ]
        }
    }
}

Human Readable Output

list of collection values demisto_store

_key _user addr
5f4e2e9c097d9e6749453536 nobody 0.0.0.0

splunk-kv-store-collection-data-delete


Deletes all data within the specified KV store collection or collections.

Base Command

splunk-kv-store-collection-data-delete

Input

Argument Name Description Required
app_name The name of the Splunk application that contains the KV store collection. For example, “search”.” Required
kv_store_collection_name A comma-separated list of KV store collections. Required

Context Output

There is no context output for this command.

Command Example

!splunk-kv-store-collection-data-delete app_name=search kv_store_collection_name=demisto_store

Human Readable Output

The values of the demisto_store were deleted successfully

splunk-kv-store-collection-delete


Deletes the specified KV stores.

Base Command

splunk-kv-store-collection-delete

Input

Argument Name Description Required
app_name The name of the Splunk application that contains the KV store. The default is “search”. Required
kv_store_name A comma-separated list of KV stores. Required

Context Output

There is no context output for this command.

Command Example

!splunk-kv-store-collection-delete app_name=search kv_store_name=demisto_store

Human Readable Output

The following KV store demisto_store were deleted successfully

splunk-kv-store-collection-search-entry


Searches for specific objects in a store. Search can be a basic key-value pair or a full query.

Base Command

splunk-kv-store-collection-search-entry

Input

Argument Name Description Required
app_name The name of the Splunk application that contains the KV store collection. The default is “search”. Required
kv_store_collection_name The name of the KV store collection Required
key The key name to search in the store. If the query argument is used, this argument will be ignored. Optional
value The value to search in the store. If the query argument is used, this argument will be ignored. Optional
query Complex query to search in the store with operators such as “and”, “or”, “not”, etc. For more information see the Splunk documentation: https://docs.splunk.com/Documentation/Splunk/8.0.3/RESTREF/RESTkvstore Optional

Context Output

Path Type Description
Splunk.KVstoreData Unknown An array of collection names. Each collection name will have an array of values, e.g., Splunk.KVstoreData.<collection_name> is a list of the data in the collection.

Command Example

!splunk-kv-store-collection-search-entry app_name=search kv_store_collection_name=demisto_store key=addr value=0.0.0.0

Context Example

{
    "Splunk": {
        "KVstoreData": {
            "demisto_store": [
                {
                    "_key": "5f4e2e9c097d9e6749453536",
                    "_user": "nobody",
                    "addr": "0.0.0.0"
                }
            ]
        }
    }
}

Human Readable Output

list of collection values demisto_store

_key _user addr
5f4e2e9c097d9e6749453536 nobody 0.0.0.0

splunk-kv-store-collection-delete-entry


Deletes the specified object in store. Search can be a basic key-value pair or a full query.

Base Command

splunk-kv-store-collection-delete-entry

Input

Argument Name Description Required
app_name The name of the Splunk application that contains the KV store collection. The default is “search”. Required
kv_store_collection_name The name of the KV store collection. Required
indicator_path The path to the indicator value in kv_store_data. Optional
key The key name to search in the store. If the query argument is used, this argument will be ignored. Optional
value The value to search in the store. If the query argument is used, this argument will be ignored. Optional
query Complex query to search in the store with operators such as “and”, “or”, “not”, etc.
For more information see the Splunk documentation: https://docs.splunk.com/Documentation/Splunk/8.0.3/RESTREF/RESTkvstore
Optional

Context Output

There is no context output for this command.

Command Example

!splunk-kv-store-collection-delete-entry app_name=search kv_store_collection_name=demisto_store key=addr value=0.0.0.0 indicator_path=addr

Human Readable Output

The values of the demisto_store were deleted successfully

get-modified-remote-data


Gets the list of finding events that were modified since the last update. This command should be used for debugging purposes, and is available from Cortex XSOAR version 6.1.

Base Command

get-modified-remote-data

Input

Argument Name Description Required
lastUpdate ISO format date with timezone, e.g., 2021-02-09T16:41:30.589575+02:00. The incident is only returned if it was modified after the last update time. Required

Context Output

There is no context output for this command.

splunk-reset-enriching-fetch-mechanism


Resets the enriching fetch mechanism.

Base Command

splunk-reset-enriching-fetch-mechanism

Input

There are no input arguments for this command.

Context Output

There is no context output for this command.

Command Example


#### Human Readable Output

>Enriching fetch mechanism was reset successfully.

### splunk-get-username-by-xsoar-user

***
Returns the Splunk's username matching the given Cortex XSOAR's username.

#### Base Command

`splunk-get-username-by-xsoar-user`

#### Input

| **Argument Name** | **Description** | **Required** |
| --- | --- | --- |
| xsoar_username | Cortex XSOAR username to match in Splunk's usernames records. | Required |

#### Context Output

| **Path** | **Type** | **Description** |
| --- | --- | --- |
| Splunk.UserMapping.XsoarUser | String | Cortex XSOAR user mapping. |
| Splunk.UserMapping.SplunkUser | String | Splunk user mapping. |

#### Command Example

```!splunk-get-username-by-xsoar-user xsoar_username=admin```

#### Context Example

{
“Splunk”: {
“UserMapping”: [
{
“SplunkUser”: “unassigned”,
“XsoarUser”: “admin”
}
]
}
}
```

Human Readable Output

Xsoar-Splunk Username Mapping

Xsoar User Splunk User
admin unassigned

splunk-kv-store-collection-create-transform


Creates the KV store collection transform.

Base Command

splunk-kv-store-collection-create-transform

Input

Argument Name Description Required
kv_store_collection_name The name of the KV store collection. Required
supported_fields A comma-delimited list of the fields supported by the collection, e.g., _key,id,name,address. If no value is specified, the KV Store collection configuration will be used. Optional
app_name The name of the Splunk application that contains the KV store collection. Default is search. Required

Context Output

There is no context output for this command.

splunk-job-share


Change job settings to share its results to all Splunk users, and change its TTL.

Base Command

splunk-job-share

Input

Argument Name Description Required
sid Comma-separated list of job IDs to share. Required
ttl Time in seconds for the job’s expiry time. Default is 1800. Optional

Context Output

There is no context output for this command.

splunk-investigation-create


Creates a new investigation in Splunk Enterprise Security.

Base Command

splunk-investigation-create

Input

Argument Name Description Required
name The name of the investigation to be created. Required
description The description of the investigation to be created. Optional
investigation_type The type of the investigation to be created (for example, default). Optional
status The status of the investigation to be created. Defaults to the out-of-the-box Splunk ES status labels; custom statuses are also supported by typing the status ID or label. Possible values are: New, In Progress, Pending, Resolved, Closed. Optional
disposition The disposition of the investigation to be created. Defaults to the out-of-the-box Splunk ES disposition labels; custom dispositions are also supported by typing the disposition ID or label. Possible values are: Undetermined, True Positive - Suspicious Activity, Benign Positive - Suspicious But Expected, False Positive - Incorrect Analytic Logic, False Positive - Inaccurate Data, Other. Optional
owner The Splunk user to assign as the owner of the investigation. Optional
urgency The urgency of the investigation to be created. Possible values are: informational, low, medium, high, critical, unknown. Optional
sensitivity The sensitivity of the investigation to be created. Possible values are: White, Green, Amber, Red, Unassigned. Optional

Context Output

Path Type Description
Splunk.Investigation.investigation_guid String The ID (GUID) of the investigation that was created.

splunk-investigation-list


Lists investigations from Splunk Enterprise Security.

Base Command

splunk-investigation-list

Input

Argument Name Description Required
investigation_ids A comma-separated list of investigation IDs (GUID or display ID such as ES-00001) to retrieve. Optional
limit The maximum number of investigations to return on the page. Maximum is 100. Default is 20. Optional
offset The pagination offset used together with limit to specify the starting point of the returned results. Optional
sort The sort expression for the returned investigations (for example, create_time:asc,status:desc). Optional
disposition A comma-separated list of disposition IDs or disposition labels to filter investigations by (for example, disposition:1,Undetermined). Optional
status A comma-separated list of status IDs or status labels to filter investigations by (for example, New,In progress). Optional
owner A comma-separated list of owners to filter investigations by. Optional
urgency A comma-separated list of urgency values to filter investigations by (for example, medium,high,critical). Valid values are informational, low, medium, high, critical, or unknown. Optional
sensitivity A comma-separated list of sensitivity values to filter investigations by (for example, Amber,Red). Valid values are White, Green, Amber, Red, or Unassigned. Optional
create_time_min The minimum (epoch) time during which investigations were created. Optional
create_time_max The maximum (epoch) time during which investigations were created. Optional
update_time_min The minimum (epoch) time during which investigations were updated. Optional
update_time_max The maximum (epoch) time during which investigations were updated. Optional

Context Output

Path Type Description
Splunk.Investigation.investigation_guid String The ID (GUID) of the investigation.
Splunk.Investigation.investigation_id String The short display ID of the investigation (for example, `ES-00001`).
Splunk.Investigation.name String The name of the investigation.
Splunk.Investigation.description String The description of the investigation.
Splunk.Investigation.investigation_type String The type of the investigation.
Splunk.Investigation.source String The detection that generated the investigation.
Splunk.Investigation.incident_origin String Where the investigation came from (for example, Splunk Enterprise Security or a risk-based alerting finding).
Splunk.Investigation.finding_id String The ID of the originating Splunk Enterprise Security finding.
Splunk.Investigation.disposition String The disposition ID of the investigation.
Splunk.Investigation.disposition_name String The disposition name of the investigation.
Splunk.Investigation.status String The status ID of the investigation.
Splunk.Investigation.status_name String The status name of the investigation.
Splunk.Investigation.owner String The person assigned to the investigation.
Splunk.Investigation.urgency String The urgency of the investigation.
Splunk.Investigation.sensitivity String The sensitivity of the investigation.
Splunk.Investigation.create_time Number The time when the investigation was created (epoch seconds).
Splunk.Investigation.update_time Number The time when the investigation was last updated (epoch seconds).
Splunk.Investigation.mc_create_time Number The time when the finding or investigation was created or imported into Splunk Enterprise Security (epoch seconds).
Splunk.Investigation.count_findings Number The number of findings (or intermediate findings) associated with this investigation or finding-based-detection (FBD) group.
Splunk.Investigation.risk_event_count Number The number of risk events associated with this investigation.
Splunk.Investigation.risk_score Number The maximum risk score for all the findings added to the investigation.
Splunk.Investigation.excluded_finding_ids Unknown A list of finding IDs (or intermediate findings in the finding groups) that are removed from the investigation.
Splunk.Investigation.attachments Unknown An array of file IDs attached directly to the investigation.
Splunk.Investigation.notes Unknown An array of note IDs added directly to the finding or investigation.
Splunk.Investigation.findings.incident_ids Unknown The added finding IDs.
Splunk.Investigation.findings.field_inheritors Unknown The added finding IDs that will inherit this investigation’s owner, status, urgency, sensitivity, and disposition values.
Splunk.Investigation.current_response_plan_phase.phase_id String The ID of the current response plan phase.
Splunk.Investigation.current_response_plan_phase.response_plan_id String The ID of the current response plan.
Splunk.Investigation.response_plans Unknown The array of response plans added to the investigation.
Splunk.Investigation.consolidated_findings Unknown The consolidated list of fields for the findings and all the findings that are added to this investigation.
Splunk.Investigation.finding Unknown The raw data of the originating finding.
Splunk.Investigation.custom_fields Unknown The custom fields in the investigation.
Splunk.Investigation.src Unknown A list of values for the `source` field.
Splunk.Investigation.dest Unknown A list of values for the `destination` field.
Splunk.Investigation.dvc Unknown A list of values for the `device` field.
Splunk.Investigation.orig_host Unknown A list of values for the `host` field.
Splunk.Investigation.src_user Unknown A list of values for the `source user` field.
Splunk.Investigation.user Unknown A list of values for the `user` field.
Splunk.Investigation.risk_object Unknown The list of entities for a finding, finding group, or investigation.
Splunk.Investigation.risk_object_type Unknown The list of risk object types for a finding, finding group, or investigation.

Additional Information

To get the HEC token

  1. Go to the Splunk UI.
  2. Under Settings > Data > Data inputs, click HTTP Event Collector.
    Screen Shot 2020-01-20 at 10 22 50

  3. Click New Token.
  4. Add all the relevant details until done.

For the HTTP Port number:
Click on Global settings (in the HTTP Event Collector page)
Screen Shot 2020-01-20 at 10 27 25

The default port is 8088.

Troubleshooting

Index Validation Issues

In some cases, the Splunk API may not return a complete list of all available indexes. If you try to submit an event to an index that you know exists but the integration reports that it cannot be found, it may be due to this Splunk issue. The integration will log an error message specifying which indexes could not be verified.

Recommended Action:

  1. Verify that the index name is spelled correctly in your request.
  2. If the index exists and is accessible in Splunk but is not found by the integration, please contact Splunk support for assistance, as this is a known limitation with the Splunk API.

Connectivity Issues

If you encounter connectivity issues while using Splunk Cloud within Cortex XSOAR8 or Cortex XSIAM you may receive the following error:

requests.exceptions.ConnectTimeout:
HTTPSConnectionPool(host='<name>.splunkcloud.com', port=8089)
: Max retries exceeded with url: /services/auth/login (Caused by ConnectTimeoutError(<urllib3.connection.HTTPSConnection object at 0x7fc389a4e170>,
 'Connection to <name>.splunkcloud.com timed out. 
(connect timeout=None)'))

To resolve this issue, add the IP addresses of Cortex XSOAR8 or Cortex XSIAM to the Splunk Cloud whitelist.
You can find the relevant IP addresses at:
Cortex XSOAR Administrator Guide
Under Used for communication between Cortex XSOAR and customer resources. Choose the IP address corresponding to your Cortex XSOAR region.

Fetch Issues

If you encounter fetch issues and you have enriching enabled, the issue may be the result of pressing the Reset the "last run" timestamp button.
Note that the way to reset the mechanism is to run the splunk-reset-enriching-fetch-mechanism command.
See here.

Large Search Results

Commands that return large data (such as splunk-search) can cause performance issues in playbooks.

Recommendation: Limit results to approximately 30,000 events, depending on the data size. You can do this through one of the following:

  • Use the event_limit argument (where available).
  • Append | head 30000 directly to your Splunk query.

splunk-configuration-stanza-create


Creates a new stanza (configuration entry) in a Splunk .conf file, optionally with attributes.

Base Command

splunk-configuration-stanza-create

Input

Argument Name Description Required
conf_file The configuration file name (without the .conf extension), for example, “transforms”, “props”, or “inputs”. Required
stanza_name The name of the new stanza to create. Required
key_value_pairs The JSON object string of attributes to set on creation, e.g., {“external_type”: “kvstore”, “collection”: “my_collection”}. If omitted, an empty stanza is created. Optional
app The name of the Splunk application namespace. Default is search. Optional
owner The Access Control List (ACL) owner for the namespace. Default is nobody. Optional

Context Output

There is no context output for this command.

splunk-configuration-stanza-list


Lists the stanzas in a .conf file, or returns the key/value content of a single stanza when stanza_name is provided.

Base Command

splunk-configuration-stanza-list

Input

Argument Name Description Required
conf_file The configuration file name (without the .conf extension). Required
stanza_name The name of the stanza to return. If provided, returns the key/value content of that single stanza. If omitted, returns the list of all stanza names in the configuration file. Optional
app The name of the Splunk application namespace. Default is search. Optional
owner The Access Control List (ACL) owner for the namespace. Default is nobody. Optional
limit The maximum number of stanzas to return when stanza_name is omitted. Range 1-500. Default is 50. Optional

Context Output

Path Type Description
Splunk.ConfigurationStanza.StanzaName String The stanza name.
Splunk.ConfigurationStanza.App String The Splunk app namespace the stanza belongs to.
Splunk.ConfigurationStanza.Owner String The Access Control List (ACL) owner of the stanza.
Splunk.ConfigurationStanza.Sharing String The sharing level of the stanza.
Splunk.ConfigurationStanza.Content Unknown The key/value content of the stanza (only returned when a single stanza_name is provided).

splunk-configuration-file-create


Creates a new, empty configuration (.conf) file in the given Splunk app namespace.

Base Command

splunk-configuration-file-create

Input

Argument Name Description Required
conf_file_name The name of the new configuration file to create (without the .conf extension). Required
app The name of the Splunk application namespace in which to create the file. Default is search. Optional
owner The Access Control List (ACL) owner for the namespace. Default is nobody. Optional

Context Output

There is no context output for this command.

splunk-configuration-file-list


Lists the configuration (.conf) files available in the given Splunk app namespace.

Base Command

splunk-configuration-file-list

Input

Argument Name Description Required
app The name of the Splunk application namespace. Default is search. Optional
owner The Access Control List (ACL) owner for the namespace. Default is nobody. Optional
limit The maximum number of configuration files to return. Range 1-500. Default is 50. Optional

Context Output

Path Type Description
Splunk.ConfigurationFile.FileName String The configuration file name (without the .conf suffix).
Splunk.ConfigurationFile.App String The Splunk app namespace the configuration file belongs to.

splunk-configuration-stanza-delete


Deletes a stanza (configuration entry) from a Splunk .conf file (for example, transforms.conf) via the Splunk REST API configuration endpoints. This is useful for cleaning up KV Store transformations that remain after deleting KV Store records.
To identify the correct stanza for deletion, use the splunk-configuration-file-list and splunk-configuration-stanza-list commands. This command is potentially harmful as it irreversibly deletes a stanza from a .conf file.

Base Command

splunk-configuration-stanza-delete

Input

Argument Name Description Required
conf_file The target configuration file name (without the .conf extension). Required
stanza_name The name of the stanza to be removed. Required
app The name of the Splunk application namespace. Default is search. Optional
owner The Access Control List (ACL) owner for the namespace. Default is nobody. Optional

Context Output

There is no context output for this command.

splunk-configuration-stanza-update


Updates (upserts) attributes on an existing stanza in a Splunk .conf file.

Base Command

splunk-configuration-stanza-update

Input

Argument Name Description Required
conf_file The configuration file name (without the .conf extension). Required
stanza_name The name of the existing stanza to update. Required
key_value_pairs The JSON object string of attributes to upsert, e.g., {“attribute_1”: “A_updated”, “new_attr”: “X”}. Existing keys are overwritten; missing keys are left untouched (Splunk’s properties endpoint is upsert-only). Required
app The name of the Splunk application namespace. Default is search. Optional
owner The Access Control List (ACL) owner for the namespace. Default is nobody. Optional

Context Output

There is no context output for this command.

<~PLATFORM>

License Requirements

The following configuration parameters require one of these licenses: Cortex XSIAM or Agentix:

  • Fetch incidents

</~PLATFORM>

Configuration parameters

  • server_url — Server URL (required)
  • authentication — (required)
  • fetchQuery — Fetch events query
  • max_fetch — Fetch Limit (Max.- 200, Recommended less than 50)
  • isFetch — Fetch incidents
  • incidentType — Incident type
  • parseFindingEventsRaw — Parse Raw Part of Finding Events
  • replaceKeys — Replace with Underscore in Incident Fields
  • first_fetch — First fetch timestamp (<number> <time unit>, e.g., 12 hours, 7 days, 3 months, 1 year)
  • fetch_event_types — Event types to fetch
  • investigations_first_fetch — First fetch timestamp (Investigations)
  • investigations_fetch_query — Investigations fetch query
  • investigations_max_fetch — Maximum investigations per fetch
  • extractFields — Extract Fields - CSV fields that will be parsed out of raw finding events
  • mirror_direction — Incident Mirroring Direction
  • close_incident — Close Mirrored Cortex XSOAR Incidents (Incoming Mirroring)
  • close_extra_labels — Additional Splunk status labels to close on mirror (Incoming Mirroring)
  • close_end_status_statuses — Enable Splunk statuses marked as "End Status" to close on mirror (Incoming Mirroring)
  • close_finding — Close Mirrored Splunk ES Events (Outgoing Mirroring)
  • unsecure — Trust any certificate (not secure)
  • proxy — Use system proxy settings
  • app — The app context of the namespace
  • enabled_enrichments — Enrichment Types
  • asset_enrich_lookup_tables — Asset enrichment lookup tables
  • identity_enrich_lookup_tables — Identity enrichment lookup tables
  • enrichment_timeout — Enrichment Timeout (Minutes)
  • num_enrichment_events — Number of Events Per Enrichment Type
  • cred_hec_token
  • hec_url — HEC BASE URL (e.g: https://localhost:8088 or https://example.splunkcloud.com/).
  • extensive_logs — Advanced: Extensive logging (for debugging purposes). Do not use this option unless advised otherwise.
  • finding_time_source — Advanced: Time type to use when fetching events
  • occurrence_look_behind — Advanced: Fetch backwards window for the events occurrence time (minutes)
  • unique_id_fields — Advanced: Unique ID Fields
  • userMapping — Enable user mapping
  • user_map_lookup_name — Users Lookup table name
  • xsoar_user_field — XSOAR user key
  • splunk_user_field — SPLUNK user key
  • note_tag_from_splunk — Note tag from Splunk
  • note_tag_to_splunk — Note tag to Splunk
  • incidentFetchInterval — Incidents Fetch Interval

Commands (33)

  • get-mapping-fields

    Query Splunk to retrieve a list of sample alerts by alert type. Used for mapping fetched incidents through the Get Schema option.

  • get-modified-remote-data

    Gets the list of finding events that were modified since the last update. This command should be used for debugging purposes, and is available from Cortex XSOAR version 6.1.

  • splunk-configuration-file-create

    Creates a new, empty configuration (.conf) file in the given Splunk app namespace.

  • splunk-configuration-file-list

    Lists the configuration (.conf) files available in the given Splunk app namespace.

  • splunk-configuration-stanza-create

    Creates a new stanza (configuration entry) in a Splunk .conf file, optionally with attributes.

  • splunk-configuration-stanza-delete

    Deletes a stanza (configuration entry) from a Splunk .conf file (for example, transforms.conf) via the Splunk REST API configuration endpoints. This is useful for cleaning up KV Store transformations that remain after deleting KV Store records. To identify the correct stanza for deletion, use the splunk-configuration-file-list and splunk-configuration-stanza-list commands. This command is potentially harmful as it irreversibly deletes a stanza from a .conf file.

  • splunk-configuration-stanza-list

    Lists the stanzas in a .conf file, or returns the key/value content of a single stanza when stanza_name is provided.

  • splunk-configuration-stanza-update

    Updates (upserts) attributes on an existing stanza in a Splunk .conf file.

  • splunk-finding-event-edit

    Updates existing finding events in Splunk ES.

  • splunk-get-indexes

    Prints all Splunk index names.

  • splunk-get-username-by-xsoar-user

    Returns the Splunk's username matching the given Cortex XSOAR's username.

  • splunk-investigation-create

    Creates a new investigation in Splunk Enterprise Security.

  • splunk-investigation-list

    Lists investigations from Splunk Enterprise Security.

  • splunk-job-create

    Creates a new search job in Splunk.

  • splunk-job-share

    Change job settings to share its results to all Splunk users, and change its TTL.

  • splunk-job-status

    Returns the status of a job.

  • splunk-kv-store-collection-add-entries

    Adds objects to a KV store utilizing the batch-save API.

  • splunk-kv-store-collection-config

    Configures the KV store fields.

  • splunk-kv-store-collection-create

    Creates a new KV store table.

  • splunk-kv-store-collection-create-transform

    Creates the KV store collection transform.

  • splunk-kv-store-collection-data-delete

    Deletes all data within the specified KV store collection or collections.

  • splunk-kv-store-collection-data-list

    Lists all data within a specific KV store collection or collections.

  • splunk-kv-store-collection-delete

    Deletes the specified KV stores.

  • splunk-kv-store-collection-delete-entry

    Deletes the specified object in store. The search can be a basic key-value pair or a full query.

  • splunk-kv-store-collection-search-entry

    Searches for specific objects in a store. The search can be a basic key-value pair or a full query.

  • splunk-kv-store-collections-list

    Lists all collections for the specified application.

  • splunk-parse-raw

    Parses the raw part of the event.

  • splunk-reset-enriching-fetch-mechanism

    Resets the enrichment mechanism of fetched findings.

  • splunk-results

    Returns the results of a previous Splunk search. You can use this command in conjunction with the splunk-job-create command.

  • splunk-search

    Searches Splunk for events. For human readable output, the table command is supported in the query argument. For example, `query=" * | table field1 field2 field3"` will generate a table with field1, field2, and field3 as headers.

  • splunk-submit-event

    Creates a new event in Splunk.

  • splunk-submit-event-hec

    Sends events to an HTTP Event Collector using the Splunk platform JSON event protocol.

  • splunk-update-investigation

    Updates existing investigations in Splunk ES. Supports updating fields such as owner, status, urgency, disposition, name, and description, adding a note, and appending finding IDs to the investigation. Note that `findings` (and `finding_times`) can only be used when exactly one investigation ID is provided in `event_ids`.

import demistomock as demisto  # noqa: F401
from CommonServerPython import *  # noqa: F401
import hashlib
import io
import json
import re
import time
import urllib.parse
from abc import ABC, abstractmethod
from collections import defaultdict
from dataclasses import dataclass, field
from datetime import datetime, timedelta, UTC
from urllib.parse import urlencode

import dateparser

import pytz
import requests
from packaging.version import InvalidVersion, Version

from splunklib import client, results
from splunklib.binding import AuthenticationError, HTTPError, namespace
from splunklib.data import Record

INTEGRATION_LOG = "SplunkPyV2- "
OUTPUT_MODE_JSON = "json"  # type of response from splunk-sdk query (json/csv/xml)
params = demisto.params()
DEFAULT_ASSET_ENRICH_TABLES = "asset_lookup_by_str,asset_lookup_by_cidr"
DEFAULT_IDENTITY_ENRICH_TABLE = "identity_lookup_expanded"
VERIFY_CERTIFICATE = not bool(params.get("unsecure"))
FETCH_LIMIT = int(params.get("max_fetch")) if params.get("max_fetch") else 50
FETCH_LIMIT = max(min(200, FETCH_LIMIT), 1)
MIRROR_LIMIT = 1000
SPLUNK_INDEXING_TIME = 60
PROBLEMATIC_CHARACTERS = [".", "(", ")", "[", "]"]
REPLACE_WITH = "_"
REPLACE_FLAG = params.get("replaceKeys", False)
PROXIES = handle_proxy()
DEFAULT_DISPOSITIONS = {
    "Unassigned": "disposition:0",
    "True Positive - Suspicious Activity": "disposition:1",
    "Benign Positive - Suspicious But Expected": "disposition:2",
    "False Positive - Incorrect Analytic Logic": "disposition:3",
    "False Positive - Inaccurate Data": "disposition:4",
    "Other": "disposition:5",
    "Undetermined": "disposition:6",
}

DEFAULT_STATUSES = {
    "Unassigned": "0",
    "Assigned": "1",
    "In Progress": "2",
    "Pending": "3",
    "Resolved": "4",
    "Closed": "5",
}

# =========== Mirroring Mechanism Globals ===========
MIRROR_DIRECTION = {"None": None, "Incoming": "In", "Outgoing": "Out", "Incoming And Outgoing": "Both"}
OUTGOING_MIRRORED_FIELDS = ["note", "status", "owner", "urgency", "reviewer", "disposition"]
ES_APP_NAME = "SplunkEnterpriseSecuritySuite"

# === Note Tag Globals ===
NOTE_TAG_TO_SPLUNK = params.get("note_tag_to_splunk", "FROM XSOAR")
NOTE_TAG_FROM_SPLUNK = params.get("note_tag_from_splunk", "FROM SPLUNK")

# =========== Enrichment Mechanism Globals ===========
ENABLED_ENRICHMENTS = params.get("enabled_enrichments", [])

DRILLDOWN_ENRICHMENT = "Drilldown"
ASSET_ENRICHMENT = "Asset"
IDENTITY_ENRICHMENT = "Identity"
SUBMITTED_FINDINGS = "submitted_findings"
EVENT_ID = "event_id"
RULE_ID = "rule_id"
ISO_FORMAT_TZ_AWARE = "%Y-%m-%dT%H:%M:%S.%f%z"  # e.g '2025-12-03T11:53:45.138540+00:00
SPLUNK_ES_EVENT_TYPE_FIELD = "splunk_es_event_type"  # tag emitted in rawJSON, consumed by classifier
INVESTIGATIONS_MAX_LIMIT = 100  # Hard cap enforced by Splunk endpoint
NOT_YET_SUBMITTED_FINDINGS = "not_yet_submitted_findings"
INFO_MIN_TIME = "info_min_time"
INFO_MAX_TIME = "info_max_time"
INCIDENTS = "incidents"
MIRRORED_ENRICHING_FINDINGS = "MIRRORED_ENRICHING_FINDINGS"
PROCESSED_MIRRORED_EVENTS = "processed_mirror_in_events_cache"
DUMMY = "dummy"
ENRICHMENTS = "enrichments"
MAX_HANDLE_FINDINGS = 20
MAX_SUBMIT_FINDINGS = 30
CACHE = "cache"
STATUS = "status"
DATA = "data"
TYPE = "type"
ID = "id"
CREATION_TIME = "creation_time"
QUERY_NAME = "query_name"
QUERY_SEARCH = "query_search"
INCIDENT_CREATED = "incident_created"

DRILLDOWN_REGEX = r'([^\s\$]+)\s*=\s*"?(\$[^\s\$\\]+\$)"?|"?(\$[^\s\$\\]+\$)"?'

ENRICHMENT_TYPE_TO_ENRICHMENT_STATUS = {
    DRILLDOWN_ENRICHMENT: "successful_drilldown_enrichment",
    ASSET_ENRICHMENT: "successful_asset_enrichment",
    IDENTITY_ENRICHMENT: "successful_identity_enrichment",
}
COMMENT_MIRRORED_FROM_XSOAR = "***Mirrored from Cortex XSOAR***"
USER_RELATED_FIELDS = ["user", "src_user"]

# =========== Not Missing Events Mechanism Globals ===========
CUSTOM_ID = "custom_id"
OCCURRED = "occurred"
INDEX_TIME = "index_time"
TIME_IS_MISSING = "time_is_missing"


# =========== Enrich User Mechanism ============
class UserMappingObject:
    def __init__(
        self,
        service: client.Service,
        should_map_user: bool,
        table_name: str = "splunk_xsoar_users",
        xsoar_user_column_name: str = "xsoar_user",
        splunk_user_column_name: str = "splunk_user",
    ):
        """
        Args:
            service (client.Service): Splunk service object.
            should_map_user (bool): Whether to map the user or not.
            table_name (str): The name of the table in Splunk.
            xsoar_user_column_name (str): The name of the column in the table that holds the XSOAR user.
            splunk_user_column_name (str): The name of the column in the table that holds the Splunk user.
        """
        self.service = service
        self.should_map = should_map_user
        self.table_name = table_name
        self.xsoar_user_column_name = xsoar_user_column_name
        self.splunk_user_column_name = splunk_user_column_name
        self._kvstore_data: list[dict[str, Any]] = []

    def _get_record(self, col: str, value_to_search: str) -> filter:
        """Gets the records with the value found in the relevant column."""
        if not self._kvstore_data:
            demisto.debug("UserMapping: kvstore data empty, initialize it")
            kvstore: client.KVStoreCollection = self.service.kvstore[self.table_name]
            self._kvstore_data = kvstore.data.query()
            demisto.debug(f"UserMapping: {self._kvstore_data=}")
        return filter(lambda row: row.get(col) == value_to_search, self._kvstore_data)

    def get_xsoar_user_by_splunk(self, splunk_user: str):
        record = list(self._get_record(self.splunk_user_column_name, splunk_user))

        if not record:
            demisto.error(
                f"UserMapping: Could not find xsoar user matching splunk's {splunk_user}. "
                f"Consider adding it to the {self.table_name} lookup."
            )
            return ""

        # assuming username is unique, so only one record is returned.
        xsoar_user = record[0].get(self.xsoar_user_column_name)

        if not xsoar_user:
            demisto.error(
                f"UserMapping: Xsoar user matching splunk's {splunk_user} is empty. Fix the record in {self.table_name} lookup."
            )
            return ""

        return xsoar_user

    def get_splunk_user_by_xsoar(self, xsoar_user: str, map_missing: bool = True):
        record = list(self._get_record(self.xsoar_user_column_name, xsoar_user))

        if not record:
            demisto.error(
                f"UserMapping: Could not find splunk user matching xsoar's {xsoar_user}. "
                f"Consider adding it to the {self.table_name} lookup."
            )
            return "unassigned" if map_missing else None

        # assuming username is unique, so only one record is returned.
        splunk_user = record[0].get(self.splunk_user_column_name)

        if not splunk_user:
            demisto.error(
                f"UserMapping: Splunk user matching Xsoar's {xsoar_user} is empty. Fix the record in {self.table_name} lookup."
            )
            return "unassigned" if map_missing else None

        return splunk_user

    def get_splunk_user_by_xsoar_command(self, args: dict[str, str]) -> CommandResults:
        xsoar_users = argToList(args.get("xsoar_username"))
        map_missing = argToBoolean(args.get("map_missing", True))

        outputs = []
        for user in xsoar_users:
            splunk_user = self.get_splunk_user_by_xsoar(user, map_missing=map_missing) if user else None
            outputs.append(
                {"XsoarUser": user, "SplunkUser": splunk_user or "Could not map splunk user, Check logs for more info."}
            )

        return CommandResults(
            outputs=outputs,
            outputs_prefix="Splunk.UserMapping",
            readable_output=tableToMarkdown("Xsoar-Splunk Username Mapping", outputs, headers=["XsoarUser", "SplunkUser"]),
        )

    def map_owner_to_xsoar_user(self, rows: list[dict]):
        """When `should_map_user` is True, rewrite each row's `owner` field from
        the Splunk username to the mapped XSOAR username.

        This is the action this helper performs — it operates on any row dict
        carrying an `owner` key (Findings, Investigations, etc.).

        Args:
            rows (list[dict]): Row dicts (Findings or Investigations) whose
                `owner` field should be translated in-place.
        """
        if self.should_map:
            demisto.debug("UserMapping: instance configured to map Splunk user to XSOAR users, trying to map.")
            for row in rows:
                if splunk_user := row.get("owner"):
                    xsoar_user = self.get_xsoar_user_by_splunk(splunk_user)
                    row["owner"] = xsoar_user
                    row_id = row.get(EVENT_ID) or row.get("investigation_id") or row.get("investigation_guid")
                    demisto.debug(f"UserMapping: 'owner' was mapped from {splunk_user} to {xsoar_user} for row {row_id}.")


class SplunkGetModifiedRemoteDataResponse(GetModifiedRemoteDataResponse):
    """get-modified-remote-data response parser

    :type modified_findings_data: ``list``
    :param modified_findings_data: The Findings that were modified since the last check.

    :type entries: ``list``
    :param entries: The entries you want to add to the war room.

    :return: No data returned
    :rtype: ``None``
    """

    def __init__(self, modified_findings_data, entries):
        self.modified_findings_data = modified_findings_data
        self.entries = entries
        extensive_log(f"mirror-in: updated findings: {self.modified_findings_data}")
        extensive_log(f"mirror-in: updated entries: {self.entries}")

    def to_entry(self):
        """Convert data to entries.

        :return: List of findings data as entries + entries (from comments and close data),
                 or [{}] if there are only entries and no modified findings.
        :rtype: ``list``
        """
        findings_entries = [
            {
                # Investigations carry `investigation_guid` and are
                # checked first; Findings carry `rule_id`. The fallback lets the same
                # response carry both event types.
                "EntryContext": {"mirrorRemoteId": data.get("investigation_guid") or data.get(RULE_ID)},
                "Contents": data,
                "Type": EntryType.NOTE,
                "ContentsFormat": EntryFormat.JSON,
            }
            for data in self.modified_findings_data
        ]

        if not findings_entries and self.entries:
            return [{}] + self.entries

        return findings_entries + self.entries


# =========== Time & Date Utilities ===========


def get_current_splunk_time(splunk_service: client.Service) -> str:
    """Get the current time from the Splunk server in ISO format with timezone.

    This query uses the gentimes command to generate a single time event, then formats it
    using strftime to get the current server time in ISO_FORMAT_TZ_AWARE format.
    The timezone offset in the result is according to the timezone configured for the user
    who owns the token that was provided for authentication.

    Args:
        splunk_service: Splunk service object

    Returns:
        Current Splunk server time as a string in ISO_FORMAT_TZ_AWARE format
        (e.g., '2025-12-03T11:53:45.138540+02:00')

    Raises:
        ValueError: If the Splunk time cannot be fetched
    """
    get_time_query = f'| gentimes start=-1 | eval clock = strftime(time(), "{ISO_FORMAT_TZ_AWARE}") | sort 1 -_time | table clock'
    search_results = splunk_service.jobs.oneshot(get_time_query, count=1, output_mode=OUTPUT_MODE_JSON)

    reader = results.JSONResultsReader(search_results)
    for item in reader:
        if isinstance(item, dict):
            return item["clock"]
        if handle_message(item):
            continue

    raise ValueError("Error: Could not fetch Splunk time")


def extract_timezone_offset_from_splunk_time(splunk_time_str: str) -> str:
    """Extract timezone offset from Splunk time string.

    Args:
        splunk_time_str: Time string in ISO_FORMAT_TZ_AWARE format
                        (e.g., '2025-12-03T11:53:45.138540+02:00')

    Returns:
        Timezone offset string (e.g., '+02:00', '-05:00')
        Returns '+00:00' if extraction fails
    """
    try:
        # Parse the datetime string to extract timezone
        dt = datetime.strptime(splunk_time_str, ISO_FORMAT_TZ_AWARE)
        # Get the timezone offset
        tz_offset = dt.strftime("%z")
        # Format as +HH:MM
        return f"{tz_offset[:3]}:{tz_offset[3:]}"
    except Exception as e:
        demisto.error(f"Failed to extract timezone from '{splunk_time_str}' using +00:00 as timezone : {e}")
        return "+00:00"  # Default to UTC


def get_splunk_timezone_offset(service: client.Service) -> str:
    """Get Splunk server timezone offset with caching.

    Retrieves the timezone offset from integration context cache.
    If not cached, queries Splunk server and caches the result.

    Args:
        service: Splunk service object

    Returns:
        Timezone offset string (e.g., '+02:00')
    """
    TIMEZONE_CACHE_KEY = "splunk_timezone_offset"

    # Try to get from cache
    integration_context = get_integration_context()
    cached_timezone = integration_context.get(TIMEZONE_CACHE_KEY)

    if cached_timezone:
        demisto.debug(f"Using cached Splunk timezone: {cached_timezone}")
        return cached_timezone

    # Not cached - query Splunk
    demisto.debug("Timezone not cached, querying Splunk server")
    try:
        splunk_time = get_current_splunk_time(service)
        timezone_offset = extract_timezone_offset_from_splunk_time(splunk_time)

        # Cache the timezone
        integration_context[TIMEZONE_CACHE_KEY] = timezone_offset
        set_integration_context(integration_context)

        demisto.debug(f"Cached Splunk timezone: {timezone_offset}")
        return timezone_offset

    except Exception as e:
        demisto.error(f"Failed to get Splunk timezone: using +00:00 as default {e}")
        return "+00:00"  # Default to UTC on error


def enforce_lookback_time(fetch_window_start_time, fetch_window_end_time, look_behind_time):
    """Verifies that the start time of the fetch is at X minutes before
    the end time, X being the number of minutes specified in the look_behind parameter.
    The reason this is needed is to ensure that events that have a significant difference
    between their index time and occurrence time in Splunk are still fetched and are not missed.

    Args:
        fetch_window_start_time (str): The current start time of the fetch.
        fetch_window_end_time (str): The current end time of the fetch.
        look_behind_time (int): The minimal difference (in minutes) that should be enforced between
                                the start time and end time.

    Returns:
        fetch_window_start_time (str): The new start time for the fetch.
    """
    start_time_datetime = datetime.strptime(fetch_window_start_time, ISO_FORMAT_TZ_AWARE)
    end_time_datetime = datetime.strptime(fetch_window_end_time, ISO_FORMAT_TZ_AWARE)
    if end_time_datetime - start_time_datetime < timedelta(minutes=look_behind_time):
        start_time_datetime = end_time_datetime - timedelta(minutes=look_behind_time)
        return datetime.strftime(start_time_datetime, ISO_FORMAT_TZ_AWARE)
    return fetch_window_start_time


def get_fetch_time_window(
    params, service, last_run_fetch_window_start_time: str | None, last_run_fetch_window_end_time: str | None
):
    """Calculate the time window (start and end times) for fetching incidents from Splunk.

    This function determines the boundaries of the fetch query time window, handling first-time fetches,
    enforcing look-behind periods for late-indexed events, and supporting both Splunk server time and local system time.

    Args:
        params: Integration parameters
        service: Splunk service object
        last_run_fetch_window_start_time: The earliest time from the last run (string in ISO_FORMAT_TZ_AWARE)
        last_run_fetch_window_end_time: The latest time from the last run (string in ISO_FORMAT_TZ_AWARE)

    Returns:
        tuple: (fetch_window_start_time, fetch_window_end_time) - both as strings in ISO_FORMAT_TZ_AWARE
    """
    fetch_window_start_time = last_run_fetch_window_start_time
    fetch_window_end_time = last_run_fetch_window_end_time

    # If this is the first fetch (no last run time), calculate it based on first_fetch parameter
    if not fetch_window_start_time:
        demisto.debug(f"[SplunkPy] First fetch - calculated earliest time: {fetch_window_start_time}")
        parse_setting = {"TIMEZONE": "UTC", "RETURN_AS_TIMEZONE_AWARE": True}
        first_fetch = params.get("first_fetch", "10 minutes")
        if parsed_time := dateparser.parse(first_fetch, settings=parse_setting):  # type: ignore[arg-type]
            fetch_window_start_time = parsed_time.strftime(ISO_FORMAT_TZ_AWARE)
        else:
            raise DemistoException(f"Failed to parse first fetch time: {first_fetch}")

    # if fetch_window_end_time is not None it's mean we are in a batch fetch iteration with offset
    # if it's none - take the current time
    if not fetch_window_end_time:
        # Get current time string - either from Splunk server or local system
        # Use Splunk server time to avoid timezone issues
        current_time_from_splunk = get_current_splunk_time(service)
        if datatime_from_splunk_as_utc := dateparser.parse(current_time_from_splunk, settings={"TIMEZONE": "UTC"}):
            fetch_window_end_time = datatime_from_splunk_as_utc.strftime(ISO_FORMAT_TZ_AWARE)
        else:
            raise DemistoException(f"Failed to parse Splunk time: {current_time_from_splunk}")

    occurrence_time_look_behind_minutes = arg_to_number(params.get("occurrence_look_behind") or 15)
    extensive_log(f"[SplunkPy] occurrence look behind is: {occurrence_time_look_behind_minutes}")

    # Enforce look-back time to ensure we don't miss late-indexed events
    fetch_window_start_time = enforce_lookback_time(
        fetch_window_start_time, fetch_window_end_time, occurrence_time_look_behind_minutes
    )

    return fetch_window_start_time, fetch_window_end_time


# =========== Helper Utilities ===========


def quote_group(text: str) -> list[str]:
    """A function that splits groups of key value pairs.
    Taking into consideration key values pairs with nested quotes.
    """

    def clean(t):
        return t.strip().rstrip(",")

    # Return strings that aren't key-valued, as is.
    if len(text.strip()) < 3 or "=" not in text:
        return [text]

    # Remove prefix & suffix wrapping quotes if present around all the text
    # For example a text could be:
    # "a="123"", we want it to be: a="123"
    text = re.sub(r"^\"([\s\S]+\")\"$", r"\1", text)

    # Some of the texts don't end with a comma so we add it to make sure
    # everything acts the same.
    if not text.rstrip().endswith(","):
        text = text.rstrip()
        text += ","

    # Fix elements that aren't key=value (`111, a="123"` => `a="123"`)
    # (^) - start of text
    # ([^=]+), - everything without equal sign and a comma at the end
    #   ('111,' above)
    text = re.sub(r"(^)([^=]+),", ",", text).lstrip(",")

    # Wrap all key values without a quote (`a=123` => `a="123"`)
    # Key part: ([^\"\,]+?=)
    #   asdf=123, here it will match 'asdf'.
    #
    # Value part: ([^\"]+?)
    #   every string without a quote or doesn't start the text.
    #   For example: asdf=123, here it will match '123'.
    #
    # End value part: (,|\")
    #   we need to decide when to end the value, in our case
    #   with a comma. We also check for quotes for this case:
    #   a="b=nested_value_without_a_wrapping_quote", as we want to
    #   wrap 'nested_value_without_a_wrapping_quote' with quotes.
    text = re.sub(r"([^\"\,]+?=)([^\"]+?)(,|\")", r'\1"\2"\3', text)

    # The basic idea here is to check that every key value ends with a `",`
    # Assuming that there are even number of quotes before
    # (some values can have deep nested quotes).
    quote_counter = 0
    rindex = 0
    lindex = 0
    groups = []
    while rindex < len(text):
        # For every quote we increment the quote counter
        # (to preserve context on the opening/closed quotes)
        if text[rindex] == '"':
            quote_counter += 1

        # A quote group ends when `",` is encountered.
        is_end_keypair = rindex > 1 and text[rindex - 1] + text[rindex] == '",'

        # If the quote_counter isn't even we shouldn't close the group,
        # for example: a="b="1",c="3""                * *
        # I'll space for readability:   a = " b = " 1 " , c ...
        #                               0 1 2 3 4 5 6 7 8 9
        # quote_counter is even:            F     T   F   T
        # On index 7 & 8 we find a potential quote closing, but as you can
        # see it isn't a valid group (because of nesting) we need to check
        # the quote counter for an even number => a closing match.
        is_even_number_of_quotes = quote_counter % 2 == 0

        # We check both conditions to find a group
        if is_end_keypair and is_even_number_of_quotes:
            # Clean the match group and append to groups
            groups.append(clean(text[lindex:rindex]))

            # Incrementing the indexes to start searching for the next group.
            lindex = rindex + 1
            rindex += 1
            quote_counter = 0

        # Continue to walk the string until we find a quote again.
        rindex += 1

    # Sometimes there aren't any quotes in the string so we can just append it
    if not groups:
        groups.append(clean(text))

    return groups


def raw_to_dict(raw: str) -> dict[str, str]:
    result: dict[str, str] = {}
    try:
        result = json.loads(raw)
    except ValueError:
        if '"message"' in raw:
            raw = raw.replace('"', "").strip("{").strip("}")
            key_val_arr = raw.split(",")
            for key_val in key_val_arr:
                single_key_val = key_val.split(":", 1)
                if len(single_key_val) <= 1:
                    single_key_val = key_val.split("=", 1)
                if len(single_key_val) > 1:
                    val = single_key_val[1]
                    key = single_key_val[0].strip()

                    result[key] = f"{result[key]},{val}" if key in tuple(result.keys()) else val
        else:
            # search for the pattern: `key="value", `
            # (the double quotes are optional)
            # we append `, ` to the end of the string to catch the last value
            groups = quote_group(raw)
            for g in groups:
                key_value = g.replace('"', "").strip()
                if key_value == "":
                    continue

                if "=" in key_value:
                    key_and_val = key_value.split("=", 1)
                    if key_and_val[0] not in result:
                        result[key_and_val[0]] = key_and_val[1]
                    else:
                        # If there are multiple values for a key, append them.
                        result[key_and_val[0]] = ", ".join([result[key_and_val[0]], key_and_val[1]])

    if REPLACE_FLAG:
        result = replace_keys(result)
    return result


def create_incident_custom_id(incident: dict[str, Any]):
    """This is used to create a custom incident ID, when fetching events that are **NOT** findings.

    Args:
        incident (dict[str, Any]): An incident created from a fetched event.

    Returns:
        str: The custom incident ID.
    """
    incident_raw_data = json.loads(incident["rawJSON"])
    fields_to_add = ["_cd", "index", "_time", "_indextime", "_raw"]
    fields_supplied_by_user = demisto.params().get("unique_id_fields") or ""
    fields_to_add.extend(fields_supplied_by_user.split(","))

    incident_custom_id = "___"
    for field_name in fields_to_add:
        if field_name in incident_raw_data:
            incident_custom_id += f"{field_name}___{incident_raw_data[field_name]}"
        elif field_name in incident:
            incident_custom_id += f"{field_name}___{incident[field_name]}"

    extensive_log(f"[SplunkPy] ID after all fields were added: {incident_custom_id}")

    unique_id = hashlib.md5(incident_custom_id.encode("utf-8")).hexdigest()  # nosec  # guardrails-disable-line
    extensive_log(f"[SplunkPy] Found incident ID is: {unique_id}")
    return unique_id


def extensive_log(message):
    if demisto.params().get("extensive_logs", False):
        demisto.debug(message)


def remove_irrelevant_incident_ids(
    last_run_fetched_ids: dict[str, dict[str, str]], window_start_time: str, window_end_time: str
) -> dict[str, Any]:
    """Remove all the IDs of the fetched incidents that are no longer in the fetch window, to prevent our
    last run object from becoming too large.

    Args:
        last_run_fetched_ids (dict[str, tuple]): The IDs incidents that were fetched in previous fetches.
        window_start_time (str): The window start time.
        window_end_time (str): The window end time.

    Returns:
        dict[str, Any]: The updated list of IDs, without irrelevant IDs.
    """
    new_last_run_fetched_ids: dict[str, dict[str, str]] = {}
    window_start_datetime = datetime.strptime(window_start_time, ISO_FORMAT_TZ_AWARE)
    demisto.debug(f"Beginning to filter irrelevant IDs with respect to window {window_start_time} - {window_end_time}")
    for incident_id, incident_occurred_time in last_run_fetched_ids.items():
        # We divided the handling of the last fetched IDs since we changed the handling of them
        # The first implementation caused IDs to be removed from the cache, even though they were still relevant
        # The second implementation now only removes the cached IDs that are not relevant to the fetch window
        extensive_log(f"[SplunkPy] Checking if {incident_id} is relevant to fetch window")
        # To handle last fetched IDs
        # Last fetched IDs hold the occurred time that they were seen, which is basically the end time of the fetch window
        # they were fetched in, and will be deleted from the last fetched IDs once they pass the fetch window
        incident_window_end_time_str = incident_occurred_time.get("occurred_time", "")
        incident_window_end_datetime = datetime.strptime(incident_window_end_time_str, ISO_FORMAT_TZ_AWARE)
        if incident_window_end_datetime >= window_start_datetime:
            # We keep the incident, since it is still in the fetch window
            extensive_log(f"[SplunkPy] Keeping {incident_id} as part of the last fetched IDs. {incident_window_end_time_str=}")
            new_last_run_fetched_ids[incident_id] = incident_occurred_time
        else:
            extensive_log(f"[SplunkPy] Removing {incident_id} from the last fetched IDs. {incident_window_end_time_str=}")

    return new_last_run_fetched_ids


def build_fetch_kwargs(
    fetch_window_start_time,
    fetch_window_end_time,
    search_offset,
    fetch_window_start_time_fieldname,
    fetch_window_end_time_fieldname,
):
    extensive_log(f"[SplunkPy] fetch_window_start_time_fieldname: {fetch_window_start_time_fieldname}")
    extensive_log(f"[SplunkPy] fetch_window_start_time: {fetch_window_start_time}")
    extensive_log(f"[SplunkPy] fetch_window_end_time_fieldname: {fetch_window_end_time_fieldname}")
    extensive_log(f"[SplunkPy] fetch_window_end_time: {fetch_window_end_time}")

    return {
        fetch_window_start_time_fieldname: fetch_window_start_time,
        fetch_window_end_time_fieldname: fetch_window_end_time,
        "count": FETCH_LIMIT,
        "offset": search_offset,
        "output_mode": OUTPUT_MODE_JSON,
    }


def build_fetch_query(params, query_param_name: str = "fetchQuery"):
    """Build the SPL fetch query for a given param key.

    Defaults to ``fetchQuery`` for the Findings flow; Investigations callers
    pass ``query_param_name="investigations_fetch_query"``. The optional
    ``extractFields`` param is appended only when a value is present — it is
    a Findings-only feature and is normally absent for the investigations path.
    """
    fetch_query = params[query_param_name]

    if extract_fields := params.get("extractFields"):
        for field in extract_fields.split(","):
            field_trimmed = field.strip()
            fetch_query = f"{fetch_query} | eval {field_trimmed}={field_trimmed}"

    return fetch_query


def fetch_findings(
    service: client.Service,
    mapper: UserMappingObject,
    cache_object: "Cache" = None,
    enrich_findings=False,
) -> "FetchResult":
    """Fetch one cycle of Splunk ES Findings.

    Returns a :class:`FetchResult` (incidents + last_run_delta) instead of
    calling ``demisto.incidents()`` / ``demisto.setLastRun()`` directly so the
    dispatcher can participate in mixed-mode cycles without double-emitting.

    In the enrichment-cache path (``enrich_findings=True`` with a
    ``cache_object``), findings accumulate in the cache_object across calls and
    a ``DUMMY`` marker is added to the returned ``last_run_delta`` so the
    dispatcher persists it.

    Each produced incident's ``rawJSON`` carries
    ``splunk_es_event_type = "Finding"`` (set in
    :meth:`Finding.create_incident`); the Classifier owns type routing.
    """
    last_run_data = demisto.getLastRun() or {}
    params = demisto.params()
    if not last_run_data:
        extensive_log("[SplunkPy] SplunkPy first run")

    earliest_time_from_last_run: str | None = last_run_data.get("next_run_earliest_time") if last_run_data else None
    latest_time_from_last_run: str | None = last_run_data.get("next_run_latest_time") if last_run_data else None
    search_offset = last_run_data.get("offset", 0) if last_run_data else 0
    extensive_log(f"[SplunkPy] SplunkPy last run is:\n {last_run_data}")

    fetch_window_start_time, fetch_window_end_time = get_fetch_time_window(
        params, service, earliest_time_from_last_run, latest_time_from_last_run
    )

    finding_time_filter_type: str = params.get("finding_time_source") or "creation time"
    if finding_time_filter_type.startswith("index time"):
        # BETA: For index time based time calculations
        fetch_window_start_time_fieldname = "index_earliest"
        fetch_window_end_time_fieldname = "index_latest"
    else:
        # Finding filter time type defaults to "creation time"
        fetch_window_start_time_fieldname = "earliest_time"
        fetch_window_end_time_fieldname = "latest_time"
    kwargs_oneshot = build_fetch_kwargs(
        fetch_window_start_time,
        fetch_window_end_time,
        search_offset,
        fetch_window_start_time_fieldname,
        fetch_window_end_time_fieldname,
    )
    fetch_query = build_fetch_query(params, query_param_name="fetchQuery")
    last_run_fetched_ids: dict[str, Any] = last_run_data.get("next_run_found_incidents_ids", {})
    if late_indexed_pagination := last_run_data.get("late_indexed_pagination"):
        # This is for handling the case when events get indexed late, and inserted in pages
        # that we have already went through
        window = f"{kwargs_oneshot.get(fetch_window_start_time_fieldname)}-{kwargs_oneshot.get(fetch_window_end_time_fieldname)}"
        demisto.debug(f"[SplunkPy] additional fetch for the window {window} to check for late indexed incidents")
        if last_run_fetched_ids:
            ids_to_exclude = [f'"{fetched_id}"' for fetched_id in last_run_fetched_ids]
            exclude_id_where = f'where not event_id in ({",".join(ids_to_exclude)})'
            fetch_query = f"{fetch_query} | {exclude_id_where}"
            kwargs_oneshot["offset"] = 0

    demisto.debug(f"[SplunkPy] fetch query = {fetch_query}")
    demisto.debug(f"[SplunkPy] oneshot query args = {kwargs_oneshot}")
    oneshotsearch_results = service.jobs.oneshot(fetch_query, **kwargs_oneshot)
    reader = results.JSONResultsReader(oneshotsearch_results)

    error_message = ""
    incidents: list[dict[str, Any]] = []
    findings = []
    incident_ids_to_add = []
    num_of_dropped = 0
    fetched_items = []
    for item in reader:
        if handle_message(item):
            if "Error" in str(item.message) or "error" in str(item.message):
                error_message = f"{error_message}\n{item.message}"
            continue
        fetched_items.append(item)
    if fetched_items:
        # enrich the fetched items with splunk notes
        finding_id_to_item = {item.get(EVENT_ID, ""): item for item in fetched_items if item.get(EVENT_ID)}
        # enrich_with_splunk_notes(service, finding_id_to_item, is_fetch=True)
        # Convert fetch_window_start_time (ISO_FORMAT_TZ_AWARE string) to an epoch float
        # so the v2 KV-store query can filter mc_notes by `update_time`.
        fetch_window_start_epoch = datetime.strptime(fetch_window_start_time, ISO_FORMAT_TZ_AWARE).timestamp()
        enrich_with_splunk_notes_v2(
            service,
            finding_id_to_item,
            last_update_splunk_timestamp=fetch_window_start_epoch,
            is_fetch=True,
        )

    for item in fetched_items:
        extensive_log(f"[SplunkPy] Incident data before parsing to finding: {item}")
        finding_incident = Finding(data=item)
        inc = finding_incident.to_incident(mapper)
        extensive_log(f"[SplunkPy] Incident data after parsing to finding: {inc}")
        incident_id = finding_incident.id or create_incident_custom_id(inc)

        if incident_id not in last_run_fetched_ids:
            incident_ids_to_add.append(incident_id)
            incidents.append(inc)
            findings.append(finding_incident)
            extensive_log(f"[SplunkPy] - Fetched incident {incident_id} to be created.")
        else:
            num_of_dropped += 1
            extensive_log(f"[SplunkPy] - Dropped incident {incident_id} due to duplication.")

    if error_message and not incident_ids_to_add:
        raise DemistoException(f"Failed to fetch incidents, check the provided query in Splunk web search - {error_message}")
    extensive_log(f"[SplunkPy] Size of last_run_fetched_ids before adding new IDs: {len(last_run_fetched_ids)}")
    for incident_id in incident_ids_to_add:
        last_run_fetched_ids[incident_id] = {"occurred_time": fetch_window_end_time}
    extensive_log(f"[SplunkPy] Size of last_run_fetched_ids after adding new IDs: {len(last_run_fetched_ids)}")

    # New way to remove IDs
    last_run_fetched_ids = remove_irrelevant_incident_ids(last_run_fetched_ids, fetch_window_start_time, fetch_window_end_time)
    extensive_log(f"[SplunkPy] Size of last_run_fetched_ids after removing old IDs: {len(last_run_fetched_ids)}")
    extensive_log(f"[SplunkPy] SplunkPy - incidents fetched on last run = {last_run_fetched_ids}")

    demisto.debug(f"SplunkPy - total number of new incidents found is: {len(incidents)}")
    demisto.debug(f"SplunkPy - total number of dropped incidents is: {num_of_dropped}")

    # Cache-mode (enrichment) accumulation.
    enrichment_cache_mode = bool(enrich_findings and cache_object)
    last_run_extras: dict[str, Any] = {}
    if enrichment_cache_mode:
        cache_object.not_yet_submitted_findings += findings  # type: ignore[union-attr]
        if DUMMY not in last_run_data:
            # we add dummy data to the last run to differentiate between the fetch-incidents triggered to the
            # fetch-incidents running as part of "Pull from instance" in Classification & Mapping, as we don't
            # want to add data to the integration context (which will ruin the logic of the cache object)
            last_run_extras[DUMMY] = DUMMY
        emitted_incidents: list[dict[str, Any]] = []
    else:
        emitted_incidents = incidents

    # We didn't get any new incidents or got less than limit,
    # so the next run's earliest time will be the fetch_window_end_time from this iteration
    if (len(incidents) + num_of_dropped) < FETCH_LIMIT:
        demisto.debug(
            f"[SplunkPy] Number of fetched incidents = {len(incidents)}, dropped = {num_of_dropped}. Sum is less"
            f" than {FETCH_LIMIT=}. Starting new fetch"
        )
        new_last_run: dict[str, Any] = {
            "next_run_earliest_time": fetch_window_end_time,
            "next_run_latest_time": None,
            "offset": 0,
            "next_run_found_incidents_ids": last_run_fetched_ids,
        }
    # we get limit findings from splunk
    # we should fetch the entire queue with offset - so set the offset, next_run_earliest_time and next_run_latest_time
    # for the next run
    else:
        demisto.debug(
            f"[SplunkPy] Number of fetched incidents = {len(incidents)}, dropped = {num_of_dropped}. Sum is"
            f" equal/greater than {FETCH_LIMIT=}. Continue pagination"
        )
        new_last_run = {
            "next_run_earliest_time": fetch_window_start_time,
            "next_run_latest_time": fetch_window_end_time,
            "offset": search_offset + FETCH_LIMIT,
            "next_run_found_incidents_ids": last_run_fetched_ids,
        }
    new_last_run["late_indexed_pagination"] = False
    # Need to fetch again this "window" to be sure no "late" indexed events are missed
    if num_of_dropped >= FETCH_LIMIT and "`notable`" in fetch_query:
        demisto.debug('Need to fetch this "window" again to make sure no "late" indexed events are missed')
        new_last_run["late_indexed_pagination"] = True
    # If we are in the process of checking late indexed events, and len(fetch_incidents) == FETCH_LIMIT,
    # that means we need to continue the process of checking late indexed events
    if len(incidents) == FETCH_LIMIT and late_indexed_pagination:
        demisto.debug(
            f"Number of valid incidents equals {FETCH_LIMIT=}, and current fetch checked for late indexed events."
            " Continue checking for late events"
        )
        new_last_run["late_indexed_pagination"] = True

    demisto.debug(
        f'SplunkPy set last run - {new_last_run["next_run_earliest_time"]=}, {new_last_run["next_run_latest_time"]=}, '
        f'{new_last_run["offset"]=}, late_indexed_pagination={new_last_run.get("late_indexed_pagination")}'
    )
    # Merge cache-mode extras (e.g. DUMMY) into the delta so the dispatcher
    # persists them in the same setLastRun call.
    new_last_run.update(last_run_extras)
    demisto.debug(f"fetch_findings: produced={len(emitted_incidents)} incidents")
    return FetchResult(incidents=emitted_incidents, last_run_delta=new_last_run)


def fetch_incidents(service: client.Service, mapper: UserMappingObject):
    """Dispatcher: build the per-event-type :class:`FetchHandler` instances
    selected by the instance parameter ``fetch_event_types`` (default
    ``Finding``), run them in registration order, merge their incidents and
    last-run deltas, then perform a single ``demisto.incidents()`` /
    ``demisto.setLastRun()`` emit at the end.

    Logging contract (single source of truth for fetch-cycle observability):
        - One INFO line per handler with its event_type, produced count and
          its last_run delta keys.
        - One INFO line at the end of the cycle summarising the totals AND
          the relevant slice of the final ``last_run`` for every event type
          actually fetched (Findings keys live at the top level; per-handler
          last-run scopes are surfaced under their ``last_run_key``).
        - Per-event-type details (windows, offsets, dedup IDs) are emitted
          here so both Findings and Investigations are reflected in the same
          log line and operators don't need to cross-reference helpers.
    """
    params = demisto.params()
    selected = argToList(params.get("fetch_event_types")) or ["Finding"]
    handlers = FetchHandlerFactory.build(selected)
    demisto.debug(f"fetch_incidents: selected={selected} handlers={[h.event_type for h in handlers]}")

    last_run = demisto.getLastRun() or {}
    new_last_run: dict[str, Any] = dict(last_run)
    all_incidents: list[dict[str, Any]] = []
    per_handler_counts: dict[str, int] = {}

    for handler in handlers:
        scope_last_run = last_run.get(handler.last_run_key, {}) if handler.last_run_key else last_run
        demisto.debug(
            f"fetch_incidents: invoking {handler.event_type} handler "
            f"scope_keys={list(scope_last_run.keys()) if isinstance(scope_last_run, dict) else []}"
        )
        try:
            result = handler.fetch(service, scope_last_run, mapper, params)
        except Exception as e:
            demisto.error(
                f"fetch_incidents: handler={type(handler).__name__} "
                f"event_type={handler.event_type} failed: {e}\n{traceback.format_exc()}"
            )
            raise
        per_handler_counts[handler.event_type] = len(result.incidents)
        demisto.info(
            f"fetch_incidents: {handler.event_type} produced={len(result.incidents)} "
            f"delta_keys={list(result.last_run_delta.keys())}"
        )
        all_incidents.extend(result.incidents)
        if handler.last_run_key:
            base = scope_last_run if isinstance(scope_last_run, dict) else {}
            new_last_run[handler.last_run_key] = {**base, **result.last_run_delta}
        else:
            new_last_run.update(result.last_run_delta)

    demisto.info(
        f"fetch_incidents: total_incidents={len(all_incidents)} per_type={per_handler_counts}"
        f" | last_run_summary={_summarize_last_run_for_log(new_last_run, [h.event_type for h in handlers])}"
    )
    demisto.incidents(all_incidents)
    demisto.setLastRun(new_last_run)


def _summarize_last_run_for_log(last_run: dict[str, Any], event_types: list[str]) -> dict[str, Any]:
    """Build a compact, log-friendly view of ``last_run`` for the dispatcher's
    end-of-cycle summary.

    For each fetched event type we surface the cursor fields the on-call
    typically wants to see (window start/end, offset, dedup-IDs count, plus
    the late-indexed pagination flag for Findings). Findings keys live at the
    top level of ``last_run``; Investigations live under
    ``last_run["investigations"]``.
    """
    summary: dict[str, Any] = {}
    if "Finding" in event_types:
        finding_ids = last_run.get("next_run_found_incidents_ids") or {}
        summary["Finding"] = {
            "earliest": last_run.get("next_run_earliest_time"),
            "latest": last_run.get("next_run_latest_time"),
            "offset": last_run.get("offset"),
            "dedup_ids": len(finding_ids) if isinstance(finding_ids, dict) else 0,
            "late_indexed_pagination": last_run.get("late_indexed_pagination"),
        }
    if "Investigation" in event_types:
        inv = last_run.get("investigations") or {}
        inv_ids = inv.get("found_incidents_ids") or {}
        summary["Investigation"] = {
            "earliest": inv.get("next_run_earliest_time"),
            "latest": inv.get("next_run_latest_time"),
            "offset": inv.get("offset"),
            "dedup_ids": len(inv_ids) if isinstance(inv_ids, dict) else 0,
        }
    return summary


# =========== Findings handler internal helpers ===========


def _fetch_findings_via_handler(service: client.Service, mapper: UserMappingObject) -> "FetchResult":
    """Internal helper used by :class:`FindingsFetchHandler`.

    Branches on whether enrichments are enabled. When off, delegates to
    :func:`fetch_findings` and propagates the resulting :class:`FetchResult`.
    When on, delegates to :func:`run_enrichment_mechanism`, which returns both
    the produced incidents and the last-run delta (containing the time-window
    cursor and ``next_run_found_incidents_ids`` for dedup) so the dispatcher
    persists them in its single final ``demisto.setLastRun`` call. This
    avoids the dispatcher's stale ``last_run`` snapshot from overwriting the
    dedup IDs that the enrichment cycle just produced.
    """
    if ENABLED_ENRICHMENTS:
        integration_context = get_integration_context() or {}
        if not demisto.getLastRun() and INCIDENTS in integration_context:
            demisto.debug(
                "fetch_incidents: last_run is empty but integration_context exists. "
                "This could be 'Pull from instance' or after 'reset last run'. "
                "If this message appears repeatedly, consider running the 'splunk-reset-enriching-fetch-mechanism' command "
                "to clear stale data and reset the enrichment mechanism."
            )
            fetch_incidents_for_mapping(integration_context)
            # Surface a DUMMY-only delta so the dispatcher's final setLastRun
            # writes only the DUMMY marker and does not overlay un-DUMMYed keys.
            return FetchResult(incidents=[], last_run_delta={DUMMY: DUMMY})
        demisto.debug("running run_enrichment_mechanism")
        enriched_incidents, enrichment_last_run_delta = run_enrichment_mechanism(service, integration_context, mapper)
        return FetchResult(
            incidents=enriched_incidents or [],
            last_run_delta=enrichment_last_run_delta or {},
        )

    demisto.debug("enrichments not enabled running fetch_findings")
    return fetch_findings(
        service=service,
        enrich_findings=False,
        mapper=mapper,
    )


# =========== Regular Fetch Mechanism ===========


# =========== Enriching Fetch Mechanism ===========


class Enrichment:
    """A class to represent an Enrichment. Each finding has 3 possible enrichment types: Drilldown, Asset & Identity

    Attributes:
        type (str): The enrichment type. Possible values are: Drilldown, Asset & Identity.
        id (str): The enrichment's job id in Splunk server.
        data (list): The enrichment's data list (events retrieved from the job's search).
        creation_time (str): The enrichment's creation time in ISO format.
        status (str): The enrichment's status.
        query_name (str): The enrichment's query name.
        query_search (str): The enrichment's query search.
    """

    FAILED = "Enrichment failed"
    EXCEEDED_TIMEOUT = "Enrichment exceed the given timeout"
    IN_PROGRESS = "Enrichment is in progress"
    SUCCESSFUL = "Enrichment successfully handled"
    HANDLED = (EXCEEDED_TIMEOUT, FAILED, SUCCESSFUL)

    def __init__(
        self, enrichment_type, status=None, enrichment_id=None, data=None, creation_time=None, query_name=None, query_search=None
    ):
        self.type = enrichment_type
        self.id = enrichment_id
        self.data = data or []
        self.creation_time = creation_time if creation_time else datetime.now(pytz.UTC).isoformat()
        self.status = status or Enrichment.IN_PROGRESS
        self.query_name = query_name
        self.query_search = query_search

    @classmethod
    def from_job(
        cls, enrichment_type: str, job: client.Job | None, query_name: str | None = None, query_search: str | None = None
    ) -> "Enrichment":
        """Creates an Enrichment object from Splunk Job object

        Args:
            enrichment_type (str): The enrichment type
            job (splunklib.client.Job): The corresponding Splunk Job
            query_name: The enrichment query name
            query_search: The enrichment query search

        Returns:
            The created enrichment (Enrichment)
        """
        if job:
            return cls(
                enrichment_type=enrichment_type, enrichment_id=job["sid"], query_name=query_name, query_search=query_search
            )
        else:
            return cls(enrichment_type=enrichment_type, status=Enrichment.FAILED)

    @classmethod
    def from_json(cls, enrichment_dict: dict[str, Any]) -> "Enrichment":
        """Deserialization method.

        Args:
            enrichment_dict (dict): The enrichment dict in JSON format.

        Returns:
            An instance of the Enrichment class constructed from JSON representation.

        """
        return cls(
            enrichment_type=enrichment_dict.get(TYPE),
            data=enrichment_dict.get(DATA),
            status=enrichment_dict.get(STATUS),
            enrichment_id=enrichment_dict.get(ID),
            creation_time=enrichment_dict.get(CREATION_TIME),
            query_name=enrichment_dict.get(QUERY_NAME),
            query_search=enrichment_dict.get(QUERY_SEARCH),
        )


# =========== Investigations helpers ===========

PLACEHOLDER = "FETCH_FILTER_PLACEHOLDER"


def to_mc_iso8601_utc(ts: "str | datetime") -> str:
    """Normalize a timestamp into the exact shape accepted by the
    `investigations` endpoint for `create_time_min` / `create_time_max`:
    `YYYY-MM-DDTHH:MM:SS.ffffffZ` (always UTC, literal `Z` suffix; canonical 6-digit microseconds).

    Behavior:
      - Numeric offsets (e.g. `+00:00`, `-05:00`) are converted to UTC and emitted with `Z`.
      - Microseconds are always emitted as 6 digits (zero-padded), so the output shape
        is stable for downstream consumers (last-run cursors, dedup keys, audit logs).
      - Naive inputs (no timezone) are rejected — the endpoint rejects them too.

    Args:
        ts: Either an ISO 8601 string or a timezone-aware ``datetime``.

    Returns:
        A canonical Mission Control timestamp string of the form
        ``YYYY-MM-DDTHH:MM:SS[.ffffff]Z``.

    Raises:
        DemistoException: If the input is a naive datetime, a naive ISO 8601 string,
            or otherwise unparseable.
    """
    # String path
    if isinstance(ts, str):
        try:
            # `datetime.fromisoformat` (Python 3.10) accepts `+00:00` style offsets
            # and (3.11+) the literal `Z` suffix as well. Normalize a trailing `Z`
            # to `+00:00` for broad compatibility.
            normalized = ts[:-1] + "+00:00" if ts.endswith("Z") else ts
            dt = datetime.fromisoformat(normalized)
        except ValueError as e:
            raise DemistoException(f"to_mc_iso8601_utc: failed to parse ISO 8601 timestamp '{ts}': {e}") from e
    elif isinstance(ts, datetime):
        dt = ts
    else:
        raise DemistoException(f"to_mc_iso8601_utc: unsupported input type {type(ts).__name__}; " "expected str or datetime.")

    if dt.tzinfo is None:
        raise DemistoException(f"to_mc_iso8601_utc: naive timestamp '{ts}' rejected.")

    # Convert to UTC.
    dt_utc = dt.astimezone(UTC)
    # Always emit canonical 6-digit microsecond precision so the output shape is
    # stable for downstream consumers (last-run cursors, dedup keys, audit logs)
    # regardless of whether the source timestamp carried sub-second precision.
    return dt_utc.strftime("%Y-%m-%dT%H:%M:%S.%f") + "Z"


def prepare_investigations_query(
    user_query: str,
    create_time_min: str,
    create_time_max: str,
    limit: int,
    offset: int,
) -> str:
    """Inject ``create_time_min`` / ``create_time_max`` / ``limit`` / ``offset`` into
    ``user_query`` by substituting the mandatory ``FETCH_FILTER_PLACEHOLDER`` token.

    Behavior:
      - The query MUST contain ``FETCH_FILTER_PLACEHOLDER``; the four managed
        params replace it at runtime.
      - ``limit`` is clamped to ``INVESTIGATIONS_MAX_LIMIT`` (100).

    Raises:
        DemistoException: if ``FETCH_FILTER_PLACEHOLDER`` is missing from
            ``user_query``.
    """
    limit = min(int(limit), INVESTIGATIONS_MAX_LIMIT)
    managed = {
        "create_time_min": create_time_min,
        "create_time_max": create_time_max,
        "limit": str(limit),
        "offset": str(offset),
    }

    if PLACEHOLDER not in user_query:
        raise DemistoException(
            f'The "Investigations fetch query" must contain the literal token `{PLACEHOLDER}`. '
            "Do not modify or remove this token — the integration replaces it at runtime with the "
            "time-range and pagination filters needed to manage the fetch cycle. To customize the "
            "query, append extra URL parameters (e.g., &status=New) without changing the placeholder. "
            'See the "Fetching investigation events" section in the integration documentation for details.'
        )

    params_str = urlencode(managed)
    demisto.debug(f"prepare_investigations_query: placeholder substitution; params={dict(managed)}")
    return user_query.replace(PLACEHOLDER, params_str)


# =========== Investigations model ===========


def collapse_dotted_keys_to_nested(data: dict[str, Any], prefix: str) -> dict[str, Any]:
    """Collapse flat dotted keys into a nested dict under a single key.

    Given a dict containing keys of the form ``"<prefix><inner_key>"`` (where
    ``prefix`` ends with a ``.``), this function pops them from ``data`` and
    re-inserts them as a nested dict under the key ``prefix`` (without the
    trailing dot). If a nested dict already exists at that key, the dotted
    keys are merged into it (dotted-key values take precedence on conflict).

    The input dict is mutated in place and also returned for convenience.

    Args:
        data: The dict to transform (mutated in place).
        prefix: The dotted-key prefix to collapse, e.g. ``"consolidated_findings."``.
            Must end with a ``.``.

    Returns:
        The same ``data`` dict, now with dotted keys collapsed.
    """
    if not prefix.endswith("."):
        prefix = f"{prefix}."
    parent_key = prefix[:-1]

    dotted_keys = [k for k in data if isinstance(k, str) and k.startswith(prefix)]
    if not dotted_keys:
        return data

    nested: dict[str, Any] = {}
    for full_key in dotted_keys:
        inner_key = full_key[len(prefix) :]
        nested[inner_key] = data.pop(full_key)

    # Preserve any pre-existing nested dict by merging dotted keys into it.
    existing = data.get(parent_key)
    if isinstance(existing, dict):
        existing.update(nested)
        nested = existing
    data[parent_key] = nested
    return data


def collapse_all_dotted_keys_to_nested(data: dict[str, Any]) -> dict[str, Any]:
    """Auto-discover every dotted-key prefix in ``data`` and collapse each one.

    Generic counterpart to :func:`collapse_dotted_keys_to_nested`: instead of
    requiring the caller to know the prefix, this scans ``data`` for keys
    containing a literal ``.``, groups them by their first-segment prefix, and
    delegates each group to :func:`collapse_dotted_keys_to_nested`.

    Only the FIRST dot is treated as the prefix separator — keys with multiple
    dots (e.g. ``a.b.c``) are collapsed under the parent ``a`` and their
    remainder (``b.c``) is preserved as the inner key, exactly mirroring the
    single-prefix variant's behaviour.

    The input dict is mutated in place and also returned for convenience.
    """
    # Snapshot first — we mutate `data` while iterating.
    prefixes = {k.split(".", 1)[0] + "." for k in list(data.keys()) if isinstance(k, str) and "." in k}
    for prefix in prefixes:
        collapse_dotted_keys_to_nested(data, prefix)
    return data


def parse_investigation(row: dict[str, Any]) -> dict[str, Any]:
    """Flatten an investigations row.

    Pure transform — does not call any Splunk service.

    The function:
      - Returns a shallow copy so the caller's row is not mutated.
      - Lifts the nested ``findings.incident_ids`` array up to top-level
        ``incident_ids`` for easier classifier/mapper consumption.
      - Serializes ``consolidated_findings`` (a nested object containing both
        scalar fields and parallel array columns) into a JSON string. The
        classifier then maps it as-is into the longText incident field
        ``splunkconsolidatedfindings``, where the ``SplunkConvertConsolidatedFindingsToMD``
        dynamic-section script renders it as Markdown in the layout.
      - Stamps ``splunk_es_event_type = "Investigation"`` for the classifier.

    Args:
        row: Raw row as returned by the investigations endpoint.

    Returns:
        A new ``dict`` with the parsed/flattened fields.
    """
    parsed: dict[str, Any] = dict(row)

    # Collapse ALL dotted-key prefixes into nested dicts (e.g.
    # `consolidated_findings.x`, `findings.y`, `<anything>.z`). This keeps
    # parse_investigation forward-compatible with new nested objects added by
    # the v2 endpoint without needing per-prefix bookkeeping.
    collapse_all_dotted_keys_to_nested(parsed)

    parsed["incident_ids"] = parsed.get("findings", {}).pop("incident_ids", [])
    consolidated = parsed.get("consolidated_findings")
    if consolidated is not None and not isinstance(consolidated, str):
        try:
            parsed["consolidated_findings"] = json.dumps(consolidated)
        except (TypeError, ValueError) as exc:
            demisto.debug(f"parse_investigation: consolidated_findings not JSON-serializable; leaving as-is. err={exc}")

    # Type tag consumed by the classifier.
    parsed[SPLUNK_ES_EVENT_TYPE_FIELD] = Investigation.event_type

    if "status_label" not in parsed:
        raw_status = str(parsed.get("status_name", "") or "")
        parsed["status_label"] = raw_status.capitalize() if raw_status else ""
    parsed.setdefault("status_end", "false")

    return parsed


class Investigation:
    """A lightweight model for a Splunk investigation row.

    Unlike :class:`Finding`, this model has no enrichment lifecycle and does not
    interact with any Splunk service. It is responsible only for turning a
    parsed row into an XSOAR incident dict.
    """

    event_type: str = "Investigation"

    def __init__(self, parsed: dict[str, Any]) -> None:
        # `parsed` is expected to be the output of `parse_investigation` —
        # i.e. it already carries the splunk_es_event_type tag. We accept the
        # raw row too for ergonomic test/use; re-running parse is idempotent.
        if parsed.get(SPLUNK_ES_EVENT_TYPE_FIELD) != self.event_type:
            parsed = parse_investigation(parsed)
        self.data: dict[str, Any] = parsed

    def get_id(self) -> str:
        """Return a stable id for dedup/display.

        Order: ``investigation_id`` → ``investigation_guid`` →
        :func:`create_incident_custom_id` fallback.
        """
        if self.data.get("investigation_id"):
            return str(self.data["investigation_id"])
        if self.data.get("investigation_guid"):
            return str(self.data["investigation_guid"])
        # Final fallback: use the dedup helper that operates on rawJSON.
        return create_incident_custom_id({"rawJSON": json.dumps(self.data)})

    def get_occurred(self) -> str:
        """Return an RFC 3339 string built from ``create_time`` (epoch float).

        Falls back to ``mc_create_time`` and finally to ``now()`` if neither
        is present.
        """
        epoch = self.data.get("create_time") or self.data.get("mc_create_time")
        if epoch is None:
            return datetime.now(UTC).strftime(ISO_FORMAT_TZ_AWARE)
        try:
            return datetime.fromtimestamp(float(epoch), tz=UTC).isoformat()
        except (TypeError, ValueError):
            return datetime.now(UTC).strftime(ISO_FORMAT_TZ_AWARE)

    def to_incident(self, mapper: "UserMappingObject") -> dict[str, Any]:
        """Build an XSOAR incident dict from this investigation.

        Notes:
            - Does NOT set ``incident["type"]``. The classifier consumes
              ``rawJSON.splunk_es_event_type`` to route to the right type.
            - ``dbotMirrorId`` is set to ``investigation_guid`` so the existing
              mirror-in/mirror-out commands (which talk to the v2 endpoint)
              keep working without any changes.
        """
        params_local = demisto.params()
        data = self.data

        name = data.get("name") or data.get("investigation_id") or "Splunk Investigation"
        incident: dict[str, Any] = {
            "name": name,
            "occurred": self.get_occurred(),
        }

        if data.get("description"):
            incident["details"] = data["description"]

        if data.get("urgency"):
            incident["severity"] = severity_to_level(data["urgency"])

        owner_value = data.get("owner")
        if owner_value and mapper.should_map and (mapped := mapper.get_xsoar_user_by_splunk(owner_value)):
            data["owner"] = mapped
            incident["owner"] = mapped

        # Mirror plumbing — same shape as Finding.create_incident:
        #   - dbotMirror* fields on the incident itself
        #   - a `mirror_*` mirror block inside rawJSON for downstream consumers
        mirror_direction = MIRROR_DIRECTION.get(params_local.get("mirror_direction") or "None")
        mirror_instance = demisto.integrationInstance()
        if data.get("investigation_guid"):
            incident["dbotMirrorId"] = data["investigation_guid"]
        incident["dbotMirrorInstance"] = mirror_instance
        incident["dbotMirrorDirection"] = mirror_direction

        data.update(
            {
                "mirror_instance": mirror_instance,
                "mirror_direction": mirror_direction,
                "mirror_tags": [NOTE_TAG_FROM_SPLUNK, NOTE_TAG_TO_SPLUNK],
            }
        )

        # Ensure the type tag survives even if a caller built `Investigation`
        # from a raw (un-parsed) row directly.
        data[SPLUNK_ES_EVENT_TYPE_FIELD] = self.event_type
        incident["rawJSON"] = json.dumps(data)
        return incident


# =========== Fetch handler factory ===========


@dataclass
class FetchResult:
    """Per-handler fetch result.

    ``incidents`` are merged by the dispatcher and emitted via
    ``demisto.incidents()``. ``last_run_delta`` is merged into the appropriate
    last-run namespace (top-level for Findings; ``last_run["investigations"]``
    for Investigations).
    """

    incidents: list[dict[str, Any]] = field(default_factory=list)
    last_run_delta: dict[str, Any] = field(default_factory=dict)


class FetchHandler(ABC):
    """Abstract base for per-event-type fetch handlers."""

    event_type: str = ""  # "Finding" | "Investigation"
    last_run_key: "str | None" = None  # None = top-level last-run namespace

    @abstractmethod
    def fetch(self, service, last_run, mapper, params) -> FetchResult:  # noqa: D401
        """Run one fetch cycle for this event type and return a FetchResult."""
        raise NotImplementedError


class FetchHandlerFactory:
    """Builds an ordered list of :class:`FetchHandler` instances per the
    user-selected ``fetch_event_types`` parameter.
    """

    _registry: dict[str, type[FetchHandler]] = {}

    @classmethod
    def register(cls, event_type: str, handler_cls: type[FetchHandler]) -> None:
        """Register a handler class under an ``event_type`` key.

        Registration order is preserved by the underlying ``dict`` (Python 3.7+),
        which the :meth:`build` method relies on for deterministic output.
        """
        cls._registry[event_type] = handler_cls

    @classmethod
    def build(cls, selected_types: "list[str] | None") -> list[FetchHandler]:
        """Instantiate handlers for ``selected_types``.

        - Empty / ``None`` selection → defaults to ``["Finding"]``.
        - Unknown event types are silently ignored.
        - Output order matches the order in ``selected_types``.
        """
        if not selected_types:
            selected_types = ["Finding"]
        ignored = [t for t in selected_types if t not in cls._registry]
        built = [cls._registry[t]() for t in selected_types if t in cls._registry]
        demisto.debug(
            "FetchHandlerFactory.build: requested="
            f"{selected_types}; built={[h.__class__.__name__ for h in built]}; "
            f"ignored={ignored}"
        )
        return built


# =========== Concrete fetch handlers ===========


class FindingsFetchHandler(FetchHandler):
    """Handler for the Splunk ES Findings fetch path.

    Wraps :func:`_fetch_findings_via_handler` so the dispatcher can treat it
    uniformly with new event types.
    """

    event_type: str = "Finding"
    last_run_key: "str | None" = None  # top-level last-run keys

    def fetch(self, service, last_run, mapper, params) -> "FetchResult":  # noqa: ARG002
        return _fetch_findings_via_handler(service, mapper)


class InvestigationsFetchHandler(FetchHandler):
    """Handler for the Splunk Investigations fetch path.

    Does NOT call :func:`run_enrichment_mechanism` — enrichment is
    Findings-only.
    """

    event_type: str = "Investigation"
    last_run_key: "str | None" = "investigations"

    # Default SPL fallback for instances upgraded without re-saving config.
    _DEFAULT_SPL: str = (
        '| rest "/servicesNS/nobody/missioncontrol/public/v2/investigations?search_format=true&FETCH_FILTER_PLACEHOLDER"'
    )

    def fetch(self, service, last_run, mapper, params) -> "FetchResult":
        try:
            return self._do_fetch(service, last_run, mapper, params)
        except Exception as e:
            demisto.error(f"handler=InvestigationsFetchHandler event_type=Investigation failed: {e}\n{traceback.format_exc()}")
            raise

    def _do_fetch(
        self,
        service: client.Service,
        last_run: dict[str, Any],
        mapper: "UserMappingObject",
        params: dict[str, Any],
    ) -> "FetchResult":
        # Time window — last_run mirrors the Findings shape under the investigations namespace.
        earliest_raw = last_run.get("next_run_earliest_time") or last_run.get("time")
        latest_raw = last_run.get("next_run_latest_time")
        # We deliberately reuse the same window helper Findings uses so the
        # `investigations_first_fetch` parameter participates in the standard
        # first-fetch path.
        params_for_window = dict(params, first_fetch=params.get("investigations_first_fetch"))
        earliest, latest = get_fetch_time_window(params_for_window, service, earliest_raw, latest_raw)

        ts_min = to_mc_iso8601_utc(earliest)
        ts_max = to_mc_iso8601_utc(latest)

        # 2. Limit + offset.
        max_fetch_raw = arg_to_number(params.get("investigations_max_fetch")) or 50
        limit = min(max_fetch_raw, INVESTIGATIONS_MAX_LIMIT)
        offset = arg_to_number(last_run.get("offset")) or 0

        demisto.debug(f"InvestigationsFetchHandler: earliest={ts_min} latest={ts_max} " f"limit={limit} offset={offset}")

        # 3. Build SPL.
        user_query = params.get("investigations_fetch_query") or self._DEFAULT_SPL
        final_spl = prepare_investigations_query(user_query, ts_min, ts_max, limit, offset)
        demisto.debug(f"investigations: prepared SPL={final_spl}")

        # The investigations endpoint filters by URL params, not earliest/latest macros, so we don't reuse build_fetch_kwargs.
        oneshot_kwargs: dict[str, Any] = {"output_mode": OUTPUT_MODE_JSON, "count": 0}
        oneshot_results = service.jobs.oneshot(final_spl, **oneshot_kwargs)
        reader = results.JSONResultsReader(oneshot_results)

        # 5. Iterate and dedup.
        rows: list[dict[str, Any]] = []
        error_message = ""
        for item in reader:
            if handle_message(item):
                if "Error" in str(item.message) or "error" in str(item.message):
                    error_message = f"{error_message}\n{item.message}"
                continue
            rows.append(item)

        # Enrich the fetched investigations with splunk notes BEFORE building incidents,
        # so `splunk_notes` is included in `rawJSON` (same pattern as findings).
        # Investigation notes reference their parent via `incident_id`, while
        # `enrich_with_splunk_notes` also accepts that key transparently.
        if rows:
            investigation_guid_to_row = {str(row["investigation_guid"]): row for row in rows if row.get("investigation_guid")}
            if investigation_guid_to_row:
                # enrich_with_splunk_notes(service, investigation_guid_to_row, is_fetch=True)
                # Convert `earliest` (ISO_FORMAT_TZ_AWARE string) to an epoch float so
                # the v2 KV-store query can filter mc_notes by `update_time`.
                fetch_window_start_epoch = datetime.strptime(earliest, ISO_FORMAT_TZ_AWARE).timestamp()
                enrich_with_splunk_notes_v2(
                    service,
                    investigation_guid_to_row,
                    last_update_splunk_timestamp=fetch_window_start_epoch,
                    is_fetch=True,
                )

        last_run_fetched_ids: dict[str, dict[str, str]] = dict(last_run.get("found_incidents_ids", {}) or {})

        incidents: list[dict[str, Any]] = []
        ids_to_add: list[str] = []
        dropped_duplicate_ids: list[str] = []
        for row in rows:
            parsed = parse_investigation(row)
            investigation = Investigation(parsed)
            incident = investigation.to_incident(mapper)
            # Dedup: prefer stable investigation_id; fall back to GUID; final
            # fallback to a content-hash custom id.
            stable_id = (
                str(parsed.get("investigation_id"))
                if parsed.get("investigation_id")
                else (
                    str(parsed.get("investigation_guid"))
                    if parsed.get("investigation_guid")
                    else create_incident_custom_id(incident)
                )
            )
            if stable_id in last_run_fetched_ids:
                dropped_duplicate_ids.append(stable_id)
                continue
            ids_to_add.append(stable_id)
            incidents.append(incident)

        if dropped_duplicate_ids:
            demisto.debug(f"investigations: dropped duplicate ids={dropped_duplicate_ids}")

        if error_message and not ids_to_add:
            raise DemistoException(
                f"Failed to fetch investigations, check the provided SPL in Splunk web search - {error_message}"
            )

        for incident_id in ids_to_add:
            last_run_fetched_ids[incident_id] = {"occurred_time": latest}

        last_run_fetched_ids = remove_irrelevant_incident_ids(last_run_fetched_ids, earliest, latest)

        # 8. Cursor: advance by `limit` when full page, reset otherwise.
        next_offset = offset + limit if len(rows) >= limit else 0
        demisto.debug(f"investigations: limit={limit} offset={offset} received={len(rows)} " f"-> next_offset={next_offset}")

        # When we advance the cursor we keep the same window (caller should
        # re-issue with the same min/max next cycle). When we reset, we
        # advance the window to `latest`.
        if next_offset > 0:
            new_time = ts_min
            new_latest = ts_max
        else:
            new_time = ts_max
            new_latest = None

        delta: dict[str, Any] = {
            "time": new_time,
            "next_run_earliest_time": new_time,
            "next_run_latest_time": new_latest,
            "offset": next_offset,
            "found_incidents_ids": last_run_fetched_ids,
        }
        demisto.info(
            f"InvestigationsFetchHandler: produced={len(incidents)} incidents; " f"last_run_delta_keys={list(delta.keys())}"
        )
        return FetchResult(incidents=incidents, last_run_delta=delta)


# Order matters: Findings first preserves existing incident ordering in mixed-mode cycles.
FetchHandlerFactory.register("Finding", FindingsFetchHandler)
FetchHandlerFactory.register("Investigation", InvestigationsFetchHandler)


class Finding:
    """A class to represent a finding (Splunk ES 8.2+).

    Attributes:
        data (dict): The finding data.
        id (str): The finding's id (event_id).
        enrichments (list): The list of all enrichments that needs to handle.
        incident_created (bool): Whether an incident created or not.
        occurred (str): The occurred time of the finding.
        custom_id (str): The custom ID of the finding (used in the fetch function).
        time_is_missing (bool): Whether the `_time` field has an empty value or not.
        index_time (str): The time the finding have been indexed.
    """

    def __init__(
        self,
        data: dict[str, Any],
        enrichments: list[Enrichment] | None = None,
        finding_id: str | None = None,
        occurred: str | None = None,
        custom_id: str | None = None,
        index_time: str | None = None,
        time_is_missing: bool | None = None,
        incident_created: bool | None = None,
    ) -> None:
        self.data = data
        self.id = finding_id or self.get_id()
        self.enrichments = enrichments or []
        self.incident_created = incident_created or False
        self.time_is_missing = time_is_missing or False
        self.index_time = index_time or self.data.get("_indextime")
        self.occurred = occurred or self.get_occurred()
        self.custom_id = custom_id or self.create_custom_id()

    def get_id(self) -> str:
        if EVENT_ID in self.data:
            return self.data[EVENT_ID]
        if ENABLED_ENRICHMENTS:
            raise Exception(
                "When using the enrichment mechanism, an event_id field is needed, and thus, "
                "one must use a fetch query of the following format: search `notable` .......\n"
                "Please re-edit the fetchQuery parameter in the integration configuration, reset "
                "the fetch mechanism using the splunk-reset-enriching-fetch-mechanism command and "
                "run the fetch again."
            )
        else:
            return ""

    @staticmethod
    def create_incident(finding_data: dict[str, Any], occurred: str, mapper: UserMappingObject) -> dict[str, Any]:
        rule_title, rule_name = "", ""
        params = demisto.params()
        if demisto.get(finding_data, "rule_title"):
            rule_title = finding_data["rule_title"]
        if demisto.get(finding_data, "rule_name"):
            rule_name = finding_data["rule_name"]
        incident: dict[str, Any] = {"name": f"{rule_title} : {rule_name}"}
        if demisto.get(finding_data, "urgency"):
            incident["severity"] = severity_to_level(finding_data["urgency"])
        if demisto.get(finding_data, "rule_description"):
            incident["details"] = finding_data["rule_description"]
        if finding_data.get("owner") and mapper.should_map and (owner := mapper.get_xsoar_user_by_splunk(finding_data["owner"])):
            finding_data["owner"] = owner
            incident["owner"] = owner

        incident["occurred"] = occurred
        finding_data = parse_finding(finding_data)
        finding_data.update(
            {
                "mirror_instance": demisto.integrationInstance(),
                "mirror_direction": MIRROR_DIRECTION.get(params.get("mirror_direction")),
                "mirror_tags": [NOTE_TAG_FROM_SPLUNK, NOTE_TAG_TO_SPLUNK],
            }
        )
        splunk_note_entries = []
        labels = []
        if params.get("parseFindingEventsRaw"):
            for key, value in raw_to_dict(finding_data["_raw"]).items():
                if not isinstance(value, str):
                    value = str(value)
                labels.append({"type": key, "value": value})
        if demisto.get(finding_data, "security_domain"):
            labels.append({"type": "security_domain", "value": finding_data["security_domain"]})
        splunk_note_entries = demisto.get(finding_data, "splunk_notes", [])
        incident["splunk_notes"] = splunk_note_entries
        labels.append({"type": "splunk_notes", "value": str(splunk_note_entries)})
        incident["labels"] = labels
        if finding_data.get(EVENT_ID):
            incident["dbotMirrorId"] = finding_data.get(EVENT_ID)
        # Tag for Classifier routing via rawJSON.splunk_es_event_type. Single tagging point.
        finding_data[SPLUNK_ES_EVENT_TYPE_FIELD] = "Finding"
        incident["rawJSON"] = json.dumps(finding_data)

        return incident

    def to_incident(self, mapper: UserMappingObject) -> dict[str, Any]:
        """Gathers all data from all finding's enrichments and return an incident"""
        self.incident_created = True

        for e in self.enrichments:
            if e.type == DRILLDOWN_ENRICHMENT:
                # A finding can have more than one drilldown search enrichment, in that case we keep the searches results in
                # a list of dictionaries - each dict contains the query detail and the search results of a drilldown search

                drilldown_enrichment_details = {
                    "query_name": e.query_name,
                    "query_search": e.query_search,
                    "query_results": e.data,
                    "enrichment_status": e.status,
                }

                if not self.data.get(e.type):  # first drilldown enrichment result to add - initiate the list
                    self.data[e.type] = [drilldown_enrichment_details]

                else:  # there are previous drilldown enrichments in the finding's data
                    self.data[e.type].append(drilldown_enrichment_details)

                if not self.data.get("successful_drilldown_enrichment"):
                    # Drilldown enrichment is successful if at least one drilldown search was successful
                    self.data["successful_drilldown_enrichment"] = e.status == Enrichment.SUCCESSFUL

            else:  # asset enrichment or identity enrichment
                self.data[e.type] = e.data
                self.data[ENRICHMENT_TYPE_TO_ENRICHMENT_STATUS[e.type]] = e.status == Enrichment.SUCCESSFUL

        return self.create_incident(
            self.data,
            self.occurred,
            mapper=mapper,
        )

    def submitted(self) -> bool:
        """Returns an indicator on whether any of the finding's enrichments was submitted or not"""
        finding_enrichment_types = {e.type for e in self.enrichments}
        return any(enrichment.status == Enrichment.IN_PROGRESS for enrichment in self.enrichments) and len(
            finding_enrichment_types
        ) == len(ENABLED_ENRICHMENTS)

        # Conditions:
        # 1. An IN_PROGRESS enrichment indicates the finding was submitted to Splunk.
        # 2. Each enabled enrichment type is represented at least once in self.enrichments
        #    (submit_finding() always appends an Enrichment regardless of success/failure).

    def failed_to_submit(self) -> bool:
        """Returns an indicator on whether all finding's enrichments were failed to submit or not"""
        finding_enrichment_types = {e.type for e in self.enrichments}
        return all(enrichment.status == Enrichment.FAILED for enrichment in self.enrichments) and len(
            finding_enrichment_types
        ) == len(ENABLED_ENRICHMENTS)

    def handled(self) -> bool:
        """Returns an indicator on whether all finding's enrichments were handled or not"""
        return all(enrichment.status in Enrichment.HANDLED for enrichment in self.enrichments) or any(
            enrichment.status == Enrichment.EXCEEDED_TIMEOUT for enrichment in self.enrichments
        )

    def get_submitted_enrichments(self) -> tuple[bool, bool, bool]:
        """Returns indicators on whether each enrichment was submitted/failed or not initiated"""
        submitted_drilldown, submitted_asset, submitted_identity = False, False, False

        for enrichment in self.enrichments:
            if enrichment.type == DRILLDOWN_ENRICHMENT:
                submitted_drilldown = True
            elif enrichment.type == ASSET_ENRICHMENT:
                submitted_asset = True
            elif enrichment.type == IDENTITY_ENRICHMENT:
                submitted_identity = True

        return submitted_drilldown, submitted_asset, submitted_identity

    def get_occurred(self) -> str:
        """Returns the occurred time, if not exists in data, returns the current fetch time"""
        if "_time" in self.data:
            finding_occurred = self.data["_time"]
        else:
            # Use-cases where fetching non-findings from Splunk

            finding_occurred = datetime.now(pytz.UTC).strftime(ISO_FORMAT_TZ_AWARE)
            self.time_is_missing = True
            demisto.debug(f"\n\n occurred time in else: {finding_occurred} \n\n")

        return finding_occurred

    def create_custom_id(self) -> str:
        """Generates a custom ID for a given finding"""
        if self.id:
            return self.id

        finding_raw_data = self.data.get("_raw", "")
        raw_hash = hashlib.md5(finding_raw_data.encode("utf-8")).hexdigest()  # nosec  # guardrails-disable-line

        if self.time_is_missing and self.index_time:
            finding_custom_id = f"{self.index_time}_{raw_hash}"  # index_time stays in epoch to differentiate
            demisto.debug("Creating finding custom id using the index time")
        else:
            finding_custom_id = f"{self.occurred}_{raw_hash}"

        return finding_custom_id

    def is_enrichment_process_exceeding_timeout(self, enrichment_timeout: int) -> bool:
        """Checks whether an enrichment process has exceeded timeout or not

        Args:
            enrichment_timeout (int): The timeout for the enrichment process

        Returns (bool): True if the enrichment process exceeded the given timeout, False otherwise
        """
        now = datetime.now(pytz.UTC)
        exceeding_timeout = False

        for enrichment in self.enrichments:
            if enrichment.status == Enrichment.IN_PROGRESS:
                creation_time_datetime = datetime.strptime(enrichment.creation_time, ISO_FORMAT_TZ_AWARE)
                if now - creation_time_datetime > timedelta(minutes=enrichment_timeout):
                    exceeding_timeout = True
                    enrichment.status = Enrichment.EXCEEDED_TIMEOUT

        return exceeding_timeout

    @classmethod
    def from_json(cls, finding_dict: dict[str, Any]) -> "Finding":
        """Deserialization method.

        Args:
            finding_dict: The finding dict in JSON format.

        Returns:
            An instance of the Finding class constructed from JSON representation.
        """
        return cls(
            data=finding_dict.get(DATA) or {},
            enrichments=list(map(Enrichment.from_json, finding_dict.get(ENRICHMENTS) or [])),
            finding_id=finding_dict.get(ID),
            custom_id=finding_dict.get(CUSTOM_ID),
            occurred=finding_dict.get(OCCURRED),
            time_is_missing=finding_dict.get(TIME_IS_MISSING),
            index_time=finding_dict.get(INDEX_TIME),
            incident_created=finding_dict.get(INCIDENT_CREATED),
        )


class Cache:
    """A class to represent the cache for the enriching fetch mechanism.

    Attributes:
        not_yet_submitted_findings (list): The list of all findings that were fetched but not yet submitted.
        submitted_findings (list): The list of all submitted findings that needs to be handled.
    """

    def __init__(
        self, not_yet_submitted_findings: list[Finding] | None = None, submitted_findings: list[Finding] | None = None
    ) -> None:
        self.not_yet_submitted_findings = not_yet_submitted_findings or []
        self.submitted_findings = submitted_findings or []

    def done_submitting(self) -> bool:
        return not self.not_yet_submitted_findings

    def done_handling(self) -> bool:
        return not self.submitted_findings

    def organize(self) -> list[Finding]:
        """This function is designated to handle unexpected behaviors in the enrichment mechanism.
         E.g. Connection error, instance disabling, etc...
         It re-organizes the cache object to the correct state of the mechanism when the exception was caught.
         If there are findings that were handled but the mechanism didn't create an incident for them, it returns them.
         This function is called in each "end" of execution of the enrichment mechanism.

        Returns:
            handled_not_created_incident (list): The list of all findings that have been handled but not created an
             incident.
        """
        not_yet_submitted, submitted, handled_not_created_incident = [], [], []

        for finding in self.not_yet_submitted_findings:
            if finding.submitted():
                if finding not in self.submitted_findings:
                    submitted.append(finding)
            elif finding.failed_to_submit():
                if not finding.incident_created:
                    handled_not_created_incident.append(finding)
            else:
                not_yet_submitted.append(finding)

        for finding in self.submitted_findings:
            if finding.handled():
                if not finding.incident_created:
                    handled_not_created_incident.append(finding)
            else:
                submitted.append(finding)

        self.not_yet_submitted_findings = not_yet_submitted
        self.submitted_findings = submitted

        return handled_not_created_incident

    @classmethod
    def from_json(cls, cache_dict: dict[str, Any]) -> "Cache":
        """Deserialization method.

        Args:
            cache_dict: The cache dict in JSON format.

        Returns:
            An instance of the Cache class constructed from JSON representation.
        """
        return cls(
            not_yet_submitted_findings=list(map(Finding.from_json, cache_dict.get(NOT_YET_SUBMITTED_FINDINGS, []))),
            submitted_findings=list(map(Finding.from_json, cache_dict.get(SUBMITTED_FINDINGS, []))),
        )

    @classmethod
    def load_from_integration_context(cls, integration_context: dict[str, Any]) -> "Cache":
        return Cache.from_json(json.loads(integration_context.get(CACHE, "{}")))

    def dump_to_integration_context(self) -> None:
        integration_context = get_integration_context()
        integration_context[CACHE] = json.dumps(self, default=lambda obj: obj.__dict__)
        set_integration_context(integration_context)


def get_fields_query_part(
    finding_data: dict[str, Any],
    prefix: str,
    fields: list[str],
    raw_dict: dict[str, Any] | None = None,
    add_backslash: bool = False,
) -> str:
    """Given the fields to search for in the findings and the prefix, creates the query part for splunk search.
    For example: if fields are ["user"], and the value of the "user" fields in the finding is ["u1", "u2"], and the
    prefix is "identity", the function returns: (identity="u1" OR identity="u2")

    Args:
        finding_data (dict): The finding.
        prefix (str): The prefix to attach to each value returned in the query.
        fields (list): The fields to search in the finding for.
        raw_dict (dict): The raw dict
        add_backslash (bool): For users that contains single backslash, we add one more

    Returns: The query part
    """
    if not raw_dict:
        raw_dict = raw_to_dict(finding_data.get("_raw", ""))
    raw_list: list = []
    for field_name in fields:
        raw_list += argToList(finding_data.get(field_name, "")) + argToList(raw_dict.get(field_name, ""))
    if add_backslash:
        raw_list = [item.replace("\\", "\\\\") for item in raw_list]
    raw_list = [f"""{prefix}="{item.strip('"')}\"""" for item in raw_list]

    if not raw_list:
        return ""
    elif len(raw_list) == 1:
        return raw_list[0]
    else:
        return f'({" OR ".join(raw_list)})'


def get_finding_field_and_value(
    raw_field: str, finding_data: dict[str, Any], raw: dict[str, Any] | None = None
) -> tuple[str, Any]:
    """Gets the value by the name of the raw_field. The raw field may contain a Splunk
    formatting suffix (e.g. "threat_match_field|s") while the actual field is "threat_match_field",
    so we compare against the base field name (the part before the first "|").

    We must NOT use a loose substring match here, because that causes field-name collisions
    (e.g. the field "src" would wrongly match the raw field "src_ip").

    Args:
        raw_field (str): The raw field
        finding_data (dict): The finding data
        raw (dict): The raw dict

    Returns: The value in the finding which is associated with raw_field

    """
    if not raw:
        raw = raw_to_dict(finding_data.get("_raw", ""))
    base_field = raw_field.split("|")[0].strip()
    if base_field in finding_data:
        return base_field, finding_data[base_field]
    if base_field in raw:
        return base_field, raw[base_field]
    demisto.error(f"Field {raw_field} was not found in the finding.")
    return "", ""


def earliest_time_exists_in_query(query: str) -> bool:
    """
    Returns True if the query contains 'earliest=' or 'earliest ='
    (any amount of whitespace around the equals sign).
    """
    if query is None:
        return False

    pattern = r"earliest\s*=\s*"
    return re.search(pattern, query) is not None


def build_drilldown_search(
    finding_data: dict[str, Any], search: str, raw_dict: dict[str, Any], is_query_name: bool = False
) -> str:
    """Replaces all needed fields in a drilldown search query, or a search query name
    Args:
        finding_data (dict): The finding data
        search (str): The drilldown search query
        raw_dict (dict): The raw dict
        is_query_name (bool): Whether the given query is a query name (default is false)

    Returns (str): A searchable drilldown search query or a parsed query name
    """
    searchable_search: list = []
    start = 0

    for match in re.finditer(DRILLDOWN_REGEX, search):
        groups = match.groups()
        prefix = groups[0]
        raw_field = (groups[1] or groups[2]).strip("$")
        field, replacement = get_finding_field_and_value(raw_field, finding_data, raw_dict)
        if not field and not replacement:
            if not is_query_name:
                demisto.error(f"Failed building drilldown search query. Field {raw_field} was not found in the finding.")
            return ""

        if prefix:
            if field in USER_RELATED_FIELDS:
                replacement = get_fields_query_part(finding_data, prefix, [field], raw_dict, add_backslash=True)
            else:
                replacement = get_fields_query_part(finding_data, prefix, [field], raw_dict)

        end = match.start()
        searchable_search.extend((search[start:end], str(replacement)))
        start = match.end()
    searchable_search.append(search[start:])  # Handling the tail of the query

    parsed_query = "".join(searchable_search)

    demisto.debug(f"Parsed query is: {parsed_query}")

    return parsed_query


def get_drilldown_timeframe(finding_data, raw) -> tuple[str, str]:
    """Sets the drilldown search timeframe data.

    Args:
        finding_data (dict): The finding
        raw (dict): The raw dict

    Returns:
        earliest_offset: The earliest time to query from.
        latest_offset: The latest time to query to.
    """
    earliest_offset = finding_data.get("drilldown_earliest", "")
    latest_offset = finding_data.get("drilldown_latest", "")
    info_min_time = raw.get(INFO_MIN_TIME, "")
    info_max_time = raw.get(INFO_MAX_TIME, "")

    if not earliest_offset or earliest_offset == f"${INFO_MIN_TIME}$":
        if info_min_time:
            earliest_offset = info_min_time
        else:
            demisto.debug("Failed retrieving info min time")
    if not latest_offset or latest_offset == f"${INFO_MAX_TIME}$":
        if info_max_time:
            latest_offset = info_max_time
        else:
            demisto.debug("Failed retrieving info max time")

    return earliest_offset, latest_offset


def escape_invalid_chars_in_drilldown_json(drilldown_search: str) -> str:
    """Goes over the drilldown search, and replace the unescaped or invalid chars.

    Args:
        drilldown_search (str): The drilldown search.

    Returns:
        str: The escaped drilldown search.
    """
    # escape the " of string from the form of 'some_key="value"' which the " char are invalid in json value
    for unescaped_val in re.findall(r"(?<==)\s*\"[^\"]*\"", drilldown_search):
        escaped_val = unescaped_val.replace('"', '\\"')
        drilldown_search = drilldown_search.replace(unescaped_val, escaped_val)

    # replace the new line (\n) with in the IN (...) condition with ','
    # Splunk replace the value of some multiline fields to the value which contain \n
    # due to the 'expandtoken' macro
    for multiline_val in re.findall(r"(?<=in|IN)\s*\([^\)]*\n[^\)]*\)", drilldown_search):
        csv_val = multiline_val.replace("\n", ",")
        drilldown_search = drilldown_search.replace(multiline_val, csv_val)
    return drilldown_search


def escape_backslashes_in_field_filters(search: str) -> str:
    """Re-escapes backslashes inside the value of `field="value"` filters in an SPL search.

    When a drilldown search arrives as a JSON string, ``json.loads`` decodes ``\\\\`` to a single
    backslash. Splunk SPL, however, requires backslashes inside a double-quoted filter value to be
    escaped (doubled) in order to match.

    To stay safe, this only touches values of genuine ``field="value"`` filters - the quoted value must
    be directly preceded by a field-name token and ``=`` (e.g. ``TaskName="..."``). This deliberately
    excludes regex/string literals that follow ``(`` or ``,`` inside SPL functions such as
    ``eval x=replace(field,"(\\)","\\\\")`` and leaves free-text / ``rex`` regex quoted strings
    untouched. It is idempotent - values that are already correctly escaped (``\\\\``) are left
    unchanged.

    Args:
        search (str): The decoded SPL drilldown search.

    Returns:
        str: The SPL search with backslashes escaped inside field filter values.
    """

    def _escape(match: re.Match) -> str:
        prefix = match.group(1)  # the field name, '=' and any surrounding whitespace
        value = match.group(2)  # the value between the double quotes
        normalized = value.replace("\\\\", "\\")  # normalize already-doubled backslashes
        escaped = normalized.replace("\\", "\\\\")  # then double every backslash (idempotent)
        return f'{prefix}"{escaped}"'

    # Anchor on a field-name token so only real `field="value"` filters match. This avoids
    # over-escaping regex literals inside function calls like replace(field,"(\\)","\\\\").
    return re.sub(r'([\w.]+\s*=\s*)"([^"]*)"', _escape, search)


def parse_drilldown_searches(drilldown_searches: list[str]) -> list[dict[str, Any]]:
    """Goes over the drilldown searches list, parses each drilldown search and converts it to a python dictionary.

    Args:
        drilldown_searches (list): The list of the drilldown searches.

    Returns:
        list[dict]: A list of the drilldown searches dictionaries.
    """
    demisto.debug("There are multiple drilldown searches to enrich, parsing each drilldown search object")
    searches: list[dict] = []

    def _fix_search_backslashes(search_obj: dict[str, Any]) -> dict[str, Any]:
        # Re-escape backslashes in the SPL search after json.loads collapsed them (XSUP-70829)
        if isinstance(search_obj, dict) and isinstance(search_obj.get("search"), str):
            search_obj["search"] = escape_backslashes_in_field_filters(search_obj["search"])
        return search_obj

    for drilldown_search in drilldown_searches:
        try:
            # drilldown_search may be a json list/dict represented as string
            drilldown_search = escape_invalid_chars_in_drilldown_json(drilldown_search)
            search = json.loads(drilldown_search)
            if isinstance(search, list):
                searches.extend(_fix_search_backslashes(s) for s in search)
            else:
                searches.append(_fix_search_backslashes(search))
        except json.JSONDecodeError as e:
            demisto.error(
                f"Caught an exception while parsing a drilldown search object."
                f"Drilldown search is: {drilldown_search}, Original Error is: {e!s}"
            )

    return searches


def get_drilldown_searches(finding_data: dict[str, Any]) -> list[dict[str, Any]]:
    """Extract the drilldown_searches from the finding_data.
    It can be a list of objects, a single object or a simple string that contains the query.

    Args:
        finding_data (dict): The finding data

    Returns: A list that contains dict/s of the drilldown data like: name, search etc or the simple search query.
    """
    # Multiple drilldown searches is a feature added to Enterprise Security v7.2.0.
    # from this version, if a user set a drilldown search, we get a list of drilldown search objects (under
    # the 'drilldown_searches' key) and submit a splunk enrichment for each one of them.
    # To maintain backwards compatibility we keep using the 'drilldown_search' key as well.

    if drilldown_search := finding_data.get("drilldown_search"):
        # The drilldown_searches are in 'old' format a simple string query.
        return [drilldown_search]
    if drilldown_search := finding_data.get("drilldown_searches", []):
        if isinstance(drilldown_search, list):
            # The drilldown_searches are a list of searches data stored as json strings:
            return parse_drilldown_searches(drilldown_search)
        else:
            # The drilldown_searches are a dict/list of the search data in a JSON string representation.
            return parse_drilldown_searches([drilldown_search])
    return []


def drilldown_enrichment(
    service: client.Service, finding_data: dict[str, Any], num_enrichment_events: int
) -> list[tuple[str | None, str | None, client.Job | None]]:
    """Performs a drilldown enrichment.
    If the finding has multiple drilldown searches, enriches all the drilldown searches.

    Args:
        service (splunklib.client.Service): Splunk service object.
        finding_data (dict): The finding data
        num_enrichment_events (int): The maximal number of events to return per enrichment type.

    Returns: A list that contains tuples of a query name, query search and the splunk job that runs the query.
             [(query_name, query_search, splunk_job)]
    """
    jobs_and_queries: list[tuple[str | None, str | None, client.Job | None]] = []
    demisto.debug(f"finding data is: {finding_data}")
    if searches := get_drilldown_searches(finding_data):
        raw_dict = raw_to_dict(finding_data.get("_raw", ""))

        total_searches = len(searches)
        demisto.debug(f"Finding {finding_data[EVENT_ID]} has {total_searches} drilldown searches to enrich")

        for i in range(total_searches):
            # Iterates over the drilldown searches of the given finding to enrich each one of them
            search = searches[i]
            demisto.debug(f"Enriches drilldown search number {i+1} out of {total_searches} for finding {finding_data[EVENT_ID]}")

            if isinstance(search, dict):
                query_name = search.get("name", "")
                query_search = search.get("search", "")
                earliest_offset = search.get("earliest") or search.get("earliest_offset", "")  # The earliest time to query from.
                latest_offset = search.get("latest") or search.get("latest_offset", "")  # The latest time to query to.

            else:
                # Got a single drilldown search under the 'drilldown_search' key (BC)
                query_search = search
                query_name = finding_data.get("drilldown_name", "")
                earliest_offset, latest_offset = get_drilldown_timeframe(finding_data, raw_dict)

            try:
                parsed_query_name = build_drilldown_search(finding_data, query_name, raw_dict, True)
                if not parsed_query_name:  # if parsing failed - keep original unparsed name
                    demisto.debug(
                        f"Failed parsing drilldown search query name, using the original "
                        f"un-parsed query name instead: {query_name}."
                    )
                    parsed_query_name = query_name
            except Exception as e:
                demisto.error(f"Caught an exception while parsing the query name, using the original query name instead: {e!s}")
                parsed_query_name = query_name

            if searchable_query := build_drilldown_search(finding_data, query_search, raw_dict):
                demisto.debug(f"Search Query was build successfully for finding {finding_data[EVENT_ID]}")

                if (earliest_offset and latest_offset) or earliest_time_exists_in_query(searchable_query):
                    kwargs = {"max_count": num_enrichment_events, "exec_mode": "normal"}
                    if latest_offset:
                        kwargs["latest_time"] = latest_offset
                    if earliest_offset:
                        kwargs["earliest_time"] = earliest_offset
                    query = build_search_query({"query": searchable_query})
                    demisto.debug(f"Drilldown query for finding {finding_data[EVENT_ID]} is: {query}")
                    try:
                        job = service.jobs.create(query, **kwargs)
                        jobs_and_queries.append((parsed_query_name, query, job))

                    except Exception as e:
                        demisto.error(f"Caught an exception in drilldown_enrichment function: {e!s}")
                else:
                    demisto.debug(f"Failed getting the drilldown timeframe for finding {finding_data[EVENT_ID]}")
                    jobs_and_queries.append((None, None, None))
            else:
                demisto.debug(
                    f"Couldn't build search query for finding {finding_data[EVENT_ID]} "
                    f"with the following drilldown search {query_search}"
                )
                jobs_and_queries.append((None, None, None))
    else:
        demisto.debug(f"drill-down was not properly configured for finding {finding_data[EVENT_ID]}")
        jobs_and_queries.append((None, None, None))

    return jobs_and_queries


def identity_enrichment(service: client.Service, finding_data: dict[str, Any], num_enrichment_events: int) -> client.Job | None:
    """Performs an identity enrichment.

    Args:
        service (splunklib.client.Service): Splunk service object
        finding_data (dict): The finding data
        num_enrichment_events (int): The maximal number of events to return per enrichment type.

    Returns: The Splunk Job
    """
    job = None
    error_msg = f"Failed submitting identity enrichment request to Splunk for finding {finding_data[EVENT_ID]}"
    if users := get_fields_query_part(
        finding_data=finding_data,
        prefix="identity",
        fields=USER_RELATED_FIELDS,
        add_backslash=True,
    ):
        tables = argToList(demisto.params().get("identity_enrich_lookup_tables", DEFAULT_IDENTITY_ENRICH_TABLE))
        query = ""
        for table in tables:
            query += f"| inputlookup {table} where {users}"
        demisto.debug(f"Identity query for finding {finding_data[EVENT_ID]}: {query}")
        try:
            kwargs = {"max_count": num_enrichment_events, "exec_mode": "normal"}
            job = service.jobs.create(query, **kwargs)
        except Exception as e:
            demisto.error(f"Caught an exception in identity_enrichment function: {e!s}")
    else:
        demisto.debug(f"No users were found in finding. {error_msg}")

    return job


def asset_enrichment(service: client.Service, finding_data: dict[str, Any], num_enrichment_events: int) -> client.Job | None:
    """Performs an asset enrichment.

    Args:
        service (splunklib.client.Service): Splunk service object
        finding_data (dict): The finding data
        num_enrichment_events (int): The maximal number of events to return per enrichment type.

    Returns: The Splunk Job
    """
    job = None
    error_msg = f"Failed submitting asset enrichment request to Splunk for finding {finding_data[EVENT_ID]}"
    if assets := get_fields_query_part(
        finding_data=finding_data,
        prefix="asset",
        fields=["src", "dest", "src_ip", "dst_ip"],
    ):
        tables = argToList(demisto.params().get("asset_enrich_lookup_tables", DEFAULT_ASSET_ENRICH_TABLES))

        query = ""
        for table in tables:
            query += f"| inputlookup append=T {table} where {assets}"
        query += "| rename _key as asset_id | stats values(*) as * by asset_id"

        demisto.debug(f"Asset query for finding {finding_data[EVENT_ID]}: {query}")
        try:
            kwargs = {"max_count": num_enrichment_events, "exec_mode": "normal"}
            job = service.jobs.create(query, **kwargs)
        except Exception as e:
            demisto.error(f"Caught an exception in asset_enrichment function: {e!s}")
    else:
        demisto.debug(f"No assets were found in finding. {error_msg}")

    return job


def handle_submitted_findings(service: client.Service, cache_object: Cache) -> list[Finding]:
    """Handles submitted findings. For each submitted finding, tries to retrieve its results, if results aren't ready,
     it moves to the next submitted finding.

    Args:
        service (splunklib.client.Service): Splunk service object.
        cache_object (Cache): The enrichment mechanism cache object

    Returns:
        handled_findings (list[Finding]): The handled Findings
    """
    handled_findings = []
    if not (enrichment_timeout := arg_to_number(str(demisto.params().get("enrichment_timeout", "5")))):
        enrichment_timeout = 5
    findings = cache_object.submitted_findings
    total = len(findings)
    demisto.debug(f"Trying to handle {len(findings[:MAX_HANDLE_FINDINGS])}/{total} open enrichments")

    for finding in findings[:MAX_HANDLE_FINDINGS]:
        if handle_submitted_finding(service, finding, enrichment_timeout):
            handled_findings.append(finding)

    cache_object.submitted_findings = [n for n in findings if n not in handled_findings]

    if handled_findings:
        demisto.debug(f"Handled {len(handled_findings)}/{total} findings.")
    return handled_findings


def handle_submitted_finding(service: client.Service, finding: Finding, enrichment_timeout: int) -> bool:
    """Handles submitted finding. If enrichment process timeout has reached, creates an incident.

    Args:
        service (splunklib.client.Service): Splunk service object
        finding (Finding): The finding
        enrichment_timeout (int): The timeout for the enrichment process

    Returns:
        finding_status (str): The status of the finding
    """
    task_status = False

    if not finding.is_enrichment_process_exceeding_timeout(enrichment_timeout):
        demisto.debug(f"Trying to handle open enrichment for finding {finding.id}")
        for enrichment in finding.enrichments:
            if enrichment.status == Enrichment.IN_PROGRESS:
                try:
                    job = client.Job(service=service, sid=enrichment.id)
                    if job.is_done():
                        demisto.debug(f"Handling {enrichment.id=} of {enrichment.type=} for finding {finding.id}")
                        for item in results.JSONResultsReader(job.results(output_mode=OUTPUT_MODE_JSON)):
                            if handle_message(item):
                                continue
                            enrichment.data.append(item)
                        enrichment.status = Enrichment.SUCCESSFUL
                        demisto.debug(
                            f"{enrichment.id=} of {enrichment.type=} for finding {finding.id} status is successful "
                            f"{len(enrichment.data)=}"
                        )
                    else:
                        demisto.debug(f"{enrichment.id=} of {enrichment.type=} for finding {finding.id} is still not done")
                except Exception as e:
                    demisto.error(
                        f"Caught an exception while retrieving {enrichment.id=} of {enrichment.type=}\
                        results for finding {finding.id}: {e!s}"
                    )

                    enrichment.status = Enrichment.FAILED
                    demisto.error(f"{enrichment.id=} of {enrichment.type=} for finding {finding.id} was failed.")

        if finding.handled():
            task_status = True
            demisto.debug(f"Handled open enrichment for finding {finding.id}.")
        else:
            demisto.debug(f"Did not finish handling open enrichment for finding {finding.id}")

    else:
        task_status = True
        demisto.debug(
            f"Open enrichment for finding {finding.id} has exceeded the enrichment timeout of {enrichment_timeout}.\
            Submitting the finding without the enrichment."
        )

    return task_status


def submit_findings(service: client.Service, cache_object: Cache) -> tuple[list[Finding], list[Finding]]:
    """Submits fetched findings to Splunk for an enrichment.

    Args:
        service (splunklib.client.Service): Splunk service object
        cache_object (Cache): The enrichment mechanism cache object

    Returns:
        tuple[list[Finding], list[Finding]]: failed_findings, submitted_findings
    """
    failed_findings, submitted_findings = [], []
    num_enrichment_events = arg_to_number(str(demisto.params().get("num_enrichment_events", "20"))) or 20
    findings = cache_object.not_yet_submitted_findings
    total = len(findings)
    if findings:
        demisto.debug(f"Enriching {len(findings[:MAX_SUBMIT_FINDINGS])}/{total} fetched findings")

    for finding in findings[:MAX_SUBMIT_FINDINGS]:
        if submit_finding(service, finding, num_enrichment_events):
            cache_object.submitted_findings.append(finding)
            submitted_findings.append(finding)
            demisto.debug(f"Submitted enrichment request to Splunk for finding {finding.id}")
        else:
            failed_findings.append(finding)
            demisto.debug(f"Incident will be created from finding {finding.id} as each enrichment submission failed")

    cache_object.not_yet_submitted_findings = [n for n in findings if n not in submitted_findings + failed_findings]

    if submitted_findings:
        demisto.debug(f"Submitted {len(submitted_findings)}/{total} findings successfully.")

    if failed_findings:
        demisto.debug(
            f"The following {len(failed_findings)} findings failed the enrichment process: \
            {[finding.id for finding in failed_findings]}, \
            creating incidents without enrichment."
        )
    return failed_findings, submitted_findings


def submit_finding(service: client.Service, finding: Finding, num_enrichment_events: int) -> bool:
    """Submits fetched finding to Splunk for an Enrichment. Three enrichments possible: Drilldown, Asset & Identity.
     If all enrichment type executions were unsuccessful, creates a regular incident, Otherwise updates the
     integration context for the next fetch to handle the submitted finding.

    Args:
        service (splunklib.client.Service): Splunk service object
        finding (Finding): The finding.
        num_enrichment_events (int): The maximal number of events to return per enrichment type.

    Returns:
        task_status (bool): True if any of the enrichment's succeeded to be submitted to Splunk, False otherwise
    """
    submitted_drilldown, submitted_asset, submitted_identity = finding.get_submitted_enrichments()

    if DRILLDOWN_ENRICHMENT in ENABLED_ENRICHMENTS and not submitted_drilldown:
        jobs_and_queries = drilldown_enrichment(service, finding.data, num_enrichment_events)
        for job_and_query in jobs_and_queries:
            finding.enrichments.append(
                Enrichment.from_job(
                    DRILLDOWN_ENRICHMENT, job=job_and_query[2], query_name=job_and_query[0], query_search=job_and_query[1]
                )
            )
    if ASSET_ENRICHMENT in ENABLED_ENRICHMENTS and not submitted_asset:
        job = asset_enrichment(service, finding.data, num_enrichment_events)
        finding.enrichments.append(Enrichment.from_job(ASSET_ENRICHMENT, job))
    if IDENTITY_ENRICHMENT in ENABLED_ENRICHMENTS and not submitted_identity:
        job = identity_enrichment(service, finding.data, num_enrichment_events)
        finding.enrichments.append(Enrichment.from_job(IDENTITY_ENRICHMENT, job))

    return finding.submitted()


def create_incidents_from_findings(findings_to_be_created: list[Finding], mapper: UserMappingObject) -> list[dict[str, Any]]:
    """Create the actual incident from the handled Findings
        in addition, taking in account the data from the integration_context (from mirror-in process)
        about Findings which was updated by mirror-in during the Enrichment time.

    Args:
        findings_to_be_created (list[Finding]): The Findings to create incidents from (handled + failed enrichment Findings).
        mapper (UserMappingObject): a UserMappingObject object

    Returns:
        incidents (list[dict]): The created incidents.
    """
    integration_context = None
    mirrored_in_findings = {}
    incidents: list[dict] = []

    if is_mirror_in_enabled():
        integration_context = get_integration_context()
        mirrored_in_findings = integration_context.get(MIRRORED_ENRICHING_FINDINGS, {})
        demisto.debug(f"found {len(mirrored_in_findings)} enriched findings updated in mirror-in")
        demisto.debug(f"{mirrored_in_findings=}")

    for finding in findings_to_be_created:
        # in case the Finding was updated in Splunk between the time of fetch and create incident,
        # we need to take the updated delta.
        if finding.id in mirrored_in_findings:
            delta = mirrored_in_findings[finding.id]
            finding.data |= delta
            del mirrored_in_findings[finding.id]

        incidents.append(finding.to_incident(mapper))
    if integration_context:
        set_integration_context(integration_context)
    return incidents


def is_mirror_in_enabled() -> bool:
    params = demisto.params()
    return MIRROR_DIRECTION.get(params.get("mirror_direction", "")) in ["Both", "In"]


def run_enrichment_mechanism(
    service: client.Service, integration_context: dict[str, Any], mapper: UserMappingObject
) -> tuple[list[dict[str, Any]], dict[str, Any]]:
    """Execute the enriching fetch mechanism
    1. We first handle submitted findings that have not been handled in the last fetch run
    2. If we finished handling and submitting all fetched findings, we fetch new findings
    3. After we finish to fetch new findings or if we have left findings that have not been submitted, we submit
       them for an enrichment to Splunk
    4. Finally and in case of an Exception, we store the current cache object state in the integration context

    Args:
        service (splunklib.client.Service): Splunk service object.
        integration_context (dict): The integration context
        mapper (UserMappingObject): The user-mapping helper.

    Returns:
        A tuple of ``(incidents, last_run_delta)``:
          * ``incidents`` — the list of incidents produced this cycle, to be
            emitted once by the dispatcher (mixed-mode safe).
          * ``last_run_delta`` — the partial last-run dict produced by the
            inner :func:`fetch_findings` call (time-window cursor +
            ``next_run_found_incidents_ids`` for dedup, plus the ``DUMMY``
            marker). The dispatcher merges this into its single final
            ``demisto.setLastRun`` so the dedup IDs survive across cycles.
            Empty when no new fetch was performed this cycle (i.e. we were
            still draining the submitted/handled pipeline).
    """
    incidents: list = []
    last_run_delta: dict[str, Any] = {}
    cache_object = Cache.load_from_integration_context(integration_context)

    try:
        handled_findings = handle_submitted_findings(service, cache_object)
        if cache_object.done_submitting() and cache_object.done_handling():
            # fetch_findings returns a FetchResult; surface the delta to the
            # caller (dispatcher) instead of persisting here, so the dispatcher's
            # final setLastRun does not overwrite our dedup IDs / time cursor
            # with its stale pre-fetch snapshot.
            findings_result = fetch_findings(
                service=service,
                cache_object=cache_object,
                enrich_findings=True,
                mapper=mapper,
            )
            last_run_delta = dict(findings_result.last_run_delta or {})
            if is_mirror_in_enabled():
                # if mirror-in enabled, we need to store in cache the fetched findings ASAP,
                # as they need to be able to update by the mirror in process
                demisto.debug("dumping the cache object direct after fetch as mirror-in enabled")
                cache_object.dump_to_integration_context()

        failed_findings, _ = submit_findings(service, cache_object)
        incidents = create_incidents_from_findings(handled_findings + failed_findings, mapper)
    except Exception as e:
        err = f"Caught an exception while executing the enriching fetch mechanism. Additional Info: {e!s}"
        demisto.error(err)
        # we throw exception only if there is no incident to create
        if not incidents:
            raise e

    finally:
        store_incidents_for_mapping(incidents)
        handled_but_not_created_incidents = cache_object.organize()
        cache_object.dump_to_integration_context()
        incidents += [finding.to_incident(mapper) for finding in handled_but_not_created_incidents]
    # Return incidents + last_run_delta so the dispatcher performs a single
    # emit / setLastRun per cycle (mixed-mode safe + dedup-IDs preserved).
    return incidents, last_run_delta


def store_incidents_for_mapping(incidents: list[dict[str, Any]]) -> None:
    """Stores ready incidents in integration context to allow the mapping to pull the incidents from the instance.
    We store at most 20 incidents.

    Args:
        incidents (list): The incidents
    """
    if incidents:
        integration_context = get_integration_context()
        integration_context[INCIDENTS] = incidents[:20]
        set_integration_context(integration_context)


def fetch_incidents_for_mapping(integration_context: dict[str, Any]) -> None:
    """Gets the stored incidents to the "Pull from instance" in Classification & Mapping (In case of enriched fetch)

    Args:
        integration_context (dict): The integration context
    """
    incidents = integration_context.get(INCIDENTS, [])
    demisto.debug(f'Retrieving {len(incidents)} incidents for "Pull from instance" in Classification & Mapping.')
    demisto.incidents(incidents)


def reset_enriching_fetch_mechanism() -> None:
    """Resets all the fields regarding the enriching fetch mechanism & the last run object"""

    # keys: INCIDENTS, CACHE, MIRRORED_ENRICHING_FINDINGS, PROCESSED_MIRRORED_EVENTS may exist in context
    set_integration_context({})
    demisto.setLastRun({})
    return_results("Enriching fetch mechanism was reset successfully.")


# =========== Mirroring Mechanism ===========


def format_splunk_note_for_xsoar(note: dict, timezone_offset: str = "+00:00") -> str:
    """Formats a Splunk note with author and timestamp in user's timezone.

    Args:
        note: The note dictionary from Splunk
        timezone_offset: Timezone offset string (e.g., '+02:00', '-05:00')

    Format:
    **author**  timestamp

    title

    note_content
    """
    # Extract author information - handle nested author object
    author = demisto.get(note, "author.username", "Unknown")

    # Extract and format timestamp with timezone
    timestamp = ""
    if time_value := note.get("update_time"):
        try:
            # Create timezone-aware datetime from epoch timestamp
            dt_utc = datetime.fromtimestamp(float(time_value), tz=pytz.UTC)

            # Parse timezone offset (e.g., '+02:00' -> hours=2, minutes=0)
            sign = 1 if timezone_offset[0] == "+" else -1
            hours = int(timezone_offset[1:3])
            minutes = int(timezone_offset[4:6])
            offset_minutes = sign * (hours * 60 + minutes)

            # Apply timezone offset
            dt_local = dt_utc + timedelta(minutes=offset_minutes)
            timestamp = dt_local.strftime("%b %d, %I:%M %p")

        except (ValueError, TypeError) as e:
            demisto.error(f"Failed to format timestamp {time_value}: {e}")
            timestamp = ""

    # Extract note content
    raw_title = note.get("title") or ""
    decoded_title = urllib.parse.unquote(raw_title)

    raw_content = note.get("content") or ""
    decoded_content = urllib.parse.unquote(raw_content)

    # Combine title and content
    note_text = decoded_title
    if decoded_content:
        note_text = f"{decoded_title}\n{decoded_content}" if decoded_title else decoded_content

    # Format with author and timestamp header
    header = f"**{author}**"
    if timestamp:
        header += f"  {timestamp}"

    # Return formatted note with header and content separated by blank line
    if note_text:
        return f"{header}\n\n{note_text}"
    else:
        return header


def get_war_room_note_entry(content: str, finding_id: str, format: str) -> dict[str, Any]:
    return {
        "EntryContext": {"mirrorRemoteId": finding_id},
        "Type": EntryType.NOTE,
        "Contents": content,
        "ContentsFormat": format,
        "Tags": [NOTE_TAG_FROM_SPLUNK],  # The list of tags to add to the entry
        "Note": True,
    }


def enrich_with_splunk_notes(
    service: client.Service,
    id_to_entity_map: dict[str, dict[str, Any]],
    last_update_splunk_timestamp: float | None = None,
    is_fetch: bool = False,
) -> list[dict[str, Any]]:
    """Get notes from Splunk with timezone-aware timestamps for findings or investigations.

    This implementation uses a search query to find note IDs associated with the entities (findings or
    investigations) from the _audit index, then retrieves the actual notes from the mc_notes KV store.

    Splunk notes reference their parent entity via either ``notable_id`` (for findings) or
    ``incident_id`` (for investigations). This function transparently supports both: for each note
    it picks whichever id field is present, so the same logic works for both entity types.

    Args:
        service (client.Service): Splunk service object
        id_to_entity_map (dict[str, dict]): Dictionary of entities (findings or investigations) keyed by entity id
        last_update_splunk_timestamp (str): Last update timestamp to filter notes (optional)
        is_fetch (bool): Whether the function is called from fetch

    Returns:
        list[dict]: The war room entries to create in XSOAR.
    """

    if not id_to_entity_map:
        return []

    # Get timezone offset (cached or from Splunk)
    timezone_offset = get_splunk_timezone_offset(service)
    demisto.debug(f"enrich_with_splunk_notes: Using timezone offset: {timezone_offset} for note timestamps")

    # Build the OR clause for the search query with all entity IDs
    entity_ids = list(id_to_entity_map.keys())
    or_clauses = [f'"{entity_id}"' for entity_id in entity_ids]
    or_clause_str = " OR ".join(or_clauses)

    # We request all notes associated with the entities, but limit the search to the last week
    # to avoid performance issues with large audit logs
    search_query = (
        f"search index=_audit source=mc_notes earliest=-7d ({or_clause_str}) "
        '| rex "(?<timestamp>[\\d.]+),(?<note_id>[\\w-]+),(?<user>[\\w_]+),(?<model>[\\w]+),(?<command>[\\w]+),(?<diff>.+)" '
        "| dedup note_id sortby -update_time"
        '| where command!="D"'  # filter out the deleted notes
        # '| table note_id, command, diff'
        "| table note_id"
    )

    demisto.debug(
        f"enrich_with_splunk_notes: Running fetch query to find the changed note IDs in {len(entity_ids)} entities, "
        f"{search_query=}"
    )

    try:
        # Execute the search query to get note IDs
        start_time = time.time()
        oneshotsearch_results = service.jobs.oneshot(
            search_query,
            output_mode=OUTPUT_MODE_JSON,
            count=0,  # No limit
        )
        reader = results.JSONResultsReader(oneshotsearch_results)

        note_ids = []
        for item in reader:
            if handle_message(item):
                continue
            if isinstance(item, dict) and item.get("note_id"):
                note_ids.append(item["note_id"])

        query_time = time.time() - start_time
        demisto.debug(
            f"enrich_with_splunk_notes: Search Note IDs from _audit completed in {query_time:.3f} sec, "
            f"found {len(note_ids)} note IDs"
        )

        if not note_ids:
            demisto.debug("enrich_with_splunk_notes: No note IDs found")
            return []

        # Now query the mc_notes KV store with the note IDs
        start_time = time.time()
        query = json.dumps({"id": {"$in": note_ids}})
        mc_notes = service.kvstore["mc_notes"].data.query(query=query)
        query_time = time.time() - start_time
        demisto.debug(
            f"enrich_with_splunk_notes: mc_notes KV store query completed in {query_time:.3f} sec, "
            f"retrieved {len(mc_notes)} notes"
        )

    except Exception as e:
        demisto.error(f"enrich_with_splunk_notes: Failed to query notes: {e}\n{traceback.format_exc()}")
        return []

    # Process the retrieved notes
    # Collect all notes grouped by entity_id for sorting
    entity_notes = defaultdict(list)
    war_room_notes = []
    for note in mc_notes:
        # Notes can reference either a finding (notable_id) or an investigation (incident_id).
        # Pick whichever id is present so this function works transparently for both entity types.
        entity_id = note.get("notable_id") or note.get("incident_id")
        if not entity_id:
            demisto.debug(f"enrich_with_splunk_notes: Skipping note without notable_id/incident_id: {note}")
            continue

        # Collect note with update time for sorting
        if entity_id in id_to_entity_map:
            markdown_content = format_splunk_note_for_xsoar(note, timezone_offset)
            entity_notes[entity_id].append({"content": markdown_content, "update_time": float(note.get("update_time", 0))})

            # Creating a XSOAR war room note for a new note ONLY
            # in fetch - we don't create Entry notes for notes
            if (
                not is_fetch
                and last_update_splunk_timestamp
                and note.get("create_time", 0) > int(last_update_splunk_timestamp)
                and COMMENT_MIRRORED_FROM_XSOAR not in markdown_content
            ):
                war_room_notes.append(get_war_room_note_entry(markdown_content, entity_id, EntryFormat.MARKDOWN))
        else:
            demisto.debug(f"enrich_with_splunk_notes: Skipping note without matching entity: {note}")

    # Sort notes by update time and update the notes in each entity
    extensive_log(f"enrich_with_splunk_notes: entity_notes = {entity_notes}")
    for entity_id, splunk_notes in entity_notes.items():
        # Sort splunk notes by update_time (newest first)
        sorted_splunk_notes = sorted(splunk_notes, key=lambda x: x["update_time"], reverse=True)  # type: ignore[arg-type,return-value]
        splunk_notes_list = [{"Note": splunk_note["content"]} for splunk_note in sorted_splunk_notes]
        # splunk_notes key mapped in the Splunk Finding/Investigation - Incoming Mapper
        id_to_entity_map[entity_id]["splunk_notes"] = splunk_notes_list
    # handle a case of deleting all the notes

    return war_room_notes


def enrich_with_splunk_notes_v2(
    service: client.Service,
    id_to_entity_map: dict[str, dict[str, Any]],
    last_update_splunk_timestamp: float | None = None,
    is_fetch: bool = False,
) -> list[dict[str, Any]]:
    """Get notes from Splunk with timezone-aware timestamps for findings or investigations.

    This v2 implementation directly queries the ``mc_notes`` KV store for notes whose
    ``update_time`` is greater than or equal to the provided ``last_update_splunk_timestamp``,
    avoiding the need to first scan the ``_audit`` index for note IDs.

    Splunk notes reference their parent entity via either ``notable_id`` (for findings) or
    ``incident_id`` (for investigations). This function transparently supports both: for each note
    it picks whichever id field is present, so the same logic works for both entity types.

    Args:
        service (client.Service): Splunk service object
        id_to_entity_map (dict[str, dict]): Dictionary of entities (findings or investigations) keyed by entity id
        last_update_splunk_timestamp (str): Last update timestamp to filter notes (optional)
        is_fetch (bool): Whether the function is called from fetch

    Returns:
        list[dict]: The war room entries to create in XSOAR.
    """

    if not id_to_entity_map:
        return []

    # Get timezone offset (cached or from Splunk)
    timezone_offset = get_splunk_timezone_offset(service)
    demisto.debug(f"enrich_with_splunk_notes_v2: Using timezone offset: {timezone_offset} for note timestamps")

    entity_ids = list(id_to_entity_map.keys())

    # Filter notes by either `notable_id` (findings) or `incident_id` (investigations)
    # so we only fetch notes that belong to the entities passed in `id_to_entity_map`.
    # The additional `update_time` filter (added below) bounds the lookup to the last 7 days.
    query_filter: dict[str, Any] = {
        "$or": [
            {"notable_id": {"$in": entity_ids}},
            {"incident_id": {"$in": entity_ids}},
        ]
    }

    # Always query notes from the last 7 days (regardless of last_update_splunk_timestamp).
    # `last_update_splunk_timestamp` is used later only to decide which notes should be
    # surfaced as XSOAR war-room note entries.
    # update_time in mc_notes is stored as a numeric epoch value.
    seven_days_ago_epoch = (datetime.now(tz=UTC) - timedelta(days=7)).timestamp()
    query_filter["update_time"] = {"$gte": seven_days_ago_epoch}

    query = json.dumps(query_filter)
    demisto.debug(
        f"enrich_with_splunk_notes_v2: Querying mc_notes KV store directly for {len(entity_ids)} entities, " f"{query=}"
    )

    try:
        start_time = time.time()
        mc_notes = service.kvstore["mc_notes"].data.query(query=query)
        query_time = time.time() - start_time
        demisto.debug(
            f"enrich_with_splunk_notes_v2: mc_notes KV store query completed in {query_time:.3f} sec, "
            f"retrieved {len(mc_notes)} notes"
        )
    except Exception as e:
        demisto.error(f"enrich_with_splunk_notes_v2: Failed to query notes: {e}\n{traceback.format_exc()}")
        return []

    if not mc_notes:
        demisto.debug("enrich_with_splunk_notes_v2: No notes found")
        return []

    # Process the retrieved notes
    # Collect all notes grouped by entity_id for sorting
    entity_notes = defaultdict(list)
    war_room_notes = []
    for note in mc_notes:
        # Notes can reference either a finding (notable_id) or an investigation (incident_id).
        # Pick whichever id is present so this function works transparently for both entity types.
        entity_id = note.get("notable_id") or note.get("incident_id")
        if not entity_id:
            demisto.debug(f"enrich_with_splunk_notes_v2: Skipping note without notable_id/incident_id: {note}")
            continue

        # Collect note with update time for sorting
        if entity_id in id_to_entity_map:
            markdown_content = format_splunk_note_for_xsoar(note, timezone_offset)
            entity_notes[entity_id].append({"content": markdown_content, "update_time": float(note.get("update_time", 0))})

            # Creating a XSOAR war room note for a new note ONLY
            # in fetch - we don't create Entry notes for notes
            if (
                not is_fetch
                and last_update_splunk_timestamp
                and float(note.get("create_time", 0)) > float(last_update_splunk_timestamp)
                and COMMENT_MIRRORED_FROM_XSOAR not in markdown_content
            ):
                war_room_notes.append(get_war_room_note_entry(markdown_content, entity_id, EntryFormat.MARKDOWN))
        else:
            demisto.debug(f"enrich_with_splunk_notes_v2: Skipping note without matching entity: {note}")

    # Sort notes by update time and update the notes in each entity
    extensive_log(f"enrich_with_splunk_notes_v2: entity_notes = {entity_notes}")
    for entity_id, splunk_notes in entity_notes.items():
        # Sort splunk notes by update_time (newest first)
        sorted_splunk_notes = sorted(splunk_notes, key=lambda x: x["update_time"], reverse=True)  # type: ignore[arg-type,return-value]
        splunk_notes_list = [{"Note": splunk_note["content"]} for splunk_note in sorted_splunk_notes]
        # splunk_notes key mapped in the Splunk Finding/Investigation - Incoming Mapper
        id_to_entity_map[entity_id]["splunk_notes"] = splunk_notes_list
    # handle a case of deleting all the notes

    return war_room_notes


def handle_enriching_findings(modified_findings: dict[str, dict[str, Any]]) -> None:
    """Store the mirror in "delta" of the findings which not yet created because of enrichment mechanism.

    Args:
        modified_findings (dict[str, str]): The Findings changes from get-modified-remote-data
    """
    try:
        integration_context = get_integration_context()
        cache_object = Cache.load_from_integration_context(integration_context)
        if enriching_findings := (cache_object.submitted_findings + cache_object.not_yet_submitted_findings):
            enriched_and_changed = [finding for finding in enriching_findings if finding.id in modified_findings]
            if enriched_and_changed:
                demisto.debug(f"mirror-in: found {len(enriched_and_changed)} submitted findings, updating delta in cache.")
                delta_map = integration_context.get(MIRRORED_ENRICHING_FINDINGS, {})
                for finding in enriched_and_changed:
                    updated_finding = modified_findings[finding.id]
                    delta = delta_map.get(finding.id, {})
                    delta |= {k: v for k, v in updated_finding.items() if finding.data.get(k) != v}
                    delta_map[finding.id] = delta
                    # delete it from the modified_findings as it still not exist in the server as incident
                    del modified_findings[finding.id]

                integration_context[MIRRORED_ENRICHING_FINDINGS] = delta_map
                extensive_log(f"delta map after mirror update: {delta_map}")
                set_integration_context(integration_context)
                demisto.debug(f"mirror-in: delta updated for the enriching findings - {[n.id for n in enriched_and_changed]}")
            else:
                demisto.debug("mirror-in: enriching findings was not updated in remote.")
        else:
            demisto.debug("mirror-in: no enriching findings found.")
    except Exception as e:
        demisto.error(f"mirror-in: failed to check for enriching findings, {e}")


def handle_closed_entities(
    modified_entities_map: dict[str, dict[str, Any]],
    close_extra_labels: list[str],
    close_end_statuses: bool,
    entries: list[dict[str, Any]],
) -> None:
    """Append a `dbotIncidentClose` entry for any remote row whose Splunk-side
    status indicates closure. The same mechanism is used for both Findings and
    Investigations: every row passed in is expected to expose the canonical
    closure keys ``status_label`` and ``status_end`` (Investigation rows are
    normalised to this shape by :func:`parse_investigation`).

    Args:
        modified_entities_map: ``{entity_id: row}`` where ``entity_id`` is the
            ``rule_id`` for Findings or ``investigation_guid`` for Investigations.
        close_extra_labels: Additional Splunk status labels that should also
            be treated as "closed".
        close_end_statuses: If True, also close when the row's ``status_end``
            is truthy.
        entries: List that will be extended in place with closure entries.
    """
    demisto.debug("Starting handling closing the entity")
    for entity_id, entity in modified_entities_map.items():
        status_label = entity.get("status_label", "")
        status_end = argToBoolean(entity.get("status_end", "false"))
        demisto.debug(
            f"handle_closed_entities: Evaluating closure for {entity_id}: status_label={status_label}, "
            f"status_end={status_end}, close_extra_labels={close_extra_labels}, "
            f"close_end_statuses={close_end_statuses}"
        )

        should_close = (status_label == "Closed") or (status_label in close_extra_labels) or (close_end_statuses and status_end)

        if should_close:
            demisto.info(
                f"handle_closed_entities: closing incident for {entity_id} "
                f"(status_label={status_label}, status_end={status_end})"
            )
            reason = (
                f'Splunk event was closed on Splunk with status "{status_label}".'
                if status_label
                else "Splunk event was closed on Splunk based on end status."
            )
            entries.append(
                {
                    "EntryContext": {"mirrorRemoteId": entity_id},
                    "Type": EntryType.NOTE,
                    "Contents": {
                        "dbotIncidentClose": True,
                        "closeReason": reason,
                    },
                    "ContentsFormat": EntryFormat.JSON,
                }
            )


# Backwards-compatible alias. Existing callers/tests that import the old name
# (``handle_closed_findings``) continue to work unchanged — the mechanism is
# identical for Findings and Investigations now.
handle_closed_findings = handle_closed_entities


def get_modified_remote_data_command(
    service: client.Service,
    args: dict[str, Any],
    close_incident: bool,
    close_end_statuses: bool,
    close_extra_labels: list[str],
    mapper: UserMappingObject,
) -> None:
    """Gets the list of the findings data that have changed since a given time

    Args:
        service (splunklib.client.Service): Splunk service object
        args (dict): The command arguments
        close_incident (bool): Indicates whether to close the corresponding XSOAR incident if the finding
            has been closed on Splunk's end.
        close_end_statuses (bool): Specifies whether "End Status" statuses on Splunk should be closed when mirroring.
        close_extra_labels (list[str]): A list of additional Splunk status labels to close during mirroring.
        mapper (UserMappingObject): mapper to map the Splunk Username to the correct XSOAR username.

    Returns:
        SplunkGetModifiedRemoteDataResponse: The response containing the list of findings changed
    """
    remote_args = GetModifiedRemoteDataArgs(args)

    # Caching Mechanism for Handling Splunk Indexing Delays:
    # 1. A 60-second buffer is subtracted from the last run time to create an overlapping query window.
    #    This ensures we catch events that were indexed late by Splunk and missed in the previous run.
    # 2. To prevent processing duplicate events from this overlap, we cache the unique key (finding_id:timestamp)
    #    of every event processed in the current run.
    # 3. This cache is stored in the integration context and loaded at the start of the next run.
    # 4. Any fetched event whose key exists in the cache is skipped as a duplicate.
    integration_context = get_integration_context()
    processed_events_cache = set(integration_context.get(PROCESSED_MIRRORED_EVENTS, []))
    demisto.debug(f"Loaded {len(processed_events_cache)} processed events from cache.")

    # Build the query with the 60-second look-behind buffer.
    last_update_dt = dateparser.parse(remote_args.last_update, settings={"TIMEZONE": "UTC"})
    if not last_update_dt:
        raise DemistoException(f"Failed to parse last update time: {remote_args.last_update}")
    original_last_update_timestamp = last_update_dt.timestamp()
    demisto.debug(f"mirror-in: {remote_args.last_update=}, {original_last_update_timestamp=}")
    last_update_splunk_timestamp = original_last_update_timestamp - SPLUNK_INDEXING_TIME

    # Selected event types — matches the FetchHandlerFactory default so instances
    # that never set the param continue to mirror Findings.
    params = demisto.params()
    selected_types = argToList(params.get("fetch_event_types")) or ["Finding"]

    modified_findings_map: dict[str, dict[str, Any]] = {}
    entries: list[dict] = []
    current_run_processed_events: set[str] = set()

    if "Finding" in selected_types:
        # Query the audit index to get modified findings
        # This query extracts last_modified_timestamp and rule_id and other findings keys from the audit logs,
        # sorts by timestamp (descending), deduplicates by rule_id
        mirror_in_search = (
            f"search index=_audit source=notable_update_rest_handler earliest={last_update_splunk_timestamp} "
            "| eval last_modified_timestamp = _time "
            "| sort last_modified_timestamp DESC "
            "| dedup rule_id "
            "| `get_current_status` "
            "| eval review_time = last_modified_timestamp"
            "| table review_time, rule_id, owner, "
            "status, status_label, status_end, disposition, disposition_label, urgency, sensitivity"
        )

        demisto.debug(f"mirror-in: performing audit log search with query: {mirror_in_search}.")

        # Execute the search query to get modified findings
        for item in results.JSONResultsReader(
            service.jobs.oneshot(query=mirror_in_search, count=MIRROR_LIMIT, output_mode=OUTPUT_MODE_JSON)
        ):
            if handle_message(item):
                continue

            # Parse the finding data from the audit log
            updated_finding = parse_finding(item, to_dict=True)

            # Deduplication Mechanism:
            # Create a unique key for the event and check against the cache of previously processed events.
            finding_id = updated_finding.get("rule_id")
            last_modified = updated_finding.get("review_time")
            if not finding_id or not last_modified:
                continue
            event_key = f"{finding_id}:{last_modified}"
            if event_key in processed_events_cache:
                extensive_log(f"mirror-in: Skipping already processed event: {event_key}")
                continue

            # This is a new event. Add it to the map for processing and to the cache for the next run.
            modified_findings_map[finding_id] = updated_finding
            current_run_processed_events.add(event_key)

        if modified_findings_map:
            # Since ES version is 8.2+, notes are in the mc_notes KV Store
            # We use the query-based approach to fetch notes
            # war_room_notes = enrich_with_splunk_notes(service, modified_findings_map, original_last_update_timestamp)
            war_room_notes = enrich_with_splunk_notes_v2(service, modified_findings_map, original_last_update_timestamp)
            entries.extend(war_room_notes)

            mapper.map_owner_to_xsoar_user(modified_findings_map.values())  # type: ignore[arg-type]

            if ENABLED_ENRICHMENTS:
                handle_enriching_findings(modified_findings_map)

            if close_incident:
                handle_closed_entities(modified_findings_map, close_extra_labels, close_end_statuses, entries)

            demisto.debug(f"mirror-in: updated finding ids: {list(modified_findings_map.keys())}")

        else:
            demisto.debug(f"mirror-in: no findings was changed since {last_update_splunk_timestamp}")
        if len(modified_findings_map) >= MIRROR_LIMIT:
            demisto.info(f"mirror-in: the number of mirrored findings reach the limit of: {MIRROR_LIMIT}")
    else:
        demisto.debug("mirror-in: 'Finding' not in fetch_event_types — skipping audit log SPL query.")

    modified_data: list[dict[str, Any]] = list(modified_findings_map.values())

    # Investigation mirror-in. Hits the v2 endpoint directly
    # (no SPL). Gated on `fetch_event_types` so instances that only fetch Findings
    # skip the extra HTTP round trip.
    if "Investigation" in selected_types:
        iso8601_z_min = to_mc_iso8601_utc(last_update_dt)
        modified_investigations = list_modified_investigations(service, iso8601_z_min, mapper=mapper)
        demisto.debug(f"mirror-in: appending {len(modified_investigations)} modified investigation rows")

        # Enrich with Splunk notes (same pattern as Findings). The notes endpoint
        # references investigations via `incident_id`, which `enrich_with_splunk_notes`
        # accepts transparently. We key by `investigation_guid` so it matches the
        # `dbotMirrorId` stamped by Investigation.to_incident().
        if modified_investigations:
            # Deduplication Mechanism (mirrors the Findings flow above):
            # Filter out investigation rows we've already mirrored in a previous tick by
            # building a unique key per (investigation_guid, update_time) and checking it
            # against `processed_events_cache`. Newly-seen events are recorded in
            # `current_run_processed_events` so the next iteration can skip them too.
            deduped_investigations: list[dict[str, Any]] = []
            skipped_event_keys: list[str] = []
            for inv in modified_investigations:
                investigation_id = inv.get("investigation_guid")
                last_modified = inv.get("update_time")
                if not investigation_id or not last_modified:
                    continue
                event_key = f"{investigation_id}:{last_modified}"
                if event_key in processed_events_cache:
                    skipped_event_keys.append(event_key)
                    continue
                deduped_investigations.append(inv)
                current_run_processed_events.add(event_key)

            if skipped_event_keys:
                extensive_log(f"mirror-in: Skipping already processed investigations: {skipped_event_keys}")

            modified_investigations = deduped_investigations

            guid_to_investigation: dict[str, dict[str, Any]] = {
                str(inv["investigation_guid"]): inv for inv in modified_investigations if inv.get("investigation_guid")
            }
            if guid_to_investigation:
                # investigation_notes = enrich_with_splunk_notes(
                #     service, guid_to_investigation, original_last_update_timestamp
                # )
                investigation_notes = enrich_with_splunk_notes_v2(service, guid_to_investigation, original_last_update_timestamp)
                entries.extend(investigation_notes)

                # Same close-on-mirror-in mechanism as Findings: investigation rows
                # are normalised by `parse_investigation` to expose `status_label` /
                # `status_end`, so the SAME `handle_closed_entities` helper handles
                # both event types. Gated on the existing `close_incident` param.
                if close_incident:
                    handle_closed_entities(guid_to_investigation, close_extra_labels, close_end_statuses, entries)

        modified_data.extend(modified_investigations)
        if len(modified_investigations) >= MIRROR_LIMIT:
            demisto.info(f"mirror-in: the number of mirrored investigations reached the limit of: {MIRROR_LIMIT}")

    # Persist the cache of events (findings + investigations) processed in this run
    # for the next iteration. Re-read the integration context immediately before
    # writing so we don't clobber updates other containers made to OTHER keys
    # between the initial get_integration_context() (early in this function) and now.
    # The PROCESSED_MIRRORED_EVENTS key itself is intentionally overwritten with this run's set.
    latest_context = get_integration_context()
    latest_context[PROCESSED_MIRRORED_EVENTS] = list(current_run_processed_events)
    set_integration_context(latest_context)

    res = SplunkGetModifiedRemoteDataResponse(modified_findings_data=modified_data, entries=entries)
    return_results(res)


def get_enterprise_security_version(service: client.Service) -> str:
    """Retrieves the installed Splunk Enterprise Security (ES) app version.

    Args:
        service (splunklib.client.Service): Splunk service object.

    Returns:
        str: The ES app version (e.g. "8.2.0"), or "unknown" if it could not be determined.
    """
    try:
        es_app = service.apps[ES_APP_NAME]
        return es_app.content.get("version", "unknown")
    except Exception as e:
        demisto.debug(f"Could not determine Enterprise Security version: {e!s}\n{traceback.format_exc()}")
        return "unknown"


def is_es_version(version: str, target_version: str) -> bool:
    """Return True when ``version`` matches ``target_version`` on major/minor.

    Uses ``packaging.version.Version`` for robust semantic comparison, relying
    on the library's own ``.major``/``.minor`` properties, and checks that the
    major/minor of ``version`` equals that of ``target_version`` (so for
    ``target_version="8.2"``, the values ``8.2``, ``8.2.0``, ``8.2.5`` all
    match, while ``8.3.0`` does not). Returns False when either version is
    unparseable (e.g. "unknown"). Used to gate version-specific behaviour such
    as the finding_time (notable_time) handling on ES ``8.2.x``.
    """
    try:
        parsed, target = Version(version), Version(target_version)
        return (parsed.major, parsed.minor) == (target.major, target.minor)
    except InvalidVersion:
        return False


def get_finding_time_for_es_notable_time(service: client.Service, data: dict[str, Any] | None) -> str | None:
    """Return the finding's event time to use as ``_time`` on ES ``>=8.2 <8.3``.

    On ES 8.2.x the v2 investigations update endpoint requires the finding's
    original event time (``_time``). On 8.3+ this is not needed. The
    finding's time is mirrored into the incident as the ``notable_time`` field
    (mapped from the incident's ``occurred`` in the outgoing mapper), so no
    extra Splunk query is required.

    Args:
        service: Splunk service object (used only to read the ES version).
        data: The mirrored incident data from ``UpdateRemoteSystemArgs.data``.

    Returns:
        The finding time string to pass as ``finding_time``, or ``None`` when the
        ES version is out of range or no time is available in the data.
    """
    es_version = get_enterprise_security_version(service)
    if not is_es_version(es_version, "8.2"):
        return None

    finding_time = (data or {}).get("notable_time")
    if not finding_time:
        demisto.debug(
            f"mirror-out: ES version {es_version} is 8.2.x but no 'time' field found in "
            "the mirrored data; proceeding without finding_time."
        )
        return None

    demisto.debug(f"mirror-out: ES version {es_version} is 8.2.x; using finding_time={finding_time}.")
    return str(finding_time)


def update_remote_system_command(
    args: dict[str, Any], params: dict[str, Any], service: client.Service, mapper: UserMappingObject
) -> str:
    """Pushes changes in XSOAR incident into the corresponding Splunk entity (finding or investigation).

    The same code path serves both event types: the remote id ('entity_id') is resolved from
    UpdateRemoteSystemArgs.remote_incident_id (which is the dbotMirrorId — set to either the
    finding's rule_id or the investigation_guid by the corresponding to_incident()).

    Args:
        args (dict): Demisto args
        params (dict): Demisto params
        service (splunklib.client.Service): Splunk service object
        mapper: UserMappingObject for user mapping

    Returns:
        entity_id (str): The remote entity id (finding or investigation).
    """
    parsed_args = UpdateRemoteSystemArgs(args)
    delta = parsed_args.delta
    entity_id = parsed_args.remote_incident_id
    entries = parsed_args.entries
    demisto.debug(f"mirroring args: entries:{parsed_args.entries} delta:{parsed_args.delta} data:{parsed_args.data}")
    # On ES >=8.2 <8.3 the v2 investigations update endpoint requires the finding's
    # original event time (_time). The notable_time is mapped from the the incident's `occured`
    # so it's available in `parsed_args.delta`.
    # return None on ES 8.3+ (or when unavailable), preserving the existing behaviour.
    finding_time = get_finding_time_for_es_notable_time(service, parsed_args.data)
    if parsed_args.incident_changed and delta:
        demisto.debug(f"Got the following delta keys {list(delta.keys())} to update incident corresponding to entity {entity_id}")

        changed_data: dict[str, Any] = {field: None for field in OUTGOING_MIRRORED_FIELDS}
        for field in delta:
            if field == "owner" and params.get("userMapping", False):
                new_owner = mapper.get_splunk_user_by_xsoar(delta["owner"]) if mapper.should_map else None
                if new_owner:
                    changed_data["owner"] = new_owner
                else:
                    demisto.error("New owner was not found while userMapping is enabled.")
            elif field in OUTGOING_MIRRORED_FIELDS:
                changed_data[field] = delta[field]

        # Close entity (finding/investigation) if relevant
        if parsed_args.inc_status == IncidentStatus.DONE and params.get("close_finding"):
            demisto.debug(f"Closing entity {entity_id}")
            changed_data["status"] = "5"

        if any(changed_data.values()):
            demisto.debug(f"Sending update request to Splunk for entity {entity_id}, data: {changed_data}")
            try:
                # Use the v2 API for field updates (handles both findings and investigations).
                demisto.debug(f"Using v2 API to update entity {entity_id}")
                response_info = update_investigation_or_finding(
                    service=service,
                    investigation_or_finding_id=entity_id,
                    owner=changed_data.get("owner"),
                    urgency=changed_data.get("urgency"),
                    status=changed_data.get("status"),
                    disposition=changed_data.get("disposition"),
                    finding_time=finding_time,
                )
                demisto.debug(f"update-remote-system for entity {entity_id} via v2 API: {response_info}")

                # Handle notes separately using the new add_investigation_note function
                if changed_data.get("note"):
                    demisto.debug(f"Adding note to entity {entity_id} via add_investigation_note")
                    try:
                        note_content = f"{changed_data['note']}\n{COMMENT_MIRRORED_FROM_XSOAR}"
                        add_investigation_note(
                            service=service,
                            investigation_or_finding_id=entity_id,
                            content=note_content,
                        )
                        demisto.debug(f"Note added successfully to entity {entity_id}")
                    except Exception as e:
                        demisto.error(f"Failed adding note to entity {entity_id}: {e!s}")

            except Exception as e:
                demisto.error(
                    f"Error in Splunk outgoing mirror for incident corresponding to entity {entity_id}. Error message: {e!s}"
                )
        else:
            demisto.debug(f"Didn't find changed data to update incident corresponding to entity {entity_id}")

    else:
        demisto.debug(f"Incident corresponding to entity {entity_id} was not changed.")

    if entries:
        for entry in entries:
            entry_tags = entry.get("tags", [])
            demisto.debug(f"Got the entry tags: {entry_tags}")
            if NOTE_TAG_TO_SPLUNK in entry_tags:
                demisto.debug("Add new note")
                note_body = f'{entry.get("contents", "")}\n{COMMENT_MIRRORED_FROM_XSOAR}'
                try:
                    add_investigation_note(
                        service=service,
                        investigation_or_finding_id=entity_id,
                        content=note_body,
                    )
                    demisto.debug(f"Note added successfully to entity {entity_id}")
                except Exception as e:
                    demisto.error(
                        f"Error in Splunk outgoing mirror for incident corresponding to entity {entity_id}. "
                        f"Error message: {e!s}"
                    )
    return entity_id


# =========== Mapping Mechanism ===========


def create_mapping_dict(total_parsed_results: list[dict[str, Any]], type_field: str) -> dict[str, Any]:
    """
    Create a {'field_name': 'fields_properties'} dict to be used as mapping schemas.
    Args:
        total_parsed_results: list. the results from the splunk search query
        type_field: str. the field that represents the type of the event or alert.
    """
    types_map = {}
    for result in total_parsed_results:
        raw_json = json.loads(result.get("rawJSON", "{}"))
        if event_type_name := raw_json.get(type_field, ""):
            types_map[event_type_name] = raw_json

    return types_map


def get_mapping_fields_command(service: client.Service, mapper: UserMappingObject, params: dict[str, Any]) -> dict[str, Any]:
    # Create the query to get unique objects
    # The logic is identical to the 'fetch_incidents' command
    type_field = "source"
    total_parsed_results = []

    # Use get_fetch_time_window to calculate the time window
    fetch_window_start_time, fetch_window_end_time = get_fetch_time_window(params, service, "", "")

    kwargs_oneshot = {
        "earliest_time": fetch_window_start_time,
        "latest_time": fetch_window_end_time,
        "count": FETCH_LIMIT,
        "offset": 0,
        "output_mode": OUTPUT_MODE_JSON,
    }

    searchquery_oneshot = params["fetchQuery"]

    if extractFields := params.get("extractFields"):
        for field in extractFields.split(","):
            field_trimmed = field.strip()
            searchquery_oneshot = f"{searchquery_oneshot} | eval {field_trimmed}={field_trimmed}"

    searchquery_oneshot = f"{searchquery_oneshot} | dedup {type_field}"
    oneshotsearch_results = service.jobs.oneshot(searchquery_oneshot, **kwargs_oneshot)
    reader = results.JSONResultsReader(oneshotsearch_results)
    for item in reader:
        if isinstance(item, dict):
            finding = Finding(data=item)
            total_parsed_results.append(finding.to_incident(mapper))
        elif handle_message(item):
            continue

    types_map = create_mapping_dict(total_parsed_results, type_field)
    return types_map


def get_cim_mapping_field_command() -> dict[str, dict[str, Any]]:
    finding = {
        "rule_name": "string",
        "rule_title": "string",
        "security_domain": "string",
        "index": "string",
        "rule_description": "string",
        "risk_score": "string",
        "host": "string",
        "host_risk_object_type": "string",
        "dest_risk_object_type": "string",
        "dest_risk_score": "string",
        "splunk_server": "string",
        "_sourcetype": "string",
        "_indextime": "string",
        "_time": "string",
        "src_risk_object_type": "string",
        "src_risk_score": "string",
        "_raw": "string",
        "urgency": "string",
        "owner": "string",
        "info_min_time": "string",
        "info_max_time": "string",
        "note": "string",
        "reviewer": "string",
        "rule_id": "string",
        "action": "string",
        "app": "string",
        "authentication_method": "string",
        "authentication_service": "string",
        "bugtraq": "string",
        "bytes": "string",
        "bytes_in": "string",
        "bytes_out": "string",
        "category": "string",
        "cert": "string",
        "change": "string",
        "change_type": "string",
        "command": "string",
        "comments": "string",
        "cookie": "string",
        "creation_time": "string",
        "cve": "string",
        "cvss": "string",
        "date": "string",
        "description": "string",
        "dest": "string",
        "dest_bunit": "string",
        "dest_category": "string",
        "dest_dns": "string",
        "dest_interface": "string",
        "dest_ip": "string",
        "dest_ip_range": "string",
        "dest_mac": "string",
        "dest_nt_domain": "string",
        "dest_nt_host": "string",
        "dest_port": "string",
        "dest_priority": "string",
        "dest_translated_ip": "string",
        "dest_translated_port": "string",
        "dest_type": "string",
        "dest_zone": "string",
        "direction": "string",
        "dlp_type": "string",
        "dns": "string",
        "duration": "string",
        "dvc": "string",
        "dvc_bunit": "string",
        "dvc_category": "string",
        "dvc_ip": "string",
        "dvc_mac": "string",
        "dvc_priority": "string",
        "dvc_zone": "string",
        "file_hash": "string",
        "file_name": "string",
        "file_path": "string",
        "file_size": "string",
        "http_content_type": "string",
        "http_method": "string",
        "http_referrer": "string",
        "http_referrer_domain": "string",
        "http_user_agent": "string",
        "icmp_code": "string",
        "icmp_type": "string",
        "id": "string",
        "ids_type": "string",
        "incident": "string",
        "ip": "string",
        "mac": "string",
        "message_id": "string",
        "message_info": "string",
        "message_priority": "string",
        "message_type": "string",
        "mitre_technique_id": "string",
        "msft": "string",
        "mskb": "string",
        "name": "string",
        "orig_dest": "string",
        "orig_recipient": "string",
        "orig_src": "string",
        "os": "string",
        "packets": "string",
        "packets_in": "string",
        "packets_out": "string",
        "parent_process": "string",
        "parent_process_id": "string",
        "parent_process_name": "string",
        "parent_process_path": "string",
        "password": "string",
        "payload": "string",
        "payload_type": "string",
        "priority": "string",
        "problem": "string",
        "process": "string",
        "process_hash": "string",
        "process_id": "string",
        "process_name": "string",
        "process_path": "string",
        "product_version": "string",
        "protocol": "string",
        "protocol_version": "string",
        "query": "string",
        "query_count": "string",
        "query_type": "string",
        "reason": "string",
        "recipient": "string",
        "recipient_count": "string",
        "recipient_domain": "string",
        "recipient_status": "string",
        "record_type": "string",
        "registry_hive": "string",
        "registry_key_name": "string",
        "registry_path": "string",
        "registry_value_data": "string",
        "registry_value_name": "string",
        "registry_value_text": "string",
        "registry_value_type": "string",
        "request_sent_time": "string",
        "request_payload": "string",
        "request_payload_type": "string",
        "response_code": "string",
        "response_payload_type": "string",
        "response_received_time": "string",
        "response_time": "string",
        "result": "string",
        "return_addr": "string",
        "rule": "string",
        "rule_action": "string",
        "sender": "string",
        "service": "string",
        "service_hash": "string",
        "service_id": "string",
        "service_name": "string",
        "service_path": "string",
        "session_id": "string",
        "sessions": "string",
        "severity": "string",
        "severity_id": "string",
        "sid": "string",
        "signature": "string",
        "signature_id": "string",
        "signature_version": "string",
        "site": "string",
        "size": "string",
        "source": "string",
        "sourcetype": "string",
        "src": "string",
        "src_bunit": "string",
        "src_category": "string",
        "src_dns": "string",
        "src_interface": "string",
        "src_ip": "string",
        "src_ip_range": "string",
        "src_mac": "string",
        "src_nt_domain": "string",
        "src_nt_host": "string",
        "src_port": "string",
        "src_priority": "string",
        "src_translated_ip": "string",
        "src_translated_port": "string",
        "src_type": "string",
        "src_user": "string",
        "src_user_bunit": "string",
        "src_user_category": "string",
        "src_user_domain": "string",
        "src_user_id": "string",
        "src_user_priority": "string",
        "src_user_role": "string",
        "src_user_type": "string",
        "src_zone": "string",
        "state": "string",
        "status": "string",
        "status_code": "string",
        "status_description": "string",
        "subject": "string",
        "tag": "string",
        "ticket_id": "string",
        "time": "string",
        "time_submitted": "string",
        "transport": "string",
        "transport_dest_port": "string",
        "type": "string",
        "uri": "string",
        "uri_path": "string",
        "uri_query": "string",
        "url": "string",
        "url_domain": "string",
        "url_length": "string",
        "user": "string",
        "user_agent": "string",
        "user_bunit": "string",
        "user_category": "string",
        "user_id": "string",
        "user_priority": "string",
        "user_role": "string",
        "user_type": "string",
        "vendor_account": "string",
        "vendor_product": "string",
        "vlan": "string",
        "xdelay": "string",
        "xref": "string",
    }

    drilldown = {
        "Drilldown": {
            "action": "string",
            "app": "string",
            "authentication_method": "string",
            "authentication_service": "string",
            "bugtraq": "string",
            "bytes": "string",
            "bytes_in": "string",
            "bytes_out": "string",
            "category": "string",
            "cert": "string",
            "change": "string",
            "change_type": "string",
            "command": "string",
            "comments": "string",
            "cookie": "string",
            "creation_time": "string",
            "cve": "string",
            "cvss": "string",
            "date": "string",
            "description": "string",
            "dest": "string",
            "dest_bunit": "string",
            "dest_category": "string",
            "dest_dns": "string",
            "dest_interface": "string",
            "dest_ip": "string",
            "dest_ip_range": "string",
            "dest_mac": "string",
            "dest_nt_domain": "string",
            "dest_nt_host": "string",
            "dest_port": "string",
            "dest_priority": "string",
            "dest_translated_ip": "string",
            "dest_translated_port": "string",
            "dest_type": "string",
            "dest_zone": "string",
            "direction": "string",
            "dlp_type": "string",
            "dns": "string",
            "duration": "string",
            "dvc": "string",
            "dvc_bunit": "string",
            "dvc_category": "string",
            "dvc_ip": "string",
            "dvc_mac": "string",
            "dvc_priority": "string",
            "dvc_zone": "string",
            "file_hash": "string",
            "file_name": "string",
            "file_path": "string",
            "file_size": "string",
            "http_content_type": "string",
            "http_method": "string",
            "http_referrer": "string",
            "http_referrer_domain": "string",
            "http_user_agent": "string",
            "icmp_code": "string",
            "icmp_type": "string",
            "id": "string",
            "ids_type": "string",
            "incident": "string",
            "ip": "string",
            "mac": "string",
            "message_id": "string",
            "message_info": "string",
            "message_priority": "string",
            "message_type": "string",
            "mitre_technique_id": "string",
            "msft": "string",
            "mskb": "string",
            "name": "string",
            "orig_dest": "string",
            "orig_recipient": "string",
            "orig_src": "string",
            "os": "string",
            "packets": "string",
            "packets_in": "string",
            "packets_out": "string",
            "parent_process": "string",
            "parent_process_id": "string",
            "parent_process_name": "string",
            "parent_process_path": "string",
            "password": "string",
            "payload": "string",
            "payload_type": "string",
            "priority": "string",
            "problem": "string",
            "process": "string",
            "process_hash": "string",
            "process_id": "string",
            "process_name": "string",
            "process_path": "string",
            "product_version": "string",
            "protocol": "string",
            "protocol_version": "string",
            "query": "string",
            "query_count": "string",
            "query_type": "string",
            "reason": "string",
            "recipient": "string",
            "recipient_count": "string",
            "recipient_domain": "string",
            "recipient_status": "string",
            "record_type": "string",
            "registry_hive": "string",
            "registry_key_name": "string",
            "registry_path": "string",
            "registry_value_data": "string",
            "registry_value_name": "string",
            "registry_value_text": "string",
            "registry_value_type": "string",
            "request_payload": "string",
            "request_payload_type": "string",
            "request_sent_time": "string",
            "response_code": "string",
            "response_payload_type": "string",
            "response_received_time": "string",
            "response_time": "string",
            "result": "string",
            "return_addr": "string",
            "rule": "string",
            "rule_action": "string",
            "sender": "string",
            "service": "string",
            "service_hash": "string",
            "service_id": "string",
            "service_name": "string",
            "service_path": "string",
            "session_id": "string",
            "sessions": "string",
            "severity": "string",
            "severity_id": "string",
            "sid": "string",
            "signature": "string",
            "signature_id": "string",
            "signature_version": "string",
            "site": "string",
            "size": "string",
            "source": "string",
            "sourcetype": "string",
            "src": "string",
            "src_bunit": "string",
            "src_category": "string",
            "src_dns": "string",
            "src_interface": "string",
            "src_ip": "string",
            "src_ip_range": "string",
            "src_mac": "string",
            "src_nt_domain": "string",
            "src_nt_host": "string",
            "src_port": "string",
            "src_priority": "string",
            "src_translated_ip": "string",
            "src_translated_port": "string",
            "src_type": "string",
            "src_user": "string",
            "src_user_bunit": "string",
            "src_user_category": "string",
            "src_user_domain": "string",
            "src_user_id": "string",
            "src_user_priority": "string",
            "src_user_role": "string",
            "src_user_type": "string",
            "src_zone": "string",
            "state": "string",
            "status": "string",
            "status_code": "string",
            "subject": "string",
            "tag": "string",
            "ticket_id": "string",
            "time": "string",
            "time_submitted": "string",
            "transport": "string",
            "transport_dest_port": "string",
            "type": "string",
            "uri": "string",
            "uri_path": "string",
            "uri_query": "string",
            "url": "string",
            "url_domain": "string",
            "url_length": "string",
            "user": "string",
            "user_agent": "string",
            "user_bunit": "string",
            "user_category": "string",
            "user_id": "string",
            "user_priority": "string",
            "user_role": "string",
            "user_type": "string",
            "vendor_account": "string",
            "vendor_product": "string",
            "vlan": "string",
            "xdelay": "string",
            "xref": "string",
        }
    }

    asset = {
        "Asset": {
            "asset": "string",
            "asset_id": "string",
            "asset_tag": "string",
            "bunit": "string",
            "category": "string",
            "city": "string",
            "country": "string",
            "dns": "string",
            "ip": "string",
            "is_expected": "string",
            "lat": "string",
            "long": "string",
            "mac": "string",
            "nt_host": "string",
            "owner": "string",
            "pci_domain": "string",
            "priority": "string",
            "requires_av": "string",
        }
    }

    identity = {
        "Identity": {
            "bunit": "string",
            "category": "string",
            "email": "string",
            "endDate": "string",
            "first": "string",
            "identity": "string",
            "identity_tag": "string",
            "last": "string",
            "managedBy": "string",
            "nick": "string",
            "phone": "string",
            "prefix": "string",
            "priority": "string",
            "startDate": "string",
            "suffix": "string",
            "watchlist": "string",
            "work_city": "string",
            "work_lat": "string",
            "work_long": "string",
        }
    }

    return {"Finding Data": finding, "Drilldown Data": drilldown, "Asset Data": asset, "Identity Data": identity}


# =========== Integration Functions & Classes ===========


RESPONSE_SIZE_WARN_THRESHOLD = 20 * 1024 * 1024  # 20 MB
RESPONSE_SIZE_WARN_MESSAGE = (
    "WARNING: Response size ({current_mb:.2f} MB) exceeds the normal usage size of {threshold_mb} MB. "
    "Consider reducing the amount of data returned by your search query. "
    "See the 'Large Search Results' section in the integration documentation for more information."
)


class ResponseSizeValidator:
    """Validates response size and reports warnings only once after all data is collected."""

    def __init__(self):
        self.validated = False

    def validate_and_report(self, data_to_return: list[dict[str, Any]]):
        """Check the size of data that will actually be returned and report warning only once.

        Args:
            data_to_return: The list of results that will be returned to the user
        """
        if self.validated:
            # Already validated and reported, don't do it again
            return

        self.validated = True

        # Calculate the actual size of the data that will be returned
        data_json = json.dumps(data_to_return)
        actual_size = len(data_json.encode("utf-8"))
        if actual_size > RESPONSE_SIZE_WARN_THRESHOLD:
            return_results(
                RESPONSE_SIZE_WARN_MESSAGE.format(
                    current_mb=actual_size / (1024 * 1024), threshold_mb=RESPONSE_SIZE_WARN_THRESHOLD / (1024 * 1024)
                )
            )


class ResponseReaderWrapper(io.RawIOBase):
    """This class was supplied as a solution for a bug in Splunk causing the search to run slowly."""

    def __init__(self, response_reader):
        self.response_reader = response_reader

    def readable(self) -> bool:
        return True

    def close(self) -> None:
        self.response_reader.close()

    def read(self, n: int) -> bytes:  # type: ignore[override]
        return self.response_reader.read(n)

    def readinto(self, b: bytearray) -> int:  # type: ignore[override]
        sz = len(b)
        data = self.response_reader.read(sz)

        # Remove non utf-8 characters to avoid decode errors in JSONResultsReader
        # See resolution section from: https://splunk.my.site.com/customer/s/article/Search-Failed-Due-to
        cleaned_data = data.decode("utf-8", errors="ignore").encode("utf-8")
        if len(cleaned_data) != len(data):  # Check if any bytes were removed
            demisto.debug(
                "Removed non utf-8 characters in incoming Splunk data:\n"
                f"Original Splunk data: {data}\n"
                f"Modified data: {cleaned_data}\n"
            )

        for idx, ch in enumerate(cleaned_data):
            b[idx] = ch

        return len(cleaned_data)


def add_investigation_note(
    service: client.Service,
    investigation_or_finding_id: str,
    content: str,
    note_type: str | None = None,
    finding_time: str | None = None,
):
    """Add a note to a Splunk investigation or finding via the v2 investigations API endpoint.

    Args:
        service: Splunk service connection
        investigation_or_finding_id: The ID of the investigation or finding
        content: The content of the note
        note_type: Optional type of the note (e.g., "Task")
        finding_time (str | None): The time associated with the finding event. When provided,
            used as the notable_time parameter on the first API call. If not provided, the first
            call is made without notable_time and falls back to notable_time="now" on failure.

    Returns:
        dict: The JSON response from the API
    """
    body = {"content": content}
    if note_type is not None:
        body["type"] = note_type

    endpoint = f"public/v2/investigations/{investigation_or_finding_id}/notes"

    demisto.debug(f"Adding note to investigation/finding {investigation_or_finding_id}")

    # Build optional kwargs for the first call: include notable_time only when finding_time is provided
    first_call_kwargs: dict[str, str] = {}
    if finding_time is not None:
        first_call_kwargs["notable_time"] = finding_time

    try:
        response = service.post(endpoint, body=json.dumps(body), **first_call_kwargs)
    except Exception as e:
        demisto.debug(
            f"Failed to add note to investigation/finding {investigation_or_finding_id} "
            f"{'with notable_time=' + finding_time if finding_time else 'without notable_time param'}, "
            f"retrying with notable_time=now. Error: {e!s}"
        )
        response = service.post(endpoint, body=json.dumps(body), notable_time="now")

    response_data = response.body.read()
    result = json.loads(response_data)
    demisto.debug(f"Note added successfully: {result}")
    return result


def update_investigation_or_finding(
    service: client.Service,
    investigation_or_finding_id: str,
    owner: str | None = None,
    urgency: str | None = None,
    status: str | None = None,
    disposition: str | None = None,
    finding_time: str | None = None,
    name: str | None = None,
    description: str | None = None,
):
    """
    Update a Splunk investigation or finding via the v2 investigations API endpoint.

    This function uses the service object to make a POST request to the
    /public/v2/investigations/:id endpoint with notable_time=now parameter.

    Args:
        service (client.Service): Splunk service object (already connected)
        investigation_or_finding_id (str): The ID of the investigation or finding to update
        owner (str | None): New owner for the investigation/finding
        urgency (str | None): New urgency level
        status (str | None): New status
        disposition (str | None): New disposition
        finding_time (str | None): The time associated with the finding event. When provided,
            used as the notable_time parameter on the first API call. If not provided, the first
            call is made without notable_time and falls back to notable_time="now" on failure.
        name (str | None): New name for the investigation (investigations only).
        description (str | None): New description for the investigation (investigations only).

    Returns:
        dict: The JSON response from the API

    Raises:
        Exception: If the API request fails
    """
    # Build the request body with only the fields that are provided
    body = {}
    if owner is not None:
        body["owner"] = owner
    if urgency is not None:
        body["urgency"] = urgency
    if status is not None:
        body["status"] = status
    if disposition is not None:
        body["disposition"] = disposition
    if name is not None:
        body["name"] = name
    if description is not None:
        body["description"] = description

    # If no fields to update, return early
    if not body:
        demisto.debug(f"No fields to update for investigation/finding {investigation_or_finding_id}")
        return {"success": False, "message": "No fields provided to update"}

    endpoint = f"public/v2/investigations/{investigation_or_finding_id}"

    demisto.debug(
        f"Updating investigation/finding {investigation_or_finding_id} via v2 API. " f"Endpoint: {endpoint}, Body: {body}"
    )

    # Build optional kwargs for the first call: include notable_time only when finding_time is provided
    first_call_kwargs: dict[str, str] = {}
    if finding_time is not None:
        first_call_kwargs["notable_time"] = finding_time

    try:
        response = service.post(endpoint, body=json.dumps(body), **first_call_kwargs)
    except Exception as e:
        demisto.debug(
            f"Failed to update investigation/finding {investigation_or_finding_id} "
            f"{'with notable_time=' + finding_time if finding_time else 'without notable_time param'}, "
            f"retrying with notable_time=now. Error: {e!s}"
        )
        response = service.post(endpoint, body=json.dumps(body), notable_time="now")

    response_data = response.body.read()
    result = json.loads(response_data)
    demisto.debug(f"Successfully updated investigation/finding {investigation_or_finding_id}: {result}")
    return result


# --------------------------------------------------------------------------- #
# Splunk ES (Mission Control) REST helpers and investigation commands         #
# --------------------------------------------------------------------------- #


def _es_rest_request(
    service: client.Service,
    method: str,
    path: str,
    body: dict | None = None,
    query: dict | None = None,
) -> dict | list:
    """Perform a request against the Splunk ES (Mission Control) public REST API.

    Reuses the connected ``splunklib`` ``Service`` object's ``post``/``get`` methods so the
    same authentication, session and namespace as the rest of the integration are used.

    Args:
        service: A connected ``splunklib.client.Service`` instance.
        method: HTTP verb. One of ``GET`` or ``POST``.
        path: Path under the ES public namespace (e.g. ``public/v2/investigations``).
        body: Optional JSON-serializable body for POST requests.
        query: Optional query string parameters for GET requests.

    Returns:
        Parsed JSON response (dict or list). Empty responses are returned as an empty dict.

    Raises:
        DemistoException: On HTTP errors or non-JSON responses.
    """
    # Ensure the service points at the Mission Control namespace used by all ES API calls.
    service.namespace = namespace(app="missioncontrol", owner="nobody")
    method_upper = method.upper()
    demisto.debug(f"ES REST request: {method_upper} {path} body={body} query={query}")
    try:
        if method_upper == "GET":
            response = service.get(path, **(query or {}))
        elif method_upper == "POST":
            response = service.post(path, body=json.dumps(body or {}))
        else:
            raise DemistoException(f"Unsupported HTTP method for ES REST helper: {method}")
    except HTTPError as exc:
        raise DemistoException(f"Splunk ES REST request failed ({method_upper} {path}): {exc!s}") from exc

    raw = response.body.read() if hasattr(response, "body") else b""
    if not raw:
        return {}
    try:
        return json.loads(raw)
    except ValueError:
        # Some endpoints return a plain string body (e.g. confirmation messages).
        return {"raw_response": raw.decode("utf-8", errors="replace")}


# Fields forwarded as-is in the POST /public/v2/investigations create payload.
INVESTIGATION_CREATE_FIELDS = (
    "name",
    "description",
    "investigation_type",
    "status",
    "disposition",
    "owner",
    "urgency",
    "sensitivity",
)


def _build_investigation_create_payload(args: dict) -> dict:
    """Build the JSON payload for ``POST /public/v2/investigations`` from command args.

    Forwards only the fields in ``INVESTIGATION_CREATE_FIELDS``: ``name``, ``description``,
    ``investigation_type``, ``status``, ``disposition``, ``owner``, ``urgency``, ``sensitivity``.
    ``name`` is required.
    """
    payload: dict[str, Any] = {}
    for key in INVESTIGATION_CREATE_FIELDS:
        value = args.get(key)
        if value:
            payload[key] = value
    if "name" not in payload:
        raise DemistoException("`name` is a required argument for splunk-investigation-create.")
    return payload


def _format_investigation_for_hr(item: dict) -> dict:
    """Build a curated, human-readable projection of an investigation dict.

    Only fields useful in the war-room table are returned, with friendly Title Case
    keys; epoch-seconds timestamps (``create_time`` / ``update_time``) are converted
    to ISO-8601 UTC via :func:`timestamp_to_datestring`. Raw ``outputs`` are not
    affected by this projection.
    """

    def _to_iso(value: Any) -> str | None:
        if value in (None, "", 0):
            return None
        try:
            return timestamp_to_datestring(int(float(value)) * 1000, date_format="%Y-%m-%dT%H:%M:%SZ", is_utc=True)
        except (TypeError, ValueError) as exc:
            demisto.debug(f"_format_investigation_for_hr: cannot format {value!r}: {exc!s}")
            return None

    return {
        "Investigation ID": item.get("investigation_id"),
        "Investigation GUID": item.get("investigation_guid"),
        "Name": item.get("name"),
        "Status": item.get("status_name"),
        "Disposition": item.get("disposition_name"),
        "Owner": item.get("owner"),
        "Urgency": item.get("urgency"),
        "Sensitivity": item.get("sensitivity"),
        "Findings Count": item.get("count_findings"),
        "Risk Score": item.get("risk_score"),
        "Create Time": _to_iso(item.get("create_time")),
        "Update Time": _to_iso(item.get("update_time")),
    }


# Ordered headers for the splunk-investigation-list HR table.
INVESTIGATION_HR_HEADERS = [
    "Investigation ID",
    "Name",
    "Status",
    "Disposition",
    "Owner",
    "Urgency",
    "Update Time",
]


def splunk_create_investigation_command(service: client.Service, args: dict) -> CommandResults:
    """Create a new investigation in Splunk ES via ``POST /public/v2/investigations``.

    Returns the new investigation's GUID in the ``Splunk.Investigation`` context.
    """
    payload = _build_investigation_create_payload(args)
    response = _es_rest_request(service, "POST", "public/v2/investigations", body=payload)

    response_dict = response if isinstance(response, dict) else {}
    investigation_guid = response_dict.get("investigation_guid")
    readable = tableToMarkdown(
        "Investigation created successfully",
        {"Investigation GUID": investigation_guid},
        headers=["Investigation GUID"],
        removeNull=True,
    )
    return CommandResults(
        outputs_prefix="Splunk.Investigation",
        outputs_key_field="investigation_guid",
        outputs=response_dict or None,
        readable_output=readable,
        raw_response=response,
    )


# Scalar query parameters accepted by GET /public/v2/investigations.
INVESTIGATION_LIST_SCALAR_PARAMS = (
    "limit",
    "offset",
    "sort",
    "create_time_min",
    "create_time_max",
    "update_time_min",
    "update_time_max",
)

# Multi-value (isArray: true) query parameters for GET /public/v2/investigations.
# Each is parsed with ``argToList()`` and re-serialized as a CSV string for the API.
INVESTIGATION_LIST_LIST_PARAMS = (
    "ids",
    "disposition",
    "status",
    "owner",
    "urgency",
    "sensitivity",
)


def splunk_list_investigations_command(service: client.Service, args: dict) -> CommandResults:
    """List investigations via ``GET /public/v2/investigations``.

    This command supports the full set of list-investigations query filters
    (``investigation_ids``, ``limit``, ``offset``, ``sort``, ``disposition``, ``status``,
    ``owner``, ``urgency``, ``sensitivity``, ``create_time_min``/``max``,
    ``update_time_min``/``max``).
    Multi-value arguments (``investigation_ids``, ``disposition``, ``status``, ``owner``,
    ``urgency``, ``sensitivity``) are parsed with ``argToList()`` and forwarded to the
    matching Splunk query parameter as a CSV string (e.g., ``ids=id1,id2``,
    ``urgency=high,critical``). Investigation IDs accept either the GUID or the display
    ID (``ES-00001``).
    """

    if "investigation_ids" in args:
        args["ids"] = args.get("investigation_ids")

    query: dict[str, Any] = {key: args.get(key) for key in INVESTIGATION_LIST_SCALAR_PARAMS}
    for key in INVESTIGATION_LIST_LIST_PARAMS:
        values = argToList(args.get(key))
        if values:
            query[key] = ",".join(str(v) for v in values)
    query = assign_params(**query)
    demisto.debug(f"splunk-investigation-list: GET public/v2/investigations query={query}")
    response = _es_rest_request(service, "GET", "public/v2/investigations", query=query)

    # Normalize the response into a list of investigation dicts.
    investigations: list[dict] = []
    if isinstance(response, list):
        investigations = [item for item in response if isinstance(item, dict)]
    elif isinstance(response, dict):
        for value in response.values():
            if isinstance(value, list):
                investigations = [item for item in value if isinstance(item, dict)]
                break
        if not investigations and ("investigation_guid" in response or "investigation_id" in response):
            investigations = [response]

    if not investigations:
        return CommandResults(readable_output="No investigations found for the provided filters.")

    hr_rows = [_format_investigation_for_hr(item) for item in investigations]
    ids_list = argToList(args.get("ids"))
    count = len(investigations)
    noun = "Investigation" if count == 1 else "Investigations"
    if ids_list:
        title = f"Splunk ES — {count} {noun} Found ({', '.join(ids_list)})"
    else:
        title = f"Splunk ES — {count} {noun} Found"
    readable = tableToMarkdown(
        title,
        hr_rows,
        headers=INVESTIGATION_HR_HEADERS,
        removeNull=True,
    )
    outputs: dict | list = investigations[0] if len(investigations) == 1 else investigations
    return CommandResults(
        outputs_prefix="Splunk.Investigation",
        outputs_key_field="investigation_guid",
        outputs=outputs,
        readable_output=readable,
        raw_response=response,
    )


def _fetch_modified_investigations_page(
    service: client.Service, update_time_min: str, limit: int, offset: int
) -> list[dict[str, Any]]:
    """Hit the investigations endpoint once and return the rows of a single page.

    Splits the HTTP round trip out of :func:`list_modified_investigations` so the
    paginating loop stays small and the helper itself is easy to mock in tests.
    """
    endpoint = "public/v2/investigations"
    params = assign_params(update_time_min=update_time_min, limit=limit, offset=offset)
    demisto.debug(f"list_modified_investigations: hitting endpoint={endpoint}?{urlencode(params)}")
    response = service.get("public/v2/investigations", app="missioncontrol", owner="nobody", **params)
    payload = json.loads(response.body.read())
    return payload or []


def list_modified_investigations(
    service: client.Service,
    update_time_min: str,
    page_size: int = INVESTIGATIONS_MAX_LIMIT,
    max_total: int = MIRROR_LIMIT,
    mapper: "UserMappingObject | None" = None,
) -> list[dict[str, Any]]:
    """Return investigation rows updated since ``update_time_min``.

    Hits the investigations endpoint directly (no SPL) and paginates with ``offset`` until
    either a short page is returned (no more rows) or ``max_total`` is reached
    — mirroring the cap Findings already enforces via :data:`MIRROR_LIMIT`.

    Errors are logged via :func:`demisto.error` and surface as an empty list so
    a single endpoint hiccup does not break mirror-in for Findings.

    Args:
        service: Splunk service connection.
        update_time_min: ISO 8601 UTC timestamp (canonical "Z" shape produced by
            :func:`to_mc_iso8601_utc`).
        page_size: Per-call cap; coerced to int and clamped to
            :data:`INVESTIGATIONS_MAX_LIMIT` (100, the v2 endpoint's hard cap).
        max_total: Aggregate cap across all pages in this invocation. Defaults to
            :data:`MIRROR_LIMIT` (1000) so investigations and findings share the
            same per-tick ceiling.
        mapper: Optional :class:`UserMappingObject` used to translate the Splunk
            ``owner`` value on each row to the matching XSOAR user. When
            ``None`` (or when the instance is not configured for user
            mapping), rows are returned with their raw Splunk ``owner``
            untouched — preserving the previous behaviour for callers that do
            not pass a mapper.

    Returns:
        A list of parsed investigation row ``dict`` objects (deduped by
        ``investigation_guid``/``investigation_id``). Each row is normalised
        through :func:`parse_investigation` so the shape matches what the
        original fetch flow produces (collapsed dotted keys, lifted
        ``incident_ids``, JSON-serialised ``consolidated_findings``,
        ``splunk_es_event_type`` tag) — keeping mirror-in payloads aligned
        with classifier/mapper expectations. Rows missing both id fields are
        silently skipped (mirror-in needs a stable id to route the update).
        Returns an empty list on any exception.
    """
    try:
        clamped_page_size = min(int(page_size), INVESTIGATIONS_MAX_LIMIT)
        collected_investigations: list[dict[str, Any]] = []
        seen_investigation_ids: set[str] = set()
        offset = 0
        while len(collected_investigations) < max_total:
            page = _fetch_modified_investigations_page(service, update_time_min, clamped_page_size, offset)
            for row in page:
                stable_id = row.get("investigation_guid") or row.get("investigation_id")
                if not stable_id:
                    continue
                if stable_id in seen_investigation_ids:
                    continue
                seen_investigation_ids.add(stable_id)
                collected_investigations.append(row)
                if len(collected_investigations) >= max_total:
                    break
            # Short page → no more rows; advance otherwise.
            if len(page) < clamped_page_size:
                break
            offset += clamped_page_size

        # Normalise each row through `parse_investigation` so every caller
        # (mirror-in today, any future caller) sees the same canonical shape
        # the original fetch flow emits — collapsed dotted keys, lifted
        # `incident_ids`, JSON-serialised `consolidated_findings`, and the
        # `splunk_es_event_type` tag. Done before owner mapping so the mapper
        # sees the same row shape it would during fetch.
        collected_investigations = [parse_investigation(row) for row in collected_investigations]

        # Owner mapping (Splunk → XSOAR), parallel to the Findings flow. The
        # helper is a no-op when the mapper is missing or `should_map` is False,
        # so this is safe to call unconditionally.
        if mapper is not None and collected_investigations:
            mapper.map_owner_to_xsoar_user(collected_investigations)

        demisto.debug(
            f"list_modified_investigations: collected {len(collected_investigations)} rows across "
            f"offset cursor (page_size={clamped_page_size}, max_total={max_total})"
        )
        return collected_investigations
    except Exception as exc:
        demisto.error(f"list_modified_investigations: failed to query v2 endpoint: {exc}\n{traceback.format_exc()}")
        return []


def severity_to_level(severity: str | None) -> int | float:
    match severity:
        case "informational":
            return 0.5
        case "critical":
            return 4
        case "high":
            return 3
        case "medium":
            return 2
        case _:
            return 1


def parse_finding(finding: dict[str, Any], to_dict: bool = False) -> dict[str, Any]:
    """Parses the finding

    Args:
        finding (OrderedDict): The finding
        to_dict (bool): Whether to cast the finding to dict or not.

    Returns (OrderedDict or dict): The parsed finding
    """
    finding = replace_keys(finding) if REPLACE_FLAG else finding
    for key, val in list(finding.items()):
        # if finding event raw fields were sent in double quotes (e.g. "DNS Destination") and the field does not exist
        # in the event, then splunk returns the field with the key as value (e.g. ("DNS Destination", "DNS Destination")
        # so we go over the fields, and check if the key equals the value and set the value to be empty string
        if key == val:
            demisto.debug(
                f"Found finding event raw field [{key}] with key that equals the value - replacing the value with empty string"
            )
            finding[key] = ""
    return dict(finding) if to_dict else finding


def requests_handler(url: str, message: dict[str, Any], **kwargs: Any) -> dict[str, Any]:
    method = message["method"].lower()
    data = message.get("body", "") if method == "post" else None
    headers = dict(message.get("headers", []))
    try:
        response = requests.request(method, url, data=data, headers=headers, verify=VERIFY_CERTIFICATE, **kwargs)
    except requests.exceptions.HTTPError as e:
        # Propagate HTTP errors via the returned response message
        response = e.response
        demisto.debug(f"Got exception while using requests handler - {e!s}")
    return {
        "status": response.status_code,
        "reason": response.reason,
        "headers": list(response.headers.items()),
        "body": io.BytesIO(response.content),
    }


def build_search_kwargs(args: dict[str, Any], polling: bool = False) -> dict[str, Any]:
    t = datetime.now(pytz.UTC) - timedelta(days=7)
    time_str = t.strftime(ISO_FORMAT_TZ_AWARE)

    kwargs_normal_search: dict[str, Any] = {
        "earliest_time": time_str,
    }
    if demisto.get(args, "earliest_time"):
        kwargs_normal_search["earliest_time"] = args["earliest_time"]
    if demisto.get(args, "latest_time"):
        kwargs_normal_search["latest_time"] = args["latest_time"]
    if demisto.get(args, "app"):
        kwargs_normal_search["app"] = args["app"]
    if argToBoolean(demisto.get(args, "fast_mode")):
        kwargs_normal_search["adhoc_search_level"] = "fast"
    kwargs_normal_search["exec_mode"] = "normal" if polling else "blocking"
    return kwargs_normal_search


def build_search_query(args: dict[str, Any]) -> str:
    query = args["query"]
    if not query.startswith("search") and not query.startswith("Search") and not query.startswith("|"):
        query = f"search {query}"
    return query


def create_entry_context(
    args: dict[str, Any],
    parsed_search_results: list[dict[str, Any]],
    dbot_scores: list[dict[str, Any]],
    status_res: CommandResults | None,
    job_id: str | None,
) -> tuple[dict[str, Any], dict[str, Any]]:
    ec = {}
    dbot_ec = {}
    number_of_results = len(parsed_search_results)

    if args.get("update_context", "true") == "true":
        ec["Splunk.Result"] = parsed_search_results
        if len(dbot_scores) > 0:
            dbot_ec["DBotScore"] = dbot_scores
        if status_res:
            ec["Splunk.JobStatus(val.SID && val.SID === obj.SID)"] = {**status_res.outputs, "TotalResults": number_of_results}  # type: ignore[dict-item, assignment]
    if job_id and not status_res:
        status = "DONE" if (number_of_results > 0) else "NO RESULTS"
        ec["Splunk.JobStatus(val.SID && val.SID === obj.SID)"] = [
            {"SID": job_id, "TotalResults": number_of_results, "Status": status}
        ]
    return ec, dbot_ec


def schedule_polling_command(command: str, args: dict[str, Any], interval_in_secs: int) -> ScheduledCommand:
    """
    Returns a ScheduledCommand object which contain the needed arguments for schedule the polling command.
    """
    return ScheduledCommand(command=command, next_run_in_seconds=interval_in_secs, args=args, timeout_in_seconds=600)


def build_search_human_readable(args: dict[str, Any], parsed_search_results: list[dict[str, Any]], sid: str | None) -> str:
    headers: str | list[str] = ""
    if parsed_search_results and len(parsed_search_results) > 0:
        if not isinstance(parsed_search_results[0], dict):
            headers = "results"
        else:
            query = args.get("query", "")
            table_args = re.findall(" table (?P<table>[^|]*)", query)
            rename_args = re.findall(" rename (?P<rename>[^|]*)", query)

            chosen_fields: list = []
            for arg_string in table_args:
                chosen_fields.extend(field.strip('"') for field in re.findall(r'((?:".*?")|(?:[^\s,]+))', arg_string) if field)
            rename_dict = {}
            for arg_string in rename_args:
                for field in re.findall(r'((?:".*?")|(?:[^\s,]+))( AS )((?:".*?")|(?:[^\s,]+))', arg_string):
                    if field:
                        rename_dict[field[0].strip('"')] = field[-1].strip('"')

            # replace renamed fields
            chosen_fields = [rename_dict.get(field, field) for field in chosen_fields]

            headers = update_headers_from_field_names(parsed_search_results, chosen_fields)

    query = args["query"].replace("`", r"\`")
    hr_headline = "Splunk Search results for query:\n"
    if sid:
        hr_headline += f"sid: {sid!s}"
    return tableToMarkdown(hr_headline, parsed_search_results, headers)


def update_headers_from_field_names(search_result: list[dict[str, Any]], chosen_fields: list[str]) -> list[str]:
    headers: list = []
    search_result_keys: set = set().union(*(list(d.keys()) for d in search_result))
    for field_name in chosen_fields:
        if field_name[-1] == "*":
            temp_field = field_name.replace("*", ".*")
            headers.extend(key for key in search_result_keys if re.search(temp_field, key))
        elif field_name in search_result_keys:
            headers.append(field_name)

    return headers


def get_current_results_batch(search_job: client.Job, batch_size: int, results_offset: int) -> Any:
    current_batch_kwargs = {
        "count": batch_size,
        "offset": results_offset,
        "output_mode": OUTPUT_MODE_JSON,
    }

    return search_job.results(**current_batch_kwargs)


def parse_batch_of_results(
    current_batch_of_results: Any, max_results_to_add: float, app: str
) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]:
    parsed_batch_results = []
    batch_dbot_scores = []
    results_reader = results.JSONResultsReader(io.BufferedReader(ResponseReaderWrapper(current_batch_of_results)))
    for item in results_reader:
        if handle_message(item):
            continue

        elif isinstance(item, dict):
            if demisto.get(item, "host"):
                batch_dbot_scores.append(
                    {"Indicator": item["host"], "Type": "hostname", "Vendor": "Splunk", "Score": 0, "isTypedIndicator": True}
                )
            if app:
                item["app"] = app
            # Normal events are returned as dicts
            parsed_batch_results.append(item)

        if len(parsed_batch_results) >= max_results_to_add:
            break
    return parsed_batch_results, batch_dbot_scores


def raise_error_for_failed_job(job: client.Job | None) -> None:
    """
    Handle the case that the search job failed due to dome reason like parsing issues etc
    raise DemistoException in case there is a fatal error in the search job.
    see https://docs.splunk.com/Documentation/Splunk/9.3.0/RESTTUT/RESTsearches#:~:text=the%20results%20returned.-,dispatchState,-dispatchState%20is%20one

    Args:
        job (Job): the created search job

    Raises:
        Exception: DemistoException in case there is a fatal error
    """
    err_msg = None
    try:
        if job and job["dispatchState"] == "FAILED":
            messages = job["messages"]
            for err_type in ["fatal", "error"]:
                if messages.get(err_type):
                    err_msg = ",".join(messages[err_type])
                    break
    except Exception:
        pass
    if err_msg:
        raise DemistoException(f"Failed to run the search in Splunk: {err_msg}")


def splunk_search_command(service: client.Service, args: dict[str, Any]) -> CommandResults | list[CommandResults]:
    query = build_search_query(args)
    polling = argToBoolean(args.get("polling", False))
    search_kwargs = build_search_kwargs(args, polling)
    job_sid = args.get("sid")
    search_job = None
    interval_in_secs = int(args.get("interval_in_seconds", 30))
    if not job_sid or not polling:
        # create a new job to search the query.
        search_job = service.jobs.create(query, **search_kwargs)
        job_sid = search_job["sid"]
        args["sid"] = job_sid
        raise_error_for_failed_job(search_job)

    status_cmd_result: CommandResults | None = None
    if polling:
        status_cmd_results = splunk_job_status(service, args)
        assert status_cmd_results  # if polling is true, status_cmd_result should not be an empty list
        status_cmd_result = status_cmd_results[0]
        status = status_cmd_result.outputs["Status"]  # type: ignore[index]
        if status.lower() != "done":
            # Job is still running, schedule the next run of the command.
            scheduled_command = schedule_polling_command("splunk-search", args, interval_in_secs)
            status_cmd_result.scheduled_command = scheduled_command
            status_cmd_result.readable_output = "Job is still running, it may take a little while..."
            return status_cmd_result
        else:
            # Get the job by its SID.
            search_job = service.job(job_sid)
    num_of_results_from_query = search_job["resultCount"] if search_job else None

    results_limit = float(args.get("event_limit", 100))
    if results_limit == 0.0:
        # In Splunk, a result limit of 0 means no limit.
        results_limit = float("inf")
    batch_size = int(args.get("batch_limit", 25000))

    results_offset = 0
    total_parsed_results: list[dict[str, Any]] = []
    dbot_scores: list[dict[str, Any]] = []

    while (
        len(total_parsed_results) < int(num_of_results_from_query)  # type: ignore[arg-type]
        and len(total_parsed_results) < results_limit
    ):
        current_batch_of_results = get_current_results_batch(search_job, batch_size, results_offset)
        max_results_to_add = results_limit - len(total_parsed_results)
        parsed_batch_results, batch_dbot_scores = parse_batch_of_results(
            current_batch_of_results, max_results_to_add, search_kwargs.get("app", "")
        )
        total_parsed_results.extend(parsed_batch_results)
        dbot_scores.extend(batch_dbot_scores)

        results_offset += batch_size

    # Validate response size only on the data that will actually be returned
    size_validator = ResponseSizeValidator()
    size_validator.validate_and_report(total_parsed_results)

    entry_context_splunk_search, entry_context_dbot_score = create_entry_context(
        args, total_parsed_results, dbot_scores, status_cmd_result, str(job_sid)
    )
    human_readable = build_search_human_readable(args, total_parsed_results, str(job_sid))
    results = [
        CommandResults(outputs=entry_context_splunk_search, raw_response=total_parsed_results, readable_output=human_readable)
    ]
    dbot_table_headers = ["Indicator", "Type", "Vendor", "Score", "isTypedIndicator"]
    if entry_context_dbot_score:
        results.append(
            CommandResults(
                outputs=entry_context_dbot_score,
                readable_output=tableToMarkdown("DBot Score", entry_context_dbot_score["DBotScore"], headers=dbot_table_headers),
            )
        )
    return results


def splunk_job_create_command(service: client.Service, args: dict[str, Any]) -> None:
    app = args.get("app", "")
    query = build_search_query(args)
    search_kwargs = {"exec_mode": "normal", "app": app}
    search_job = service.jobs.create(query, **search_kwargs)

    return_results(
        CommandResults(
            outputs_prefix="Splunk",
            readable_output=f"Splunk Job created with SID: {search_job.sid}",
            outputs={"Job": search_job.sid},
        )
    )


def splunk_results_command(service: client.Service, args: dict[str, Any]) -> str | None:
    res = []
    sid = args.get("sid", "")
    limit = int(args.get("limit", "100"))
    try:
        job = service.job(sid)
    except HTTPError as error:
        msg = error.message if hasattr(error, "message") else str(error)
        if error.status == 404:
            return f"Found no job for sid: {sid}"
        else:
            return_error(msg, error)
    else:
        for result in results.JSONResultsReader(job.results(count=limit, output_mode=OUTPUT_MODE_JSON)):
            if isinstance(result, results.Message):
                res.append({"Splunk message": json.dumps(result.message)})
            elif isinstance(result, dict):
                # Normal events are returned as dicts
                res.append(result)
        return_results(
            CommandResults(
                raw_response=json.dumps(res),
                content_format=EntryFormat.JSON,
            )
        )
    return None


def splunk_get_indexes_command(service: client.Service, app: str = "-"):
    search_query = f"""| rest "/servicesNS/nobody/{app}/data/indexes/?count=-1&offset=0"
    | eval name=title, count=totalEventCount
    | table name, count"""

    indexesNames = []

    # Try the first approach: REST API query
    try:
        demisto.debug("Attempting to get indexes using REST API query approach")
        for item in results.JSONResultsReader(service.jobs.oneshot(query=search_query, output_mode=OUTPUT_MODE_JSON)):
            if handle_message(item):
                continue
            indexesNames.append(item)
        demisto.debug(f"Successfully retrieved {len(indexesNames)} indexes using REST API query approach")
    except Exception as e:
        # Log the error and fall back to the second approach
        demisto.error(f"Failed to get indexes using REST API query approach: {e!s}")
        demisto.debug("Falling back to direct API approach using service.indexes")

        try:
            # Second approach: Direct API using service.indexes
            indexes = service.indexes
            for index in indexes:
                index_json = {"name": index.name, "count": index["totalEventCount"]}
                indexesNames.append(index_json)
            demisto.debug(f"Successfully retrieved {len(indexesNames)} indexes using direct API approach")
        except Exception as fallback_error:
            demisto.error(f"Failed to get indexes using direct API approach: {fallback_error!s}")
            raise DemistoException(
                f"Failed to retrieve indexes using both methods. " f"REST API error: {e!s}. Direct API error: {fallback_error!s}"
            )

    return_results(
        CommandResults(
            content_format=EntryFormat.JSON,
            raw_response=json.dumps(indexesNames),
            readable_output=tableToMarkdown("Splunk Indexes names", indexesNames, ""),
        )
    )


def splunk_submit_event_command(service: client.Service, args: dict[str, Any]) -> None:
    try:
        index = service.indexes[args["index"]]
    except KeyError:
        return_error(f'Found no Splunk index: {args["index"]}')

    else:
        data = args["data"]
        data_formatted = data.encode("utf8")
        r = index.submit(data_formatted, sourcetype=args["sourcetype"], host=args["host"])
        return_results(f"Event was created in Splunk index: {r.name}")


def get_events_from_file(entry_id: str) -> str:
    """
    Retrieves event data from a file in Demisto based on a specified entry ID as a string.

    Args:
        entry_id (int): The entry ID corresponding to the file containing event data.

    Returns:
        str: The content of the file as a string.
    """
    get_file_path_res = demisto.getFilePath(entry_id)
    file_path = get_file_path_res["path"]
    with open(file_path, encoding="utf-8") as file_data:
        return file_data.read()


def parse_fields(fields: str | None) -> dict[str, Any] | None:
    """
    Parses the `fields` input into a dictionary.

    - If `fields` is a valid JSON string, it is converted into the corresponding dictionary.
    - If `fields` is not valid JSON, it is wrapped as a dictionary with a single key-value pair,
    where the key is `"fields"` and the value is the original `fields` string.

    Examples:
    1. Input: '{"severity": "INFO", "category": "test2, test2"}'
       Output: {"severity": "INFO", "category": "test2, test2"}

    2. Input: 'severity: INFO, category: test2, test2'
       Output: {"fields": "severity: INFO, category: test2, test2"}
    """
    if fields:
        try:
            parsed_fields = json.loads(fields)
        except Exception:
            demisto.debug("Fields provided are not valid JSON; treating as a single field")
            parsed_fields = {"fields": fields}
        return parsed_fields
    return None


def splunk_submit_event_hec(
    hec_token: str | None,
    baseurl: str,
    event: str | None,
    fields: str | None,
    host: str | None,
    index: str | None,
    source_type: str | None,
    source: str | None,
    time_: str | None,
    request_channel: str | None,
    batch_event_data: str | None,
    entry_id: str | None,
    service: client.Service,
) -> requests.Response:
    if hec_token is None:
        raise Exception("The HEC Token was not provided")

    if batch_event_data:
        events = batch_event_data

    elif entry_id:
        demisto.debug(f"{INTEGRATION_LOG} - loading events data from file with {entry_id=}")
        events = get_events_from_file(entry_id)

    else:
        parsed_fields = parse_fields(fields)

        events = assign_params(
            event=event, host=host, fields=parsed_fields, index=index, sourcetype=source_type, source=source, time=time_
        )

    headers = {
        "Authorization": f"Splunk {hec_token}",
        "Content-Type": "application/json",
    }
    if request_channel:
        headers["X-Splunk-Request-Channel"] = request_channel

    data = ""
    if entry_id or batch_event_data:
        data = events
    else:
        data = json.dumps(events)

    return requests.post(
        f"{baseurl}/services/collector/event",
        data=data,
        headers=headers,
        verify=VERIFY_CERTIFICATE,
    )


def splunk_submit_event_hec_command(params: dict[str, Any], service: client.Service, args: dict[str, Any]) -> None:
    hec_token = params.get("cred_hec_token", {}).get("password")
    baseurl = params.get("hec_url")
    if baseurl is None:
        raise Exception("The HEC URL was not provided.")

    event = args.get("event")
    host = args.get("host")
    fields = args.get("fields")
    index = args.get("index")
    source_type = args.get("source_type")
    source = args.get("source")
    time_ = args.get("time")
    request_channel = args.get("request_channel")
    batch_event_data = args.get("batch_event_data")
    entry_id = args.get("entry_id")

    if not event and not batch_event_data and not entry_id:
        raise DemistoException(
            "Invalid input: Please specify one of the following arguments: `event`, `batch_event_data`, or `entry_id`."
        )

    response_info = splunk_submit_event_hec(
        hec_token,
        baseurl,
        event,
        fields,
        host,
        index,
        source_type,
        source,
        time_,
        request_channel,
        batch_event_data,
        entry_id,
        service,
    )

    if "Success" not in response_info.text:
        return_error(f"Could not send event to Splunk {response_info.text}")
    else:
        response_dict = json.loads(response_info.text)
        if response_dict and "ackId" in response_dict:
            return_results(f"The events were sent successfully to Splunk. AckID: {response_dict['ackId']}")
        else:
            return_results("The events were sent successfully to Splunk.")


def add_findings_to_investigation(
    service: client.Service,
    investigation_id: str,
    finding_ids: list[str],
    finding_times: list[str] | None = None,
):
    """Append (link) one or more findings to an existing Splunk investigation in a single request.


    The operation is append-only and does not replace existing findings on the investigation.

    Args:
        service: Splunk service connection (already connected, namespace pre-set to ``missioncontrol/nobody``).
        investigation_id: The ID of the investigation to add the finding(s) to.
        finding_ids: A list of finding IDs to append to the investigation.
        finding_times: Optional list of times for findings added to the investigation.
            Value can be in relative, ISO, or epoch time.

    Returns:
        dict: The JSON response from the API (empty dict if the response body is empty).

    Raises:
        DemistoException: If ``finding_times`` is provided with a different length than ``finding_ids``.
        Exception: Propagates any underlying HTTP error so the caller can decide how to surface it.
    """
    if not finding_ids:
        return {}
    if finding_times is not None and len(finding_times) != len(finding_ids):
        raise DemistoException("'finding_times' must have the same length as 'finding_ids' when provided.")

    endpoint = f"public/v2/investigations/{investigation_id}/findings"
    body: dict[str, Any] = {"finding_ids": finding_ids}
    if finding_times:
        body["finding_times"] = finding_times

    demisto.debug(f"Adding {len(finding_ids)} finding(s) to investigation {investigation_id} via endpoint {endpoint}")

    response = service.post(endpoint, body=json.dumps(body))
    response_data = response.body.read()
    result = json.loads(response_data) if response_data else {}
    demisto.debug(f"Findings added to investigation {investigation_id}: {result}")
    return result


def _resolve_update_investigation_args(args: dict) -> dict[str, Any]:
    """Pull and normalize the updatable fields from ``args`` for ``splunk-update-investigation``.

    Maps the human-readable ``status`` / ``disposition`` labels to their Splunk IDs (when known)
    so the rest of the command handler can stay agnostic of label-vs-id concerns.
    """
    status = args.get("status")
    if status and status in DEFAULT_STATUSES:
        status = DEFAULT_STATUSES[status]

    disposition = args.get("disposition")
    if disposition and disposition in DEFAULT_DISPOSITIONS:
        disposition = DEFAULT_DISPOSITIONS[disposition]

    return {
        "owner": args.get("owner"),
        "urgency": args.get("urgency"),
        "status": status,
        "disposition": disposition,
        "name": args.get("name"),
        "description": args.get("description"),
        "note": args.get("note"),
        "findings": argToList(args.get("findings")),
        "finding_times": argToList(args.get("finding_times")),
    }


def splunk_update_investigation_command(service: client.Service, args: dict) -> CommandResults:
    """Update one or more Splunk ES investigations.

    Supports updating ``owner``, ``urgency``, ``status``, ``disposition``, ``name``, ``description``,
    adding a ``note``, and appending ``findings`` (finding IDs) to the investigation. At least one
    of these updatable fields must be supplied; otherwise a :class:`DemistoException` is raised.

    Args:
        service: Splunk service connection (namespace pre-set to ``missioncontrol/nobody`` by the caller).
        args: Command arguments. Must contain ``event_ids`` (CSV of investigation IDs).

    Returns:
        CommandResults: A plain-text summary mirroring the ``splunk-finding-event-edit`` output style.
    """
    investigation_ids = argToList(args.get("event_ids"))
    if not investigation_ids:
        raise DemistoException("event_ids parameter is required")

    resolved = _resolve_update_investigation_args(args)
    update_field_keys = ("owner", "urgency", "status", "disposition", "name", "description")
    updated_fields = {k: v for k, v in resolved.items() if k in update_field_keys and v is not None}
    note = resolved["note"]
    findings = resolved["findings"]
    finding_times = resolved["finding_times"]

    if finding_times and not findings:
        raise DemistoException("'finding_times' was provided without 'findings'. Provide 'findings' alongside 'finding_times'.")

    # The Splunk add-findings API targets a single investigation per call. Guard against attempting to
    # add the same findings to multiple investigations in one command, which would be unsafe and
    # is not what the caller intends. This must run BEFORE any HTTP call so we never partially apply.
    if (findings or finding_times) and len(investigation_ids) != 1:
        raise DemistoException(
            "'findings' (and 'finding_times') can only be used when updating a single investigation. "
            "Provide exactly one investigation ID in 'event_ids'."
        )

    if not updated_fields and not note and not findings:
        raise DemistoException(
            "At least one of 'owner', 'urgency', 'status', 'disposition', 'name', 'description', "
            "'note' or 'findings' must be provided to update an investigation."
        )

    successes: list[str] = []
    errors: list[str] = []

    for raw_investigation_id in investigation_ids:
        investigation_id = raw_investigation_id.strip()
        try:
            if updated_fields:
                update_investigation_or_finding(service=service, investigation_or_finding_id=investigation_id, **updated_fields)
            if note:
                add_investigation_note(service=service, investigation_or_finding_id=investigation_id, content=note)
            if findings:
                cleaned = [fid.strip() for fid in findings if fid and fid.strip()]
                if cleaned:
                    add_findings_kwargs: dict[str, Any] = {
                        "service": service,
                        "investigation_id": investigation_id,
                        "finding_ids": cleaned,
                    }
                    if finding_times:
                        add_findings_kwargs["finding_times"] = finding_times
                    add_findings_to_investigation(**add_findings_kwargs)
            successes.append(f"Successfully updated Splunk ES event {investigation_id}")
        except Exception as e:
            err = f"Failed to update Splunk ES event {investigation_id}: {e!s}"
            demisto.error(err)
            errors.append(err)

    if successes and not errors:
        readable = "Splunk ES events updated successfully:\n" + "\n".join(successes)
    elif successes and errors:
        readable = (
            "Splunk ES events partially updated:\n" "Successes:\n" + "\n".join(successes) + "\n" "Errors:\n" + "\n".join(errors)
        )
    else:
        raise DemistoException("Failed to update all Splunk ES events:\n" + "\n".join(errors))

    return CommandResults(readable_output=readable)


def splunk_edit_event_command(service: client.Service, args: dict) -> None:
    """Edit finding or investigations events in Splunk ES using the v2 investigations API.

    Args:
        service: Splunk service connection
        args: Command arguments containing event_ids and fields to update
    """
    event_ids = argToList(args.get("event_ids"))
    if not event_ids:
        return_error("event_ids parameter is required")
        return

    # Prepare the fields to update
    status = args.get("status")
    urgency = args.get("urgency")
    owner = args.get("owner")
    disposition = args.get("disposition")

    # Map the status label to the status id if needed
    if status and status in DEFAULT_STATUSES:
        status = DEFAULT_STATUSES[status]

    # Map the disposition label to the disposition id if needed
    if disposition and disposition in DEFAULT_DISPOSITIONS:
        disposition = DEFAULT_DISPOSITIONS[disposition]

    note = args.get("note")
    finding_time = args.get("finding_time")

    # Track results for each event ID
    results = []
    errors = []

    for event_id in event_ids:
        event_id = event_id.strip()
        try:
            # Update the finding using the v2 API
            update_investigation_or_finding(
                service=service,
                investigation_or_finding_id=event_id,
                owner=owner,
                urgency=urgency,
                status=status,
                disposition=disposition,
                finding_time=finding_time,
            )

            # Add note separately if provided
            if note:
                try:
                    add_investigation_note(
                        service=service,
                        investigation_or_finding_id=event_id,
                        content=note,
                        finding_time=finding_time,
                    )
                    results.append(f"Successfully updated Splunk ES event {event_id} (including note)")
                except Exception as e:
                    demisto.error(f"Failed to add note to Splunk ES event {event_id}: {e!s}")
                    results.append(f"Successfully updated Splunk ES event {event_id} (note failed: {e!s})")
            else:
                results.append(f"Successfully updated Splunk ES event {event_id}")

        except Exception as e:
            error_msg = f"Failed to update Splunk ES event {event_id}: {e!s}"
            demisto.error(error_msg)
            errors.append(error_msg)

    # Prepare the output message
    if results and not errors:
        return_results("Splunk ES events updated successfully:\n" + "\n".join(results))
    elif results and errors:
        return_results(
            "Splunk ES events partially updated:\n" "Successes:\n" + "\n".join(results) + "\n" "Errors:\n" + "\n".join(errors)
        )
    else:
        return_error("Failed to update all Splunk ES events:\n" + "\n".join(errors))


def splunk_job_status(service: client.Service, args: dict[str, Any]) -> list[CommandResults]:
    sids = argToList(args.get("sid"))
    job_results = []
    for sid in sids:
        try:
            job = service.job(sid)
        except HTTPError as error:
            if str(error) == "HTTP 404 Not Found -- Unknown sid.":
                job_results.append(CommandResults(readable_output=f"Not found job for SID: {sid}"))
            else:
                job_results.append(
                    CommandResults(readable_output=f"Querying splunk for SID: {sid} resulted in the following error {str(error)}")
                )
        else:
            status = job.state.content.get("dispatchState")
            entry_context = {"SID": sid, "Status": status}
            human_readable = tableToMarkdown("Splunk Job Status", entry_context)
            job_results.append(
                CommandResults(
                    outputs=entry_context,
                    readable_output=human_readable,
                    outputs_prefix="Splunk.JobStatus",
                    outputs_key_field="SID",
                )
            )
    return job_results


def splunk_job_share(service: client.Service, args: dict[str, Any]) -> list[CommandResults]:  # pragma: no cover
    sids = argToList(args.get("sid"))
    try:
        ttl = int(args.get("ttl", 1800))
    except ValueError:
        return_error(f"Input error: Invalid TTL provided, '{args.get('ttl')}'. Must be a valid integer.")

    job_results = []
    for sid in sids:
        try:
            job = service.job(sid)
        except HTTPError as error:
            if str(error) == "HTTP 404 Not Found -- Unknown sid.":
                job_results.append(CommandResults(readable_output=f"Not found job for SID: {sid}"))
            else:
                job_results.append(
                    CommandResults(readable_output=f"Querying splunk for SID: {sid} resulted in the following error {str(error)}")
                )
        else:
            try:
                ttl_results = True
                job.set_ttl(ttl)  # extend time-to-live for results
            except HTTPError as error:
                job_results.append(
                    CommandResults(
                        readable_output=f"Error increasing TTL for SID: {sid} resulted in the following error {str(error)}"
                    )
                )
                ttl_results = False
            try:
                share_results = True
                endpoint = f"search/jobs/{sid}/acl"
                service.post(endpoint, **{"sharing": "global", "perms.read": "*"})
            except HTTPError as error:
                job_results.append(
                    CommandResults(
                        readable_output=f"Error changing permissions for SID: {sid} resulted in the following error {str(error)}"
                    )
                )
                share_results = False

            entry_context = {"SID": sid, "TTL updated": str(ttl_results), "Sharing updated": str(share_results)}
            human_readable = tableToMarkdown("Splunk Job Updates", entry_context)
            job_results.append(
                CommandResults(
                    outputs=entry_context,
                    readable_output=human_readable,
                    outputs_prefix="Splunk.JobUpdates",
                    outputs_key_field="SID",
                )
            )
    return job_results


def splunk_parse_raw_command(args: dict[str, Any]) -> None:
    raw = args.get("raw", "")
    raw_dict = raw_to_dict(raw)
    return_results(
        CommandResults(
            outputs_prefix="Splunk.Raw.Parsed",
            raw_response=json.dumps(raw_dict),
            outputs=raw_dict,
            content_format=EntryFormat.JSON,
        )
    )


def test_module(service: client.Service, params: dict[str, Any]) -> None:
    try:
        # validate connection
        service.info()
    except AuthenticationError:
        return_error("Authentication error, please validate your credentials.")

    # validate fetch
    if params.get("isFetch"):
        t = datetime.now(pytz.UTC) - timedelta(days=3)
        time = t.strftime(ISO_FORMAT_TZ_AWARE)
        kwargs = {"count": 1, "earliest_time": time, "output_mode": OUTPUT_MODE_JSON}
        query = params["fetchQuery"]
        try:
            has_event_id = False
            for item in results.JSONResultsReader(service.jobs.oneshot(query, **kwargs)):
                if isinstance(item, results.Message):
                    continue

                if EVENT_ID not in item:
                    if MIRROR_DIRECTION.get(params.get("mirror_direction", "")):
                        return_error("Cannot mirror incidents if fetch query does not use the `notable` macro.")
                    if ENABLED_ENRICHMENTS:
                        return_error(
                            "When using the enrichment mechanism, an event_id field is needed, and thus, "
                            "one must use a fetch query of the following format: search `notable` .......\n"
                            "Please re-edit the fetchQuery parameter in the integration configuration, reset "
                            "the fetch mechanism using the splunk-reset-enriching-fetch-mechanism command and "
                            "run the fetch again."
                        )
                else:
                    has_event_id = True

        except HTTPError as error:
            return_error(str(error))

        selected_event_types = argToList(params.get("fetch_event_types")) or ["Finding"]
        if "Investigation" in selected_event_types:
            investigations_query = params.get("investigations_fetch_query") or InvestigationsFetchHandler._DEFAULT_SPL
            # Bounded probe window: 1 day back, single record. Real fetch parameters
            # (window, pagination) are not needed here — we only want to validate
            # the SPL/placeholder/connectivity round-trip.
            probe_min = to_mc_iso8601_utc((datetime.now(UTC) - timedelta(days=1)).strftime(ISO_FORMAT_TZ_AWARE))
            probe_max = to_mc_iso8601_utc(datetime.now(UTC).strftime(ISO_FORMAT_TZ_AWARE))
            try:
                probe_spl = prepare_investigations_query(investigations_query, probe_min, probe_max, limit=1, offset=0)
            except DemistoException as e:
                # Re-raise with the parameter name prefix so customers know which field to fix.
                raise DemistoException(f"'Investigations fetch query' parameter is invalid: {e}") from e

            try:
                for _ in results.JSONResultsReader(service.jobs.oneshot(probe_spl, output_mode=OUTPUT_MODE_JSON, count=1)):
                    # We only need to confirm Splunk accepts and executes the query;
                    # the contents of the (at most one) row are irrelevant here.
                    break
            except HTTPError as error:
                raise DemistoException(
                    f"'Investigations fetch query' parameter is invalid: Splunk rejected the query. {error}"
                ) from error

        # Validate custom ID generation for queries without `notable` macro
        if not has_event_id and "`notable`" not in query:
            try:
                demisto.debug("Running duplicate incident ID validation test")
                test_kwargs = {"count": 10, "earliest_time": time, "output_mode": OUTPUT_MODE_JSON}
                test_items = [
                    item
                    for item in results.JSONResultsReader(service.jobs.oneshot(query, **test_kwargs))
                    if not isinstance(item, results.Message)
                ]

                if len(test_items) >= 2:
                    custom_ids = [
                        create_incident_custom_id(
                            Finding(data=item).to_incident(
                                UserMappingObject(service, False, "splunk_xsoar_users", "xsoar_user", "splunk_user")
                            )
                        )
                        for item in test_items
                    ]

                    if len(set(custom_ids)) < len(custom_ids):
                        return_error(
                            f"Duplicate incident IDs detected in test ({len(custom_ids) - len(set(custom_ids))} duplicates).\n\n"
                            "IMPACT: Incidents with duplicate IDs will be incorrectly identified as already fetched during the "
                            "fetch process.\n"
                            "This means these incidents will be dropped and will NOT be created in XSOAR, resulting in missing "
                            "incidents.\n\n"
                            "CAUSE: The integration generates incident IDs from fields: _cd, index, _time, _indextime, _raw.\n"
                            "These fields may not provide unique values in your fetch query.\n\n"
                            "SOLUTION: Configure the 'Unique ID Fields' parameter in the integration settings.\n"
                            "Add a comma-separated list of additional fields "
                            "that will create unique combinations for each incident.\n"
                            "Example: source,host,unique_field\n\n"
                            "The parameter can be found in the integration configuration under:\n"
                            "Advanced Settings -> Unique ID Fields"
                        )
                    demisto.debug(
                        f"Duplicate ID validation passed: {len(set(custom_ids))} unique IDs from {len(custom_ids)} incidents"
                    )
            except Exception as e:
                demisto.debug(f"Could not complete duplicate ID validation: {e!s}")
    if params.get("hec_url"):
        headers = {"Content-Type": "application/json"}
        try:
            requests.get(params.get("hec_url", "") + "/services/collector/health", headers=headers, verify=VERIFY_CERTIFICATE)
        except Exception as e:
            return_error("Could not connect to HEC server. Make sure URL and token are correct.", e)


def replace_keys(data: dict[str, Any]) -> dict[str, Any]:
    if not isinstance(data, dict):
        return data
    for key in list(data.keys()):
        value = data.pop(key)
        for character in PROBLEMATIC_CHARACTERS:
            key = key.replace(character, REPLACE_WITH)

        data[key] = value
    return data


def kv_store_collection_create(service: client.Service, args: dict[str, Any]) -> CommandResults:
    try:
        service.kvstore.create(args["kv_store_name"])
    except HTTPError as error:
        if error.status == 409 and error.reason == "Conflict":
            raise DemistoException(
                f"KV store collection {service.namespace['app']} already exists.",
            ) from error
        raise

    return CommandResults(
        readable_output=f"KV store collection {service.namespace['app']} created successfully",
    )


def kv_store_collection_config(service: client.Service, args: dict[str, Any]) -> CommandResults:
    app = service.namespace["app"]
    kv_store_collection_name = args["kv_store_collection_name"]
    kv_store_fields = args["kv_store_fields"].split(",")
    for key_val in kv_store_fields:
        try:
            _key, val = key_val.split("=", 1)
        except ValueError:
            return_error(f"error when trying to parse {key_val} you possibly forgot to add the field type.")
        else:
            if _key.startswith("index."):
                service.kvstore[kv_store_collection_name].update_index(_key.replace("index.", ""), val)
            else:
                service.kvstore[kv_store_collection_name].update_field(_key.replace("field.", ""), val)
    return CommandResults(readable_output=f"KV store collection {app} configured successfully")


def kv_store_collection_create_transform(service: client.Service, args: dict[str, Any]) -> CommandResults:
    collection_name = args["kv_store_collection_name"]
    fields = args.get("supported_fields")
    if not fields:
        kv_store = service.kvstore[collection_name]
        default_keys = get_keys_and_types(kv_store).keys()
        if not default_keys:
            raise DemistoException("Please provide supported_fields or run first splunk-kv-store-collection-config")
        default_keys = (key.replace("field.", "").replace("index.", "") for key in default_keys)
        fields = f"_key,{','.join(default_keys)}"

    transforms = service.confs["transforms"]
    params = {"external_type": "kvstore", "collection": collection_name, "namespace": service.namespace, "fields_list": fields}
    transforms.create(name=collection_name, **params)
    return CommandResults(readable_output=f"KV store collection transforms {collection_name} created successfully")


def parse_key_value_pairs(key_value_pairs: str | None) -> dict[str, Any]:
    """Parses the key_value_pairs argument (a JSON object string) into a dict of stanza attributes."""
    if not key_value_pairs:
        return {}
    try:
        parsed = json.loads(key_value_pairs)
    except (ValueError, TypeError) as e:
        raise DemistoException(
            'The "key_value_pairs" argument must be a valid JSON object string, '
            'e.g. {"external_type": "kvstore", "collection": "my_collection"}.'
        ) from e
    if not isinstance(parsed, dict):
        raise DemistoException('The "key_value_pairs" argument must be a JSON object (dict), not a list or scalar.')
    return parsed


def get_configuration_file(service: client.Service, conf_file_name: str) -> client.ConfigurationFile:
    """Returns the ConfigurationFile object for the given conf file name, raising a clear error if not found."""
    try:
        return service.confs[conf_file_name]
    except KeyError as e:
        raise DemistoException(f"Configuration file '{conf_file_name}' was not found.") from e


def splunk_configuration_file_list(service: client.Service, args: dict[str, Any]) -> CommandResults:
    """Lists the configuration (.conf) files available in the given Splunk app namespace."""
    app = args.get("app", "search")
    limit = arg_to_number(args.get("limit"))
    if limit is None:
        limit = 50
    conf_files = [{"FileName": conf.name, "App": app} for conf in service.confs.list(count=limit)]
    readable_output = tableToMarkdown(
        name=f"Configuration files in app '{app}'", t=conf_files, headers=["FileName", "App"], removeNull=True
    )
    return CommandResults(
        outputs_prefix="Splunk.ConfigurationFile",
        outputs_key_field="FileName",
        outputs=conf_files,
        readable_output=readable_output,
        raw_response=conf_files,
    )


def splunk_configuration_file_create(service: client.Service, args: dict[str, Any]) -> CommandResults:
    """Creates a new, empty configuration (.conf) file in the given Splunk app namespace."""
    conf_file_name = args["conf_file_name"]
    app = args.get("app", "search")
    service.confs.create(conf_file_name)
    return CommandResults(readable_output=f"Configuration file '{conf_file_name}' in app '{app}' was created successfully.")


def splunk_configuration_stanza_create(service: client.Service, args: dict[str, Any]) -> CommandResults:
    """Creates a new stanza (configuration entry) in a Splunk .conf file, optionally with attributes."""
    conf_file_name = args["conf_file"]
    stanza_name = args["stanza_name"]
    key_value_pairs = parse_key_value_pairs(args.get("key_value_pairs"))
    conf_file = get_configuration_file(service, conf_file_name)
    stanza = conf_file.create(stanza_name)
    if key_value_pairs:
        stanza.submit(key_value_pairs)
    return CommandResults(
        readable_output=f"Stanza '{stanza_name}' in configuration file '{conf_file_name}' was created successfully."
    )


def splunk_configuration_stanza_list(service: client.Service, args: dict[str, Any]) -> CommandResults:
    """Lists the stanzas in a .conf file, or returns the key/value content of a single stanza when stanza_name is given."""
    conf_file_name = args["conf_file"]
    stanza_name = args.get("stanza_name")
    app = args.get("app", "search")
    owner = args.get("owner", "nobody")
    limit = arg_to_number(args.get("limit"))
    if limit is None:
        limit = 50
    conf_file = get_configuration_file(service, conf_file_name)

    if stanza_name:
        if stanza_name not in conf_file:
            raise DemistoException(f"Stanza '{stanza_name}' was not found in configuration file '{conf_file_name}'.")
        stanza = conf_file[stanza_name]
        content = dict(stanza.content.items())
        access = getattr(stanza, "access", None) or {}
        outputs = {
            "StanzaName": stanza.name,
            "App": access.get("app", app),
            "Owner": access.get("owner", owner),
            "Sharing": access.get("sharing"),
            "Content": content,
        }
        readable_output = tableToMarkdown(
            name=f"Stanza '{stanza_name}' in configuration file '{conf_file_name}'", t=content, removeNull=True
        )
        return CommandResults(
            outputs_prefix="Splunk.ConfigurationStanza",
            outputs_key_field="StanzaName",
            outputs=outputs,
            readable_output=readable_output,
            raw_response=content,
        )

    stanzas = []
    for stanza in conf_file.list(count=limit):
        access = getattr(stanza, "access", None) or {}
        stanzas.append(
            {
                "StanzaName": stanza.name,
                "App": access.get("app", app),
                "Owner": access.get("owner", owner),
                "Sharing": access.get("sharing"),
            }
        )
    readable_output = tableToMarkdown(
        name=f"Stanzas in configuration file '{conf_file_name}'",
        t=stanzas,
        headers=["StanzaName", "App", "Owner", "Sharing"],
        headerTransform=pascalToSpace,
        removeNull=True,
    )
    return CommandResults(
        outputs_prefix="Splunk.ConfigurationStanza",
        outputs_key_field="StanzaName",
        outputs=stanzas,
        readable_output=readable_output,
        raw_response=stanzas,
    )


def splunk_configuration_stanza_update(service: client.Service, args: dict[str, Any]) -> CommandResults:
    """Upserts attributes on an existing stanza in a Splunk .conf file (existing keys overwritten, others untouched)."""
    conf_file_name = args["conf_file"]
    stanza_name = args["stanza_name"]
    key_value_pairs = parse_key_value_pairs(args.get("key_value_pairs"))
    if not key_value_pairs:
        raise DemistoException('The "key_value_pairs" argument is required and must contain at least one attribute to update.')
    conf_file = get_configuration_file(service, conf_file_name)
    if stanza_name not in conf_file:
        raise DemistoException(f"Stanza '{stanza_name}' was not found in configuration file '{conf_file_name}'.")
    conf_file[stanza_name].submit(key_value_pairs)
    return CommandResults(
        readable_output=f"Stanza '{stanza_name}' in configuration file '{conf_file_name}' was updated successfully."
    )


def splunk_configuration_stanza_delete(service: client.Service, args: dict[str, Any]) -> CommandResults:
    """Deletes a stanza (configuration entry) from a Splunk .conf file via the configs/conf-{file} REST endpoint.

    This complements splunk-kv-store-collection-delete-entry, which removes records from the KV Store
    but leaves the corresponding stanza in the .conf file (e.g. transforms.conf) intact.
    """
    conf_file_name = args["conf_file"]
    stanza_name = args["stanza_name"]
    conf_file = get_configuration_file(service, conf_file_name)
    if stanza_name not in conf_file:
        raise DemistoException(f"Stanza '{stanza_name}' was not found in configuration file '{conf_file_name}'.")
    conf_file.delete(stanza_name)
    return CommandResults(
        readable_output=f"Stanza '{stanza_name}' in configuration file '{conf_file_name}' was deleted successfully."
    )


def batch_kv_upload(kv_data_service_client: client.KVStoreCollectionData, json_data: str) -> dict[str, Any]:
    if json_data.startswith("[") and json_data.endswith("]"):
        record: Record = kv_data_service_client._post(
            "batch_save", headers=client.KVStoreCollectionData.JSON_HEADER, body=json_data.encode("utf-8")
        )
        return dict(record.items())
    elif json_data.startswith("{") and json_data.endswith("}"):
        return kv_data_service_client.insert(json_data.encode("utf-8"))
    else:
        raise DemistoException(
            'kv_store_data argument should be in json format. (e.g. {"key": "value"} or [{"key": "value"}, {"key": "value"}]'
        )


def kv_store_collection_add_entries(service: client.Service, args: dict[str, Any]) -> None:
    kv_store_data = args.get("kv_store_data", "")
    kv_store_collection_name = args["kv_store_collection_name"]
    indicator_path = args.get("indicator_path")
    batch_kv_upload(service.kvstore[kv_store_collection_name].data, kv_store_data)
    indicators_timeline = None
    if indicator_path:
        kv_store_data = json.loads(kv_store_data)
        indicators = extract_indicator(indicator_path, kv_store_data if isinstance(kv_store_data, list) else [kv_store_data])
        indicators_timeline = IndicatorsTimeline(
            indicators=indicators,
            category="Integration Update",
            message=f"Indicator added to {kv_store_collection_name} store in Splunk",
        )
    return_results(
        CommandResults(readable_output=f"Data added to {kv_store_collection_name}", indicators_timeline=indicators_timeline)
    )


def kv_store_collections_list(service: client.Service) -> None:
    app_name = service.namespace["app"]
    names = [x.name for x in service.kvstore.iter()]
    readable_output = "list of collection names {}\n| name |\n| --- |\n|{}|".format(app_name, "|\n|".join(names))
    return_results(
        CommandResults(outputs_prefix="Splunk.CollectionList", outputs=names, readable_output=readable_output, raw_response=names)
    )


def kv_store_collection_data_delete(service: client.Service, args: dict[str, Any]) -> None:
    kv_store_collection_name = args["kv_store_collection_name"].split(",")
    for store in kv_store_collection_name:
        service.kvstore[store].data.delete()
    return_results(f"The values of the {args['kv_store_collection_name']} were deleted successfully")


def kv_store_collection_delete(service: client.Service, args: dict[str, Any]) -> CommandResults:
    kv_store_names = args["kv_store_name"]
    for store in kv_store_names.split(","):
        service.kvstore[store].delete()
    return CommandResults(readable_output=f"The following KV store {kv_store_names} were deleted successfully.")


def build_kv_store_query(kv_store: client.KVStoreCollection, args: dict[str, Any]) -> str | dict[str, Any]:
    if "key" in args and "value" in args:
        _type = get_key_type(kv_store, args["key"])
        args["value"] = _type(args["value"]) if _type else args["value"]
        return json.dumps({args["key"]: args["value"]})
    elif "limit" in args:
        return {"limit": args["limit"]}
    else:
        return args.get("query", "{}")


def kv_store_collection_data(service: client.Service, args: dict[str, Any]) -> None:
    stores = args["kv_store_collection_name"].split(",")

    for i, store_res in enumerate(get_store_data(service)):
        store = service.kvstore[stores[i]]

        if store_res:
            readable_output = tableToMarkdown(name=f"list of collection values {store.name}", t=store_res)
            return_results(
                CommandResults(
                    outputs_prefix="Splunk.KVstoreData",
                    outputs={store.name: store_res},
                    readable_output=readable_output,
                    raw_response=store_res,
                )
            )
        else:
            return_results(get_kv_store_config(store))


def kv_store_collection_delete_entry(service: client.Service, args: dict[str, Any]) -> None:
    store_name = args["kv_store_collection_name"]
    indicator_path = args.get("indicator_path")
    store: client.KVStoreCollection = service.kvstore[store_name]
    query = build_kv_store_query(store, args)
    store_res = next(get_store_data(service))
    indicators = extract_indicator(indicator_path, store_res) if indicator_path else []
    store.data.delete(query=query)
    indicators_timeline = (
        IndicatorsTimeline(
            indicators=indicators, category="Integration Update", message=f"Indicator deleted from {store_name} store in Splunk"
        )
        if indicators
        else None
    )
    return_results(
        CommandResults(
            readable_output=f"The values of the {store_name} were deleted successfully", indicators_timeline=indicators_timeline
        )
    )


def check_error(service: client.Service, args: dict[str, Any]) -> None:
    app = args.get("app_name")
    store_name = args.get("kv_store_collection_name")
    if app not in service.apps:
        raise DemistoException("app not found")
    elif store_name and store_name not in service.kvstore:
        raise DemistoException("KV Store not found")


def get_key_type(kv_store: client.KVStoreCollection, _key: str) -> type | None:
    keys_and_types = get_keys_and_types(kv_store)
    types = {"number": float, "string": str, "cidr": str, "boolean": bool, "time": str}
    index = f"index.{_key}"
    field = f"field.{_key}"
    val_type = keys_and_types.get(field) or keys_and_types.get(index) or ""
    return types.get(val_type)


def get_keys_and_types(kv_store: client.KVStoreCollection) -> dict[str, str]:
    keys = kv_store.content()
    for key_name in list(keys.keys()):
        if not (key_name.startswith(("field.", "index."))):
            del keys[key_name]
    return keys


def get_kv_store_config(kv_store: client.KVStoreCollection) -> str:
    keys = get_keys_and_types(kv_store)
    readable = [f"#### configuration for {kv_store.name} store", "| field name | type |", "| --- | --- |"]
    readable.extend(f"| {_key} | {val} |" for _key, val in keys.items())
    return "\n".join(readable)


def extract_indicator(indicator_path: str, _dict_objects: list[dict[str, Any]]) -> list[str]:
    indicators = []
    indicator_paths = indicator_path.split(".")
    for indicator_obj in _dict_objects:
        indicator = ""
        for path in indicator_paths:
            indicator = indicator_obj.get(path, {})
        indicators.append(str(indicator))
    return indicators


def get_store_data(service: client.Service) -> Any:
    args = demisto.args()
    stores = args["kv_store_collection_name"].split(",")

    for store in stores:
        kvstore: client.KVStoreCollection = service.kvstore[store]
        query = build_kv_store_query(kvstore, args)
        if isinstance(query, str):
            query = {"query": query}
        yield kvstore.data.query(**query)


def get_connection_args(params: dict[str, Any]) -> dict[str, Any]:
    """
    This function gets the connection arguments: host, port, app, and verify.
    Parses the server_url parameter to extract host and port, with port 8089 as default.

    Returns: connection args
    """
    server_url = params.get("server_url", "")
    # If URL doesn't have a scheme, add one for proper parsing
    if not server_url.startswith(("http://", "https://")):
        server_url = f"https://{server_url}"
    parsed = urllib.parse.urlparse(server_url)

    # Extract host (hostname or netloc without port)
    host = parsed.hostname or parsed.netloc.split(":")[0]
    host = host.rstrip("/")

    # Extract port or use default 8089
    port = parsed.port if parsed.port else 8089

    app = params.get("app", "-")
    return {
        "host": host,
        "port": port,
        "app": app or "-",
        "verify": VERIFY_CERTIFICATE,
        "retries": 3,
        "retryDelay": 3,
    }


def handle_message(item: results.Message | dict) -> bool:
    """Checks if the response from JSONResultsReader is a message object.
        The message can be info etc.
        such as: "the test table is empty"

    Args:
        item (results.Message | dict): The item to be checked. It can be either a `results.Message`
            object or a dictionary.

    Returns:
        bool: Returns `True` if the item is an instance of `results.Message`, `False` otherwise.

    """
    if isinstance(item, results.Message):
        demisto.info(f"Splunk-SDK message: {item.message}")
        return True
    return False


def main() -> None:  # pragma: no cover
    command = demisto.command()
    params = demisto.params()
    args = demisto.args()

    if command == "splunk-parse-raw":
        splunk_parse_raw_command(args)
        sys.exit(0)
    service = None
    proxy = argToBoolean(params.get("proxy", False))

    connection_args = get_connection_args(params)
    password = params["authentication"]["password"]
    connection_args["splunkToken"] = password
    connection_args["autologin"] = True

    if proxy:
        handle_proxy()

    # Validate that the note tags are different
    if NOTE_TAG_TO_SPLUNK == NOTE_TAG_FROM_SPLUNK:
        raise DemistoException("Note Tag to Splunk and Note Tag from Splunk cannot have the same value.")

    connection_args["handler"] = requests_handler

    if (service := client.connect(**connection_args)) is None:
        return_error("Could not connect to Splunk")

    mapper = UserMappingObject(
        service,
        params.get("userMapping"),
        params.get("user_map_lookup_name"),
        params.get("xsoar_user_field"),
        params.get("splunk_user_field"),
    )

    # The command command holds the command sent from the user.
    if command == "test-module":
        test_module(service, params)
        return_results("ok")
    elif command == "splunk-reset-enriching-fetch-mechanism":
        reset_enriching_fetch_mechanism()
    elif command == "splunk-search":
        return_results(splunk_search_command(service, args))
    elif command == "splunk-job-create":
        splunk_job_create_command(service, args)
    elif command == "splunk-results":
        splunk_results_command(service, args)
    elif command == "splunk-get-indexes":
        splunk_get_indexes_command(service, app=connection_args.get("app", "-"))
    elif command == "fetch-incidents":
        demisto.info("########### FETCH #############")
        fetch_incidents(service, mapper)
        extensive_log("[SplunkPy] Fetch Incidents was successfully executed.")
    elif command == "splunk-submit-event":
        splunk_submit_event_command(service, args)
    elif command == "splunk-finding-event-edit" and service is not None:
        service.namespace = namespace(app="missioncontrol", owner="nobody")
        splunk_edit_event_command(service, args)
    elif command == "splunk-update-investigation" and service is not None:
        service.namespace = namespace(app="missioncontrol", owner="nobody")
        return_results(splunk_update_investigation_command(service, args))
    elif command == "splunk-investigation-create" and service is not None:
        return_results(splunk_create_investigation_command(service, args))
    elif command == "splunk-investigation-list" and service is not None:
        return_results(splunk_list_investigations_command(service, args))
    elif command == "splunk-submit-event-hec":
        splunk_submit_event_hec_command(params, service, args)
    elif command == "splunk-job-status":
        return_results(splunk_job_status(service, args))
    elif command == "splunk-job-share":
        return_results(splunk_job_share(service, args))
    elif command.startswith("splunk-configuration-") and service is not None:
        # All configuration commands share the same namespace scoping (app + owner), so it is set once here.
        service.namespace = namespace(app=args.get("app", "search"), owner=args.get("owner", "nobody"), sharing="app")

        if command == "splunk-configuration-file-list":
            return_results(splunk_configuration_file_list(service, args))
        elif command == "splunk-configuration-file-create":
            return_results(splunk_configuration_file_create(service, args))
        elif command == "splunk-configuration-stanza-create":
            return_results(splunk_configuration_stanza_create(service, args))
        elif command == "splunk-configuration-stanza-list":
            return_results(splunk_configuration_stanza_list(service, args))
        elif command == "splunk-configuration-stanza-update":
            return_results(splunk_configuration_stanza_update(service, args))
        elif command == "splunk-configuration-stanza-delete":
            return_results(splunk_configuration_stanza_delete(service, args))
    elif command.startswith("splunk-kv-") and service is not None:
        app = args.get("app_name", "search")
        service.namespace = namespace(app=app, owner="nobody", sharing="app")
        check_error(service, args)

        if command == "splunk-kv-store-collection-create":
            return_results(kv_store_collection_create(service, args))
        elif command == "splunk-kv-store-collection-config":
            return_results(kv_store_collection_config(service, args))
        elif command == "splunk-kv-store-collection-create-transform":
            return_results(kv_store_collection_create_transform(service, args))
        elif command == "splunk-kv-store-collection-delete":
            return_results(kv_store_collection_delete(service, args))
        elif command == "splunk-kv-store-collections-list":
            kv_store_collections_list(service)
        elif command == "splunk-kv-store-collection-add-entries":
            kv_store_collection_add_entries(service, args)
        elif command in ["splunk-kv-store-collection-data-list", "splunk-kv-store-collection-search-entry"]:
            kv_store_collection_data(service, args)
        elif command == "splunk-kv-store-collection-data-delete":
            kv_store_collection_data_delete(service, args)
        elif command == "splunk-kv-store-collection-delete-entry":
            kv_store_collection_delete_entry(service, args)

    elif command == "get-mapping-fields":
        return_results(get_mapping_fields_command(service, mapper, params))
    elif command == "get-remote-data":
        raise NotImplementedError(f"the {command} command is not implemented, use get-modified-remote-data instead.")
    elif command == "get-modified-remote-data":
        demisto.info("########### MIRROR IN #############")
        try:
            get_modified_remote_data_command(
                service=service,
                args=args,
                close_incident=params.get("close_incident"),
                close_end_statuses=params.get("close_end_status_statuses"),
                close_extra_labels=argToList(params.get("close_extra_labels", "")),
                mapper=mapper,
            )
        except Exception as e:
            return_error(f"An error occurred during the Mirror In - in get_modified_remote_data_command: {e}")
    elif command == "update-remote-system" and service is not None:
        demisto.info("########### MIRROR OUT #############")
        service.namespace = namespace(app="missioncontrol", owner="nobody")
        return_results(update_remote_system_command(args, params, service, mapper))
    elif command == "splunk-get-username-by-xsoar-user":
        return_results(mapper.get_splunk_user_by_xsoar_command(args))
    else:
        raise NotImplementedError(f"Command not implemented: {command}")


if __name__ in ["__main__", "__builtin__", "builtins"]:
    main()