OpenAi ChatGPT v3

This integration assists security professionals with security investigations, threat hunting, and anomaly detection by leveraging OpenAI GPT models' natural language conversation capabilities.

Messaging and Conferencing · OpenAI

Details

IDOpenAi ChatGPT v3
ProviderOpenAI
CategoryMessaging and Conferencing
From Version6.0.0
Docker Imagedemisto/parse-emails:0.1.48.10120494
Supported ModulesAgentix XSIAM

README

OpenAI GPT

Instance Configuration

  • Generate an API Key

    1. Sign up or log in to OpenAI developer platform.
    2. Generate a new API key at OpenAI developer platform - api-keys.
  • Choose a GPT model to interact with

    1. This integration supports only the ‘Chat Completions’ endpoint. Therefore, you can only configure models that support this endpoint (https://api.openai.com/v1/chat/completions).

    2. For tasks requiring deep understanding and extensive inputs, opt for more advanced models (e.g. gpt-4). These models offer a larger context window, allowing them to process bigger documents, and provide more refined and comprehensive responses.
      The more elementary models (e.g. gpt-3.5) often provide shallower answers and input analysis.
      Refer to Models overview for more information.

  • Text generation setting (Optional)

    1. max-tokens: The maximum number of tokens that can be generated for the response. (Allows controlling tokens’ consumption). Default: unset.
    2. temperature: Sets the randomness in responses. Lower values (closer to 0) produce more deterministic and consistent outputs, while higher values (up to 2) increase randomness and variety. It is generally recommended altering this or top_p but not both. Default: 1.
    3. top_p: Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Lower values (closer to 0) result in more focused outputs, while higher values (closer to 1) increase diversity. It is generally recommended altering this or temperature but not both. Default: unset.
  • Event Collector — Generate API Keys

    1. Admin API Key (required for OpenAI Audit logs): generate from the OpenAI Platform admin console. Used to call /v1/organization/audit_logs.
    2. Compliance API Key (required for any Compliance event type): generate from the ChatGPT Platform. Used to call /v1/compliance/workspaces/{workspace_id}/....
    3. Workspace ID (required for any Compliance event type): the UUID of the compliance workspace whose events you want to collect.
  • Event Collector — Select event types to fetch

    Toggle Fetch events, then select one or more Events types to fetch:

    User-facing label Source Required credentials
    OpenAI Audit logs OpenAI Platform — Admin API Admin API Key
    Conversation Messages ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    Apps ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    Apps Auth ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    Compliance Audit ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    Auth ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    Codex ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    ChatGPT ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    Codex Security ChatGPT Platform — Compliance API Compliance API Key + Workspace ID
    Workspace Agents ChatGPT Platform — Compliance API Compliance API Key + Workspace ID

    Selecting an event type without its matching credentials raises an informative error at instance test time, naming the missing parameter.

  • Event Collector — Datasets

    Each Event Collector stream lands in its own Cortex dataset:

    Stream Vendor Product Dataset
    OpenAI Audit logs openai chatgpt_audit openai_chatgpt_audit_raw
    Compliance logs (all) openai chatgpt_compliance openai_chatgpt_compliance_raw
  • Event Collector — Tuning (Optional)

    Parameter Default Description
    Maximum number of OpenAI Audit events per fetch 1000 Cap on Audit events ingested per fetch cycle.
    Maximum number of Compliance events per fetch 900 Cap on Compliance events ingested per fetch cycle.
    Events Fetch Interval 1 minute How often the scheduled fetch runs.
    ChatGPT Server URL https://api.chatgpt.com Base URL of the ChatGPT Compliance API. Override only for non-default tenants.
  • Click ‘Test’

Commands

You can execute these commands from the Cortex XSOAR CLI, as part of an automation, or in a playbook.

gpt-send-message


Send a message as a prompt to the GPT model.

!gpt-send-message message="<MESSAGE_TEXT>"

Input

Argument Name Description Required
message The message to send to the GPT model wrapped with quotes. Yes
reset_conversation_history Whether to reset conversation history or keep it as context for the sent message. (Conversation history is not reset by default). No
max_tokens The maximum number of tokens that can be generated for the response. Overrides text generation setting for the specific message sent. No
temperature Sets the randomness in responses. Overrides text generation setting for the specific message sent. No
top_p Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Overrides text generation setting for the specific message sent. No

gpt-check-email-body


Check email body for possible security issues.

!gpt-check-email-body entryId="<ENTRY_ID_OF_UPLOADED_EML_FILE>"

Input

Argument Name Description Required
entryId Entry ID of an uploaded .eml file from the context window. Yes
additionalInstructions Provide additional instructions for the GPT model when analyzing the email body. No
max_tokens The maximum number of tokens that can be generated for the response. Overrides text generation setting for the specific message sent. No
temperature Sets the randomness in responses. Overrides text generation setting for the specific message sent. No
top_p Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Overrides text generation setting for the specific message sent. No

gpt-check-email-header


Check email body for possible security issues.

!gpt-check-email-header entryId="<ENTRY_ID_OF_UPLOADED_EML_FILE>"

Input

Argument Name Description Required
entryId Entry ID of an uploaded .eml file from context window. Yes
additionalInstructions Provide additional instructions for the GPT model when analyzing the email headers. No
max_tokens The maximum number of tokens that can be generated for the response. Overrides text generation setting for the specific message sent. No
temperature Sets the randomness in responses. Overrides text generation setting for the specific message sent. No
top_p Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Overrides text generation setting for the specific message sent. No

gpt-analyze-email-header


Analyze email headers for potential security issues using the OpenAI Responses API. This command uses the Responses API which is recommended for all new projects (instead of gpt-check-email-header which uses the Chat Completions API).

!gpt-analyze-email-header entry_id="3@123" additional_instructions="Pay close attention to SPF/DKIM."

Input

Argument Name Description Required
entry_id Entry ID of an uploaded .eml file. Yes
additional_instructions Additional instructions or security issue to focus on. Substituted into the prompt template. No
max_tokens The maximum number of tokens that can be generated for the response. Maps internally to the API body field max_output_tokens. No
temperature Sets the randomness in responses. Lower values (closer to 0) produce more deterministic and consistent outputs, while higher values (up to 2) increase randomness and variety. No
top_p Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Range 0–1. No
reasoning_effort Reasoning effort level for reasoning models (o1, o3, o4, gpt-5). Controls how much thinking the model does before responding. Possible values: low, medium, high. No

Context Output

Path Type Description
OpenAiChatGPTV3.Response Unknown The conversation state including the response_id.
OpenAiChatGPTV3.Response.user String The prompt sent to the model.
OpenAiChatGPTV3.Response.assistant String The assistant response text.
OpenAiChatGPTV3.Response.response_id String The OpenAI response ID.

Human Readable Output

Two war-room entries are produced:

  1. A table of the parsed email headers.
  2. The AI verdict followed by a token-usage table. A Reasoning tokens row appears in the usage table when a reasoning model is used.

gpt-analyze-email-body


Analyze email body for potential security risks using the OpenAI Responses API. This command uses the Responses API which is recommended for all new projects (instead of gpt-check-email-body which uses the Chat Completions API).

!gpt-analyze-email-body entry_id="3@123"

Input

Argument Name Description Required
entry_id Entry ID of an uploaded .eml file. Yes
additional_instructions Additional instructions or security issue to focus on. Substituted into the prompt template. No
max_tokens The maximum number of tokens that can be generated for the response. Maps internally to the API body field max_output_tokens. No
temperature Sets the randomness in responses. Lower values (closer to 0) produce more deterministic and consistent outputs, while higher values (up to 2) increase randomness and variety. No
top_p Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Range 0–1. No
reasoning_effort Reasoning effort level for reasoning models (o1, o3, o4, gpt-5). Controls how much thinking the model does before responding. Possible values: low, medium, high. No

Context Output

Path Type Description
OpenAiChatGPTV3.Response Unknown The conversation state including the response_id.
OpenAiChatGPTV3.Response.user String The prompt sent to the model.
OpenAiChatGPTV3.Response.assistant String The assistant response text.
OpenAiChatGPTV3.Response.response_id String The OpenAI response ID.

Human Readable Output

Two war-room entries are produced:

  1. A table of the parsed email body (text and HTML).
  2. The AI verdict followed by a token-usage table. A Reasoning tokens row appears in the usage table when a reasoning model is used.

gpt-create-soc-email-template


Create an email template out of the conversation context to be sent from the SOC.

!gpt-create-soc-email-template

Input

Argument Name Description Required
additionalInstructions Provide additional instructions for the GPT model when analyzing the email headers. No
max_tokens The maximum number of tokens that can be generated for the response. Overrides text generation setting for the specific message sent. No
temperature Sets the randomness in responses. Overrides text generation setting for the specific message sent. No
top_p Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Overrides text generation setting for the specific message sent. No

gpt-draft-soc-email


Draft a SOC email template using the OpenAI Responses API. This command uses the Responses API which is recommended for all new projects (instead of gpt-create-soc-email-template which uses the Chat Completions API). Consumes prior conversation context by design (e.g. from a preceding gpt-analyze-email-body call).

Cortex XSOAR sequence (typical phishing flow)

!gpt-analyze-email-body entry_id="3@123"
…assistant returns analysis…
!gpt-draft-soc-email additional_instructions="Notify the user the email was quarantined."

Input

Argument Name Description Required
additional_instructions Specific issue or focus area to weave into the template. Substituted into the prompt template. No
max_tokens The maximum number of tokens that can be generated for the response. Maps internally to the API body field max_output_tokens. No
temperature Sets the randomness in responses. Lower values (closer to 0) produce more deterministic and consistent outputs, while higher values (up to 2) increase randomness and variety. No
top_p Enables nucleus sampling where only the top ‘p’ percent of probable tokens are considered. Range 0–1. No
reasoning_effort Reasoning effort level for reasoning models (o1, o3, o4, gpt-5). Controls how much thinking the model does before responding. Possible values: low, medium, high. No

Context Output

Path Type Description
OpenAiChatGPTV3.Response Unknown The conversation state including the response_id.
OpenAiChatGPTV3.Response.user String The prompt sent to the model.
OpenAiChatGPTV3.Response.assistant String The assistant response text.
OpenAiChatGPTV3.Response.response_id String The OpenAI response ID.

Human Readable Output

Two war-room entries are produced:

  1. The SOC email template context output (replace_existing=True — running twice overwrites the previous draft).
  2. The AI-generated template followed by a token-usage table. A Reasoning tokens row appears in the usage table when a reasoning model is used.

openai-get-events


Manually fetch a bounded batch of Audit and/or Compliance events for development/debugging. Does NOT advance the persisted last_run cursor, so it is safe to run against production tenants. Use should_push_events=true to additionally ingest the fetched events into the matching Cortex dataset.

Base Command

openai-get-events

Input

Argument Name Description Required
event_type The event type(s) to fetch. Comma-separated list. Possible values: OpenAI Audit logs, Conversation Messages, Apps, Apps Auth, Compliance Audit, Auth, Codex, ChatGPT, Codex Security, Workspace Agents. Defaults to the values configured in the integration parameters. No
limit Maximum number of events to return per stream. Default: 50. No
start_time Lookback start time for the fetch. Supports ISO 8601 or relative time (e.g., 3 days ago, 2099-01-01T00:00:00Z). No
should_push_events If true, the command also pushes the retrieved events to Cortex (Audit -> openai_chatgpt_audit_raw, Compliance -> openai_chatgpt_compliance_raw). Possible values: true, false. Default: false. No

Context Output

Path Type Description
OpenAI.Event.id String The unique identifier of the event.
OpenAI.Event._event_type String The upstream event_type for Compliance events. Left empty for Audit events.
OpenAI.Event.source_log_type String The source log type used by downstream parsing rules.
OpenAI.Event._time Date The event timestamp in ISO 8601 format.

Human Readable Output

OpenAI GPT Events

id _event_type source_log_type _time
FAKE_AUDIT_EVENT_001   openai_audit_logs 2099-01-01T00:00:00Z
FAKE_LISTING_002 AUDIT_LOG compliance_audit_log 2099-01-02T00:00:00Z

gpt-create-response


Sends a message to the OpenAI Responses API and receives the generated response. This command uses the Responses API which is recommended for all new projects (instead of gpt-send-message which uses the Chat Completions API). Supports multi-turn conversations via previous_response_id, reasoning effort control for o-series and gpt-5 models, and background execution.

Base Command

gpt-create-response

Input

Argument Name Description Required
message The user message to send. Required
reset_conversation_history Whether to discard the existing conversation context and start fresh. Possible values are: yes, no. Default is no. Optional
max_tokens The maximum number of output tokens. Falls back to instance config. Maps internally to the API body field max_output_tokens. Optional
temperature The randomness level in responses. Falls back to instance config. Range 0-2. Lower values produce more deterministic outputs, while higher values increase variety. Optional
top_p The nucleus sampling threshold. Falls back to instance config. Range 0-1. Lower values result in more focused outputs, while higher values increase diversity. Optional
reasoning_effort The reasoning effort level. Honored only for reasoning families (o1, o3, o4, gpt-5); silently dropped on others. Default medium. Possible values are: none, minimal, low, medium, high, xhigh. Optional
background Whether to run the model response in the background. When true, the command uses polling to wait for the response to complete. Possible values are: true, false. Optional
compact_threshold The token threshold at which compaction should be triggered for this entry. Minimum 1000. Optional
model The model to use. Use the gpt-list-models command to see available models. Falls back to instance config. Optional

Context Output

Path Type Description
OpenAiChatGPTV3.Response Unknown The conversation state, which includes the response_id for multi-turn continuity.
OpenAiChatGPTV3.Response.user String The user message sent.
OpenAiChatGPTV3.Response.assistant String The assistant response text.
OpenAiChatGPTV3.Response.response_id String The OpenAI response ID used for multi-turn conversation continuity.

gpt-list-models


Lists all models available to the configured API key. Lets users discover models per their actual API-key tier without redeploying the integration when OpenAI ships new ones.

Base Command

gpt-list-models

Input

| Argument Name | Description | Required |
| — | — | — |

Context Output

Path Type Description
OpenAiChatGPTV3.Model.Id String The model identifier (e.g., gpt-4, gpt-3.5-turbo).
OpenAiChatGPTV3.Model.Created Number The Unix timestamp when the model was created.
OpenAiChatGPTV3.Model.OwnedBy String The organization or entity that owns the model.

gpt-create-moderation


Runs text or an image through the OpenAI Moderations API and returns per-category flagging results. Exactly one of text, entry_id, or image_url must be provided.

Base Command

gpt-create-moderation

Input

Argument Name Description Required
text A comma-separated list of text strings to moderate. Exactly one of text, entry_id, or image_url must be provided. Optional
entry_id The war-room entry ID of an uploaded image file. The file is base64-encoded internally and posted as a data URL. Exactly one of text, entry_id, or image_url must be provided. Optional
image_url The publicly reachable HTTP(S) URL of an image (limited to 20 MB). Exactly one of text, entry_id, or image_url must be provided. Optional
model The moderation model to use. Possible values are: omni-moderation-latest, omni-moderation-2024-09-26. Default is omni-moderation-latest. Optional

Context Output

Path Type Description
OpenAiChatGPTV3.Moderation.Input.input_type String The type of input that was moderated (text, image, or image_url).
OpenAiChatGPTV3.Moderation.Input.input_value String The value of the input that was moderated.
OpenAiChatGPTV3.Moderation.Flagged Boolean Whether the content was flagged by the moderation model.
OpenAiChatGPTV3.Moderation.Categories Unknown The object of boolean values indicating which categories were flagged.
OpenAiChatGPTV3.Moderation.CategoryScores Unknown The object of float values indicating the confidence score for each category.

<~PLATFORM>

License Requirements

The following configuration parameters require the Cortex XSIAM license:

  • Fetch events

</~PLATFORM>

Configuration parameters

  • url — Server URL (required)
  • chatgpt_api_url — ChatGPT Server URL
  • apikey
  • admin_api_key — Admin API Key
  • compliance_api_key — Compliance API Key
  • workspace_id — Workspace ID
  • model-select — Model
  • model-freetext — Model (Optional - overrides selected choice)
  • max_tokens — Max tokens
  • temperature — Temperature
  • top_p — Top P
  • insecure — Trust any certificate (not secure)
  • proxy — Use system proxy settings
  • isFetchEvents — Fetch events
  • event_types_to_fetch — Events types to fetch
  • audit_max_fetch — Maximum number of OpenAI Audit events per fetch
  • compliance_max_fetch — Maximum number of Compliance events per fetch
  • eventFetchInterval — Events Fetch Interval

Commands (11)

  • gpt-analyze-email-body

    Analyzes email body for potential security risks using the OpenAI Responses API. This is the Responses-API counterpart of gpt-check-email-body (which uses the Chat Completions API).

  • gpt-analyze-email-header

    Analyzes email headers for potential security issues using the OpenAI Responses API. This is the Responses-API counterpart of gpt-check-email-header (which uses the Chat Completions API).

  • gpt-check-email-body

    Checks the email body for possible security issues. Enables you to ask subsequent questions on the provided information using the 'gpt-send-message' command, and resets the conversation context by default.

  • gpt-check-email-header

    Checking email header for possible security issues. It is possible to keep asking questions on the provided info using 'gpt-send-message'. Resets conversation context by default.

  • gpt-create-moderation

    Runs text or an image through the OpenAI Moderations API and returns per-category flagging results. Exactly one of text, entry_id, or image_url must be provided.

  • gpt-create-response

    Sends a message to the OpenAI Responses API and receives the generated response. This command uses the Responses API which is recommended for all new projects (instead of gpt-send-message which uses the Chat Completions API). Supports multi-turn conversations via previous_response_id, reasoning effort control for o-series and gpt-5 models, and background execution.

  • gpt-create-soc-email-template

    Create an email template out of the conversation context to be sent from the SOC.

  • gpt-draft-soc-email

    Drafts a SOC email template using the OpenAI Responses API. This command uses the Responses API which is recommended for all new projects (instead of gpt-create-soc-email-template which uses the Chat Completions API). Consumes prior conversation context by design (e.g. from a preceding gpt-analyze-email-body call).

  • gpt-list-models

    Lists all models available to the configured API key. Lets users discover models per their actual API-key tier without redeploying the integration when OpenAI ships new ones.

  • gpt-send-message

    Send a plain message to the selected GPT model and receive the generated response.

  • openai-get-events

    Manually retrieves events from OpenAI. Use this command for development and debugging only, as it may produce duplicate events, exceed API rate limits, or disrupt the fetch mechanism.

import base64
import json
import mimetypes

import demistomock as demisto  # noqa: F401
import parse_emails
import urllib3
from CommonServerPython import *  # noqa: F401

from CommonServerUserPython import *  # noqa

from concurrent.futures import Future, ThreadPoolExecutor, as_completed
from dataclasses import dataclass, field, replace
from datetime import datetime, timedelta, UTC
from typing import Any
from collections.abc import Callable

# Disable insecure warnings
urllib3.disable_warnings()


# region Constants - Chat / Email
# =================================
# Existing GPT chat / email constants
# =================================
INTEGRATION_NAME = "OpenAI GPT"


class Config:
    """Global static configuration shared across all integration features."""

    DATE_FORMAT = "%Y-%m-%dT%H:%M:%SZ"

    VENDOR = "openai"
    PRODUCT_AUDIT = "admin_audit"
    PRODUCT_COMPLIANCE = "chatgpt_compliance"

    DEFAULT_CHATGPT_URL = "https://api.chatgpt.com"

    RETRY_POLICY: dict = {
        "retries": 3,
        "status_list_to_retry": [429, 500, 502, 503, 504],
        "backoff_factor": 2,
        "raise_on_status": True,
    }

    DEFAULT_AUDIT_MAX_FETCH = 1000
    DEFAULT_COMPLIANCE_MAX_FETCH = 900
    DEFAULT_GET_EVENTS_LIMIT = 50
    AUDIT_PAGE_SIZE = 100
    COMPLIANCE_PAGE_SIZE = 100
    MAX_PAGES_PER_FETCH = 50  # Safety cap on pagination loops.

    DEFAULT_FIRST_FETCH = "1 hour ago"

    # Test-module probe: per-stream max-events ceiling when test_module exercises the collector
    # via the same fetch_stream pipeline (mirrors Koi's TEST_MODULE_MAX_EVENTS=1).
    TEST_MODULE_MAX_EVENTS = 1


class Stream:
    """Stream identifiers used by the parallel fetch dispatcher and dataset routing."""

    AUDIT = "audit"
    COMPLIANCE = "compliance"


EML_FILE_PREFIX = ".eml"


class ResponsesConfig:
    """Configuration constants for the OpenAI Responses API and background polling."""

    # Regex prefixes that identify reasoning-capable model families.
    # When the model name starts with any of these, the usage table includes a "Reasoning tokens" row.
    REASONING_MODEL_PREFIXES = ("o1", "o3", "o4", "gpt-5")

    # Terminal statuses for background responses polling.
    BACKGROUND_TERMINAL_STATUSES = frozenset({"completed", "failed", "cancelled", "incomplete"})
    BACKGROUND_PENDING_STATUSES = frozenset({"queued", "in_progress"})

    DEFAULT_POLLING_INTERVAL_SECS = 10
    DEFAULT_POLLING_TIMEOUT_SECS = 600


class ApiPaths:
    """Centralized OpenAI API endpoint paths.

    Chat-completions is hosted on `api.openai.com`; audit/compliance live on `api.chatgpt.com`.
    """

    CHAT_COMPLETIONS = "v1/chat/completions"
    RESPONSES = "v1/responses"
    MODELS = "v1/models"
    MODERATIONS = "v1/moderations"
    AUDIT_LOGS = "v1/organization/audit_logs"

    @classmethod
    def compliance_logs(cls, workspace_id: str) -> str:
        return f"v1/compliance/workspaces/{workspace_id}/logs"

    @classmethod
    def compliance_log_content(cls, workspace_id: str, log_id: str) -> str:
        return f"v1/compliance/workspaces/{workspace_id}/logs/{log_id}"


class EventType:
    """User-facing event-type labels used by the multi-select fetch parameter."""

    AUDIT = "OpenAI Audit logs"
    CONVERSATION_MESSAGE = "Conversation Messages"
    APP_LOG = "Apps"
    APP_AUTH_LOG = "Apps Auth"
    AUDIT_LOG = "Compliance Audit"
    AUTH_LOG = "Auth"
    CODEX_LOG = "Codex"
    CHATGPT_PLUGIN_SPREADSHEET = "ChatGPT"
    CODEX_SECURITY_LOG = "Codex Security"
    CUSTOM_AGENTS_LOG = "Workspace Agents"


class ComplianceEvent:
    """Upstream `event_type` query values used by `/v1/compliance/.../logs`.

    These are the canonical strings sent on the wire and stored as `_event_type` on
    the resulting events; they map to the user-facing `EventType` labels above.
    """

    CONVERSATION_MESSAGE = "CONVERSATION_MESSAGE"
    APP_LOG = "APP_LOG"
    APP_AUTH_LOG = "APP_AUTH_LOG"
    AUDIT_LOG = "AUDIT_LOG"
    AUTH_LOG = "AUTH_LOG"
    CODEX_LOG = "CODEX_LOG"
    CHATGPT_PLUGIN_SPREADSHEET = "CHATGPT_PLUGIN_SPREADSHEET"
    CODEX_SECURITY_LOG = "CODEX_SECURITY_LOG"
    CUSTOM_AGENTS_LOG = "CUSTOM_AGENTS_LOG"


class SourceLogType:
    """`source_log_type` values written to Compliance events so a single shared dataset can be
    disambiguated downstream by parsing/modeling rules. Audit events are NOT given a
    `source_log_type` because they land in a dedicated dataset.
    """

    CONVERSATION_MESSAGE = "conversation_message"
    COMPLIANCE_AUDIT_LOG = "compliance_audit_log"
    AUTH_LOG = "auth_log"
    APP_AUTH_LOG = "app_auth_log"
    APP_LOG = "app_log"
    CODEX_LOG = "codex_log"
    CHATGPT_PLUGIN_SPREADSHEET = "chatgpt_plugin_spreadsheet"
    CODEX_SECURITY_LOG = "codex_security_log"
    CUSTOM_AGENT_LOG = "custom_agent_log"


# Mapping: user-facing label -> upstream event_type query value (for /v1/compliance/.../logs).
EVENT_TYPE_LABEL_TO_API: dict[str, str] = {
    EventType.CONVERSATION_MESSAGE: ComplianceEvent.CONVERSATION_MESSAGE,
    EventType.APP_LOG: ComplianceEvent.APP_LOG,
    EventType.APP_AUTH_LOG: ComplianceEvent.APP_AUTH_LOG,
    EventType.AUDIT_LOG: ComplianceEvent.AUDIT_LOG,
    EventType.AUTH_LOG: ComplianceEvent.AUTH_LOG,
    EventType.CODEX_LOG: ComplianceEvent.CODEX_LOG,
    EventType.CHATGPT_PLUGIN_SPREADSHEET: ComplianceEvent.CHATGPT_PLUGIN_SPREADSHEET,
    EventType.CODEX_SECURITY_LOG: ComplianceEvent.CODEX_SECURITY_LOG,
    EventType.CUSTOM_AGENTS_LOG: ComplianceEvent.CUSTOM_AGENTS_LOG,
}

# Mapping: upstream event_type value -> source_log_type used downstream by parsing/modeling rules.
COMPLIANCE_EVENT_TYPE_TO_SOURCE_LOG_TYPE: dict[str, str] = {
    ComplianceEvent.CONVERSATION_MESSAGE: SourceLogType.CONVERSATION_MESSAGE,
    ComplianceEvent.AUDIT_LOG: SourceLogType.COMPLIANCE_AUDIT_LOG,
    ComplianceEvent.AUTH_LOG: SourceLogType.AUTH_LOG,
    ComplianceEvent.APP_AUTH_LOG: SourceLogType.APP_AUTH_LOG,
    ComplianceEvent.APP_LOG: SourceLogType.APP_LOG,
    ComplianceEvent.CODEX_LOG: SourceLogType.CODEX_LOG,
    ComplianceEvent.CHATGPT_PLUGIN_SPREADSHEET: SourceLogType.CHATGPT_PLUGIN_SPREADSHEET,
    ComplianceEvent.CODEX_SECURITY_LOG: SourceLogType.CODEX_SECURITY_LOG,
    ComplianceEvent.CUSTOM_AGENTS_LOG: SourceLogType.CUSTOM_AGENT_LOG,
}


class LastRunKey:
    """Keys used to persist per-stream pagination state across fetch cycles.

    The Audit and Compliance streams use *separate* keys so each stream's pagination state
    is independent (one stream failing or being disabled never affects the other).
    """

    # --- Audit stream (cursor-based pagination via the `after` query param) ---
    AUDIT_AFTER = "audit_after"  # opaque cursor (last_id from API), passed verbatim as `after=` next run.
    AUDIT_FIRST_FETCH_SEED = "audit_first_fetch_seed"  # Unix-seconds seed kept across empty first-fetches.

    # --- Compliance stream (time-based pagination via the `after` query param + per-id dedup) ---
    COMPLIANCE_LAST_END_TIME = "compliance_last_end_time"  # ISO timestamp echoed by the listing response.
    COMPLIANCE_LAST_IDS = "compliance_last_ids"  # listing IDs seen at last_end_time, deduped on next run.


CHECK_EMAIL_HEADERS_PROMPT = """
I have a set of email headers.
Analyze these headers for any potential security issues such as spoofing, phishing attempts, or other malicious activity.
Please identify any suspicious fields, explain why they might be concerning, and suggest any further actions that could be taken \
to investigate or mitigate these issues.
Additional instructions: {}

'''
{}
'''

Please, review each header, highlighting any red flags and explaining the potential risks associated with them.
Make you answer very concise and easily readable, with references to the email headers if there are, otherwise do not refer to \
hypothetical problems.
"""

CHECK_EMAIL_BODY_PROMPT = """
I have this email body that I suspect may contain security risks such as phishing links, suspicious attachments,
or signs of social engineering. Please analyze the content of this email body, identify any elements that may pose security
threats, and explain why these elements are concerning. Also, suggest any steps that could be taken to further verify these risks
or protect against these threats.
{}
'''
{}
'''

Highlight potential security risks, and explain the implications of such risks.
Make you answer very concise and easily readable, with references to the email body if there are, otherwise do not refer to \
hypothetical problems.
"""

CREATE_SOC_EMAIL_TEMPLATE_PROMPT = """
Based on the details provided in our conversation and any specific instructions you have been given,
create a professional email template suitable for a Security Operations Center (SOC).
The template should be adaptable, clearly structured, and include placeholders for specific incident details,
recommendations for action, and any necessary escalation points.
Please ensure the tone is appropriate for communication within a cybersecurity context.
{}
"""


class ArgAndParamNames:
    MODEL = "model"
    MESSAGE = "message"
    RESET_CONVERSATION_HISTORY = "reset_conversation_history"
    ENTRY_ID = "entry_id"
    ADDITIONAL_INSTRUCTIONS = "additional_instructions"
    MAX_TOKENS = "max_tokens"
    TEMPERATURE = "temperature"
    TOP_P = "top_p"


class Roles:
    ASSISTANT = "assistant"
    USER = "user"


class EmailParts:
    HEADERS = "headers"
    BODY = "body"


class OpenAiClient(BaseClient):
    """OpenAI HTTP client wrapping three API surfaces:

    - `api.openai.com` (`Connect` section, via inherited `base_url` + `url_suffix=`):
        * Chat Completions - uses `apikey`.
        * Audit Logs - uses `admin_api_key`.
    - `api.chatgpt.com` (`Connect - Compliance` section, via `self.chatgpt_base_url` + `full_url=`):
        * Compliance Logs - uses `compliance_api_key` and `workspace_id`.
    """

    def __init__(
        self,
        url: str,
        api_key: str,
        model: str,
        proxy: bool,
        verify: bool,
        admin_api_key: str = "",
        compliance_api_key: str = "",
        chatgpt_base_url: str = Config.DEFAULT_CHATGPT_URL,
    ):
        super().__init__(base_url=url, proxy=proxy, verify=verify)

        self.api_key = api_key
        self.model = model
        self.admin_api_key = admin_api_key
        self.compliance_api_key = compliance_api_key
        self.chatgpt_base_url = chatgpt_base_url.rstrip("/") + "/"
        self.headers = {"Authorization": f"Bearer {self.api_key}", "Content-Type": "application/json"}

    def get_chat_completions(
        self, chat_context: List[dict[str, str]], completion_params: dict[str, str | None]
    ) -> dict[str, Any]:
        """Gets the response to a chat_completions request using the OpenAI API."""
        if not self.api_key:
            raise DemistoException(
                "API Key is required for the chat-completion commands "
                "(gpt-send-message, gpt-check-email-header, gpt-check-email-body, gpt-create-soc-email-template). "
                "Configure the 'API Key' integration parameter and try again."
            )

        options: Dict[str, Any] = {ArgAndParamNames.MODEL: self.model}
        max_tokens = completion_params.get(ArgAndParamNames.MAX_TOKENS, None)
        if max_tokens:
            options[ArgAndParamNames.MAX_TOKENS] = int(max_tokens)

        temperature = completion_params.get(ArgAndParamNames.TEMPERATURE, None)
        if temperature:
            options[ArgAndParamNames.TEMPERATURE] = float(temperature)

        top_p = completion_params.get(ArgAndParamNames.TOP_P, None)
        if top_p:
            options[ArgAndParamNames.TOP_P] = float(top_p)

        options["messages"] = chat_context
        demisto.debug(
            f"[API Chat Completions] Calling | model={options.get(ArgAndParamNames.MODEL)} | "
            f"messages_count={len(chat_context)} | "
            f"max_tokens={options.get(ArgAndParamNames.MAX_TOKENS)} | "
            f"temperature={options.get(ArgAndParamNames.TEMPERATURE)} | "
            f"top_p={options.get(ArgAndParamNames.TOP_P)}"
        )
        return self._http_request(method="POST", url_suffix=ApiPaths.CHAT_COMPLETIONS, json_data=options, headers=self.headers)

    # region commands using api.openai.com and ApiKey
    def list_models(self) -> dict[str, Any]:
        """List all models available to the configured API key (GET /v1/models).

        Returns:
            The parsed JSON response dict from the API containing a ``data`` list.
        """
        if not self.api_key:
            raise DemistoException(
                "API Key is required for the gpt-list-models command. "
                "Configure the 'API Key' integration parameter and try again."
            )
        demisto.debug("[API Models] Listing available models.")
        return self._http_request(
            method="GET",
            url_suffix=ApiPaths.MODELS,
            headers=self.headers,
            retries=3,
            status_list_to_retry=[429, 500, 502, 503, 504],
            backoff_factor=5,
        )

    def create_moderation(self, body: dict[str, Any]) -> dict[str, Any]:
        """Call the OpenAI Moderations API (POST /v1/moderations).

        Args:
            body: The full JSON body containing ``model`` and ``input``.

        Returns:
            The parsed JSON response dict from the API.
        """
        if not self.api_key:
            raise DemistoException(
                "API Key is required for the gpt-create-moderation command. "
                "Configure the 'API Key' integration parameter and try again."
            )
        demisto.debug(f"[API Moderations] Calling | model={body.get('model')}")
        return self._http_request(
            method="POST",
            url_suffix=ApiPaths.MODERATIONS,
            json_data=body,
            headers=self.headers,
            retries=3,
            status_list_to_retry=[429, 500, 502, 503, 504],
            backoff_factor=5,
        )

    def create_response(self, body: dict[str, Any]) -> dict[str, Any]:
        """Call the OpenAI Responses API (POST /v1/responses).

        Uses the regular OpenAI API key and the Server URL base.

        Args:
            body: The full JSON body to send (model, input, and optional params).

        Returns:
            The parsed JSON response dict from the API.
        """
        if not self.api_key:
            raise DemistoException(
                "API Key is required for the Responses API commands "
                "(gpt-create-response, gpt-analyze-email-header, gpt-analyze-email-body, gpt-draft-soc-email). "
                "Configure the 'API Key' integration parameter and try again."
            )
        demisto.debug(
            f"[API Responses] Calling | model={body.get('model')} | "
            f"max_output_tokens={body.get('max_output_tokens')} | "
            f"temperature={body.get('temperature')} | "
            f"top_p={body.get('top_p')} | "
            f"reasoning_effort={body.get('reasoning', {}).get('effort') if body.get('reasoning') else None} | "
            f"background={body.get('background')}"
        )
        return self._http_request(
            method="POST",
            url_suffix=ApiPaths.RESPONSES,
            json_data=body,
            headers=self.headers,
            retries=3,
            status_list_to_retry=[429, 500, 502, 503, 504],
            backoff_factor=5,
        )

    def get_response(self, response_id: str) -> dict[str, Any]:
        """Retrieve an existing response by ID (GET /v1/responses/{response_id}).

        Used for polling background responses until they reach a terminal state.

        Args:
            response_id: The response ID returned from a ``POST /v1/responses`` call.

        Returns:
            The parsed JSON response dict from the API.
        """
        if not self.api_key:
            raise DemistoException(
                "API Key is required for polling Responses API results. "
                "Configure the 'API Key' integration parameter and try again."
            )
        demisto.debug(f"[API Responses] Polling | response_id={response_id}")
        return self._http_request(
            method="GET",
            url_suffix=f"{ApiPaths.RESPONSES}/{response_id}",
            headers=self.headers,
            retries=3,
            status_list_to_retry=[429, 500, 502, 503, 504],
            backoff_factor=5,
        )

    # endregion

    # region Event Collector - Audit Logs (Admin API)
    def get_audit_logs(
        self,
        after: str | None = None,
        limit: int = Config.AUDIT_PAGE_SIZE,
        effective_at_gt: int | None = None,
    ) -> dict[str, Any]:
        """Fetch one cursor-based page of audit logs.

        On the first ever run, `after` is omitted and `effective_at_gt` (Unix seconds) bounds
        the starting point. On subsequent runs, `after` is the previous response's `last_id`.

        Returns the response envelope `{data, has_more, last_id}`.
        """
        if not self.admin_api_key:
            raise DemistoException("Admin API Key is required to fetch OpenAI Audit logs.")

        params: dict[str, Any] = {"limit": min(limit, Config.AUDIT_PAGE_SIZE), "order": "asc"}
        if after:
            params["after"] = after
        elif effective_at_gt is not None:
            # No cursor yet (first ever run) - constrain the starting point by time.
            params["effective_at[gt]"] = effective_at_gt

        headers = {"Authorization": f"Bearer {self.admin_api_key}", "Accept": "application/json"}
        demisto.debug(
            f"[API Audit] Fetching audit logs page | limit={params['limit']} | "
            f"after_cursor_set={bool(after)} | effective_at_gt_set={effective_at_gt is not None}"
        )

        # Fetch as text first so we can defensively handle both single-JSON and concatenated-JSON
        # response shapes. The OpenAI Audit API has been observed in production returning newline-
        # separated multi-document responses ("Extra data: line 1 column N" from `response.json()`),
        # so we parse the body ourselves and tolerate both shapes.
        raw_body = self._http_request(
            method="GET",
            url_suffix=ApiPaths.AUDIT_LOGS,
            params=params,
            headers=headers,
            resp_type="text",
            **Config.RETRY_POLICY,
        )
        response = _parse_json_or_concatenated(raw_body, log_prefix="[API Audit]")
        if not isinstance(response, dict):
            # Concatenated-JSON path returned a list of records - wrap in the standard envelope so
            # callers consistently see {data, has_more, last_id}.
            response = {"data": response, "has_more": False, "last_id": None}
        page_size_returned = len(response.get("data") or [])
        demisto.debug(
            f"[API Audit] Page received | events_count={page_size_returned} | "
            f"has_more={response.get('has_more')} | last_id_set={bool(response.get('last_id'))}"
        )
        return response

    # endregion

    # region Event Collector - Compliance Logs (ChatGPT Platform - `Connect - Compliance` section)
    def _compliance_headers(self) -> dict[str, str]:
        """Build the auth headers for the ChatGPT Compliance API.

        The credential MUST carry the `Bearer` scheme - a bare key is not a well-formed
        `Authorization` value and the API rejects it with 401 "Access token is missing"
        without ever evaluating the key.
        """
        return {"Authorization": f"Bearer {self.compliance_api_key}", "Accept": "application/json"}

    def list_compliance_logs(
        self,
        workspace_id: str,
        event_types: list[str],
        after: str,
        limit: int | None = None,
    ) -> dict[str, Any]:
        """List compliance log entries (step 1 of the two-step compliance flow).

        Returns a normalized `{data, last_end_time, has_more}` dict where `last_end_time` is the
        ISO 8601 upper bound for this page (used as `after=` on the next run) and `data` is a
        list of entry descriptors carrying at least `id`, `event_type`, `end_time`.
        """
        if not self.compliance_api_key:
            raise DemistoException("Compliance API Key is required to fetch OpenAI Compliance logs.")
        if not workspace_id:
            raise DemistoException("Workspace ID is required to fetch OpenAI Compliance logs.")

        # Build query params; `event_type` repeats per value.
        params: list[tuple[str, Any]] = [("after", after)]
        for et in event_types:
            params.append(("event_type", et))
        if limit is not None:
            # The Compliance API rejects limit > 100 with HTTP 422; clamp defensively so the per-fetch
            # `compliance_max_fetch` (which can be larger) never leaks straight into the wire `limit`.
            effective_limit = min(limit, Config.COMPLIANCE_PAGE_SIZE)
            params.append(("limit", effective_limit))

        full_url = self.chatgpt_base_url + ApiPaths.compliance_logs(workspace_id)
        headers = self._compliance_headers()
        demisto.debug(
            f"[API Compliance List] Listing logs | url={full_url} | auth_scheme=Bearer | "
            f"event_types_count={len(event_types)} | after_set={bool(after)} | limit={limit}"
        )

        # Fetch as text first so we can defensively handle single-JSON, concatenated-JSON, and empty
        # response shapes. The Compliance Listing API has been observed in production returning
        # newline-separated multi-document responses ("Expecting value: line 2 column 1" from
        # `response.json()`), so we parse the body ourselves and tolerate all three shapes.
        raw_body = self._http_request(
            method="GET",
            full_url=full_url,
            params=params,
            headers=headers,
            resp_type="text",
            **Config.RETRY_POLICY,
        )
        response = _parse_json_or_concatenated(raw_body, log_prefix="[API Compliance List]")
        # Normalize the response shape - the API may return either a bare list or a dict with `data`/`last_end_time`.
        if isinstance(response, list):
            normalized: dict[str, Any] = {"data": response, "last_end_time": None}
        elif isinstance(response, dict):
            normalized = {
                "data": response.get("data", []) or [],
                "last_end_time": response.get("last_end_time"),
                "has_more": response.get("has_more"),
            }
        else:
            normalized = {"data": [], "last_end_time": None}
        demisto.debug(
            f"[API Compliance List] Listing returned {len(normalized['data'])} entry(ies) | "
            f"last_end_time_set={bool(normalized.get('last_end_time'))}"
        )
        return normalized

    def get_compliance_log_content(self, workspace_id: str, log_id: str) -> list[dict[str, Any]]:
        """Fetch the content for one compliance log entry (step 2 of the two-step flow).

        The response body is concatenated JSON / JSONL (not a single JSON document), so we
        fetch raw text and parse it via `parse_concatenated_json`.
        """
        if not self.compliance_api_key:
            raise DemistoException("Compliance API Key is required to fetch OpenAI Compliance log content.")

        full_url = self.chatgpt_base_url + ApiPaths.compliance_log_content(workspace_id, log_id)
        headers = self._compliance_headers()
        demisto.debug(f"[API Compliance Content] Fetching content for one log entry | url={full_url} | auth_scheme=Bearer")

        # The response body is a stream of concatenated JSON objects (or a JSONL file) - fetch raw text.
        raw_body = self._http_request(method="GET", full_url=full_url, headers=headers, resp_type="text", **Config.RETRY_POLICY)
        records = parse_concatenated_json(raw_body)
        demisto.debug(f"[API Compliance Content] Parsed {len(records)} record(s) from response body.")
        return records

    # endregion

    # region Event Collector - Cortex ingestion
    def send_events(self, events: list[dict], product: str) -> None:
        """Send events to XSIAM under `Config.VENDOR` and the given `product` (dataset suffix).

        Args:
            events: List of event dicts to send.
            product: One of `Config.PRODUCT_AUDIT` or `Config.PRODUCT_COMPLIANCE`.
        """
        if not events:
            demisto.debug(f"[API Send] No events to send for product={product}.")
            return
        demisto.debug(f"[API Send] Sending {len(events)} event(s) to XSIAM | vendor={Config.VENDOR} | product={product}")
        send_events_to_xsiam(events=events, vendor=Config.VENDOR, product=product)
        demisto.debug(f"[API Send] Successfully sent {len(events)} event(s) | vendor={Config.VENDOR} | product={product}")

    # endregion


def setup_args(args: Dict[str, Any], params: Dict[str, Any]):
    """Backfill model-configuration args from the instance params when not provided as command args."""
    for key in (ArgAndParamNames.MAX_TOKENS, ArgAndParamNames.TEMPERATURE, ArgAndParamNames.TOP_P):
        if not args.get(key) and params.get(key):
            args[key] = params.get(key)


def _resolve_model(args: dict[str, Any], client: OpenAiClient) -> str:
    """Resolve the model name from command args, falling back to the instance-level default.

    Args:
        args: Command arguments from ``demisto.args()``.
        client: The configured ``OpenAiClient`` instance (carries the instance-level model).

    Returns:
        The resolved model name.

    Raises:
        DemistoException: If no model is specified anywhere.
    """
    model: str = args.get(ArgAndParamNames.MODEL, "") or client.model
    if not model:
        raise DemistoException("No model specified. Provide it as a command argument or configure it in the instance settings.")
    return model


def _build_responses_api_body(
    args: dict[str, Any],
    params: dict[str, Any],
    model: str,
    prompt: str,
) -> dict[str, Any]:
    """Build a Responses API request body from command args and instance params.

    Centralises the repeated pattern of resolving ``max_output_tokens``,
    ``temperature``, ``top_p``, and ``reasoning_effort`` from *args* (with
    *params* as fallback) and assembling them into the ``POST /v1/responses``
    request body.

    Args:
        args: Command arguments from ``demisto.args()``.
        params: Instance parameters from ``demisto.params()``.
        model: The resolved model name (use ``_resolve_model`` first).
        prompt: The input text / prompt to send.

    Returns:
        A dict ready to be passed to ``client.create_response(body)``.
    """
    body: dict[str, Any] = {
        "model": model,
        "input": prompt,
    }

    # max_output_tokens: command arg takes precedence, then instance param
    max_output_tokens = args.get(ArgAndParamNames.MAX_TOKENS)
    if max_output_tokens is None:
        max_output_tokens = params.get(ArgAndParamNames.MAX_TOKENS)
    if max_output_tokens is not None:
        body["max_output_tokens"] = int(max_output_tokens)

    # temperature: command arg takes precedence, then instance param
    temperature = args.get(ArgAndParamNames.TEMPERATURE)
    if temperature is None:
        temperature = params.get(ArgAndParamNames.TEMPERATURE)
    if temperature is not None:
        body["temperature"] = float(temperature)

    # top_p: command arg takes precedence, then instance param
    top_p = args.get(ArgAndParamNames.TOP_P)
    if top_p is None:
        top_p = params.get(ArgAndParamNames.TOP_P)
    if top_p is not None:
        body["top_p"] = float(top_p)

    # reasoning_effort (only for reasoning model families)
    reasoning_effort = args.get("reasoning_effort")
    if reasoning_effort is not None:
        body["reasoning"] = {"effort": reasoning_effort}

    return body


def conversation_to_chat_context(conversation: List[dict[str, str]]) -> List[dict[str, str]]:
    """A 'Conversation' list that was retrieved from 'demisto.context()' is formatted to be more intuitive for XSOAR users
    and is formatted as: [
                            {'user': '<USER_MESSAGE_0>, 'assistant': '<ASSISTANT_MESSAGE_0>},
                            {'user': '<USER_MESSAGE_1>', 'assistant': '<ASSISTANT_MESSAGE_1>'},
                             ...
                        ].

    The conversational format that is supported by the 'Chat Completions' endpoint is a sequence of messages,
     labeled with roles:
        [
            {'role': 'user', 'content': '<USER_MESSAGE_0>'},
            {'role': 'assistant', 'content': '<ASSISTANT_MESSAGE_0>'},
            {'role': 'user', 'content': '<USER_MESSAGE_1>'},
            {'role': 'assistant', 'content': '<ASSISTANT_MESSAGE_1>'},
            ...
        ]

    Therefore, it has to be transformed.
    """

    chat_context = []
    demisto.debug(f"[Chat Context] Expanding {len(conversation)} turn(s) from prior conversation.")
    for element in conversation:
        chat_context.append({"role": Roles.USER, "content": element.get(Roles.USER, "")})
        chat_context.append({"role": Roles.ASSISTANT, "content": element.get(Roles.ASSISTANT, "")})

    return chat_context


def get_chat_context(reset_conversation_history: bool, message: str) -> List[dict[str, str]]:
    """
    Retrieves the existing chat conversation history from the incident context, if exists.
    If `reset_conversation_history` is True, or if no conversation history exists, it initializes a new conversation list
    with the given message and returns it.

    Args:
        reset_conversation_history (bool): Flag to determine whether to reset the existing conversation history.
        message (str): The new message to be added to the conversation.

    Returns:
        List[Dict[str, str]]: The updated conversation history with the new message appended.
    """
    # Retrieve or initialize conversation history based on the context and reset flag
    conversation = demisto.context().get("OpenAiChatGPTV3", {}).get("Conversation")

    if reset_conversation_history or not conversation:
        conversation = []
        demisto.debug("[Chat Context] History reset or initialized as empty.")
    else:
        demisto.debug(
            f"[Chat Context] Loaded prior conversation from context | "
            f"type={type(conversation).__name__} | turns={len(conversation) if hasattr(conversation, '__len__') else 'n/a'}"
        )

    # Create the chat context which is suitable with the required format for a 'chat-completions' request.
    chat_context = conversation_to_chat_context(conversation)
    chat_context.append({"role": Roles.USER, "content": message})
    demisto.debug(f"[Chat Context] Appended new user message | total_messages={len(chat_context)}")
    return chat_context


def extract_assistant_message(response: dict[str, Any]) -> str:
    """Extract the assistant's reply content from a chat-completions response."""
    choices: list = response.get("choices") or []
    if not choices:
        raise DemistoException("Could not retrieve message from response: 'choices' field is empty or missing.")

    message: dict = choices[0].get("message") or {}
    if not message:
        raise DemistoException("Could not retrieve message from response: 'choices[0].message' field is empty or missing.")

    response_content: str = message.get("content") or ""
    if not response_content:
        raise DemistoException("Could not retrieve message from response: 'choices[0].message.content' is empty.")

    demisto.debug(f"[Chat Response] Extracted assistant message | length={len(response_content)} chars")
    return response_content


def get_email_parts(entry_id: str) -> tuple[List[dict[str, str]] | None, str | None, str | None, str | None]:
    """
    Extracts and parses the headers, text body, and HTML body from an .eml file identified by a given entry ID.

    Args:
    - entry_id (str): The unique identifier for the uploaded .eml file in the war room.

    Returns:
    - tuple[List[Dict[str, str]] | None, str | None, str | None]: A tuple containing three elements:
        - headers (List[Dict[str, str]] | None): A list of dictionaries where each dictionary represents an email header.
        - text_body (str | None): The plain text body of the email, if available.
        - html_body (str | None): The HTML body of the email, if available.
    """
    if not entry_id:
        DemistoException("Provide an entryId of an uploaded '.eml' file.")

    demisto.debug("[Email Parts] Resolving uploaded .eml file path from entry_id.")
    get_file_path_res = demisto.getFilePath(entry_id)
    file_path = get_file_path_res["path"]
    file_name = get_file_path_res["name"]

    if not file_name.endswith(EML_FILE_PREFIX):
        DemistoException("Provided 'entry_id' does not point to a valid '.eml' file.")

    email_parser = parse_emails.EmailParser(file_path=file_path)
    email_parser.parse()

    headers, text_body, html_body = (
        email_parser.parsed_email.get("Headers", None),
        email_parser.parsed_email.get("Text", None),
        email_parser.parsed_email.get("HTML", None),
    )
    demisto.debug(
        f"[Email Parts] Parsed .eml | headers_count={len(headers) if headers else 0} | "
        f"text_body_present={bool(text_body)} | html_body_present={bool(html_body)}"
    )
    return headers, text_body, html_body, file_name


def check_email_part(email_part: str, client: OpenAiClient, args: dict[str, Any]) -> CommandResults:
    """
    Checks email parts (headers/body) for potential security issues using predefined prompts
    ('CHECK_EMAIL_HEADERS_PROMPT', 'CHECK_EMAIL_BODY_PROMPT') that are sent to the GPT model.
    """
    entry_id: str = args.get(ArgAndParamNames.ENTRY_ID, "")
    email_headers, email_text_body, email_html_body, file_name = get_email_parts(entry_id)
    additional_instructions = (
        (f"openai-gpt check_email_part Additional instructions: {ArgAndParamNames.ADDITIONAL_INSTRUCTIONS}\n")
        if args.get(ArgAndParamNames.ADDITIONAL_INSTRUCTIONS, "")
        else ""
    )

    if email_part == EmailParts.HEADERS:
        demisto.debug(f"[Email Check] Checking email headers | headers_present={bool(email_headers)}")
        if email_headers:
            email_headers_formatted = {
                header["name"]: header["value"] for header in email_headers if "name" in header and "value" in header
            }
            readable_input = tableToMarkdown(name=f"{file_name} headers:", t=email_headers_formatted, sort_headers=False)
            check_email_part_message = CHECK_EMAIL_HEADERS_PROMPT.format(additional_instructions, readable_input)

        else:
            raise DemistoException("'parse_emails' did not extract any email headers from the provided file..")
    elif email_part == EmailParts.BODY:
        demisto.debug(
            f"[Email Check] Checking email body | "
            f"text_body_present={bool(email_text_body)} | html_body_present={bool(email_html_body)}"
        )

        if not email_text_body and not email_html_body:
            raise DemistoException("'email_parser' did not extract any email body from the provided file.")

        email_text_body = email_text_body if email_text_body else ""
        email_html_body = email_html_body if email_html_body else ""

        email_body = {"Body/Text": email_text_body, "HTML/Text": email_html_body}

        readable_input = tableToMarkdown(name=f"{file_name} body:", t=email_body, sort_headers=False)
        check_email_part_message = CHECK_EMAIL_BODY_PROMPT.format(additional_instructions, readable_input)
    else:
        raise DemistoException("Invalid email part to check provided.")

    demisto.debug(f"[Email Check] Built prompt | length={len(check_email_part_message or '')} chars")

    # Starting a new conversation as of a new topic discussed.
    args.update({ArgAndParamNames.RESET_CONVERSATION_HISTORY: "yes", ArgAndParamNames.MESSAGE: check_email_part_message})
    send_message_command_results, response = send_message_command(client, args)

    # Displaying the analyzed email part to the war room and setting the context for the email checking response
    # prior to returning the 'send-message-command' results and the entire conversation to the context.
    return_results(
        CommandResults(
            readable_output=readable_input,
            outputs_prefix="OpenAiChatGPTV3.Email" + email_part.capitalize(),
            outputs={"Email" + email_part.capitalize(): readable_input, "Response": response},
            replace_existing=True,
        )
    )
    return send_message_command_results


""" COMMAND FUNCTIONS """


def test_module(client: OpenAiClient, params: dict) -> str:
    """Probe chat-completions (when its API Key is set) and each selected/credentialed event-collector stream.

    Each capability is independent: an instance using only chat commands needs the API Key; an
    instance using only the event collector needs the matching Admin/Compliance keys (and
    workspace_id for compliance). At least ONE configured capability is required for the test
    to be meaningful.
    """
    demisto.debug("[Test Module] Starting test-module probes.")

    if client.api_key:
        demisto.debug("[Test Module] Probing chat-completions endpoint...")
        try:
            chat_message = {"role": "user", "content": ""}
            completion_params = {
                ArgAndParamNames.MAX_TOKENS: params.get(ArgAndParamNames.MAX_TOKENS, None),
                ArgAndParamNames.TEMPERATURE: params.get(ArgAndParamNames.TEMPERATURE, None),
                ArgAndParamNames.TOP_P: params.get(ArgAndParamNames.TOP_P, None),
            }
            client.get_chat_completions(chat_context=[chat_message], completion_params=completion_params)
        except DemistoException as e:
            if "Forbidden" in str(e) or "Authorization" in str(e):
                demisto.error(f"[Test Module] Chat-completions probe failed with auth error: {e}")
                return "Authorization Error: make sure API Key is correctly set"
            demisto.error(f"[Test Module] Chat-completions probe raised non-auth error: {e}")
            raise
        demisto.debug("[Test Module] Chat-completions probe passed.")

        demisto.debug("[Test Module] Probing Responses API endpoint...")
        try:
            response_params = {
                ArgAndParamNames.MODEL: client.model or "gpt-3.5-turbo",
                ArgAndParamNames.TEMPERATURE: params.get(ArgAndParamNames.TEMPERATURE, None),
                ArgAndParamNames.TOP_P: params.get(ArgAndParamNames.TOP_P, None),
                "max_output_tokens": params.get(ArgAndParamNames.MAX_TOKENS, None),
                "input": "Present random english sentence",
            }
            client.create_response(body=response_params)
        except DemistoException as e:
            if "Forbidden" in str(e) or "Authorization" in str(e):
                demisto.error(f"[Test Module] Responses API probe failed with auth error: {e}")
                return "Authorization Error: make sure API Key is correctly set"
            demisto.error(f"[Test Module] Responses API probe raised non-auth error: {e}")
            raise
        demisto.debug("[Test Module] Responses API probe passed.")
    else:
        demisto.debug("[Test Module] Chat-completions and Responses API probes skipped (no API Key configured).")

    collector_params = parse_collector_params(params)
    streams_to_probe = [
        stream
        for stream in collector_params.streams_to_run()
        if (stream == Stream.AUDIT and client.admin_api_key)
        or (stream == Stream.COMPLIANCE and client.compliance_api_key and collector_params.workspace_id)
    ]
    if not streams_to_probe:
        if not client.api_key:
            raise DemistoException(
                "No capability is configured: provide either the 'API Key' (for chat-completion commands) "
                "or an 'Admin API Key' / 'Compliance API Key' + 'Workspace ID' (for the event collector). "
                "At least one is required."
            )
        demisto.debug(
            "[Test Module] No collector streams to probe (none selected, or missing credentials/workspace_id)."
            " Returning 'ok' after chat-completions probe only."
        )
        return "ok"

    probe_params = replace(
        collector_params,
        audit_max_fetch=Config.TEST_MODULE_MAX_EVENTS,
        compliance_max_fetch=Config.TEST_MODULE_MAX_EVENTS,
    )
    demisto.debug(f"[Test Module] Probing collector streams (max_events={Config.TEST_MODULE_MAX_EVENTS}): {streams_to_probe}")

    for stream in streams_to_probe:
        demisto.debug(f"[Test Module] Probing {stream} stream...")
        try:
            fetch_stream(client=client, stream=stream, last_run={}, collector_params=probe_params)
        except Exception as exc:
            demisto.error(f"[Test Module] {stream} probe FAILED: {exc}")
            raise DemistoException(f"Test failed for the '{stream}' event-collector stream: {exc}") from exc
        demisto.debug(f"[Test Module] {stream} probe succeeded.")

    demisto.debug("[Test Module] All probes passed - returning 'ok'.")
    return "ok"


def send_message_command(client: OpenAiClient, args: dict[str, Any]) -> tuple[CommandResults, dict[str, Any]]:
    """
    Sending a message with conversation context to an OpenAI GPT model and retrieving the generated response.
    """
    demisto.debug("[Command gpt-send-message] triggered")
    message = args.get(ArgAndParamNames.MESSAGE, "")
    if not message:
        raise ValueError("Message not provided")

    completion_params = {
        ArgAndParamNames.MAX_TOKENS: args.get(ArgAndParamNames.MAX_TOKENS, None),
        ArgAndParamNames.TEMPERATURE: args.get(ArgAndParamNames.TEMPERATURE, None),
        ArgAndParamNames.TOP_P: args.get(ArgAndParamNames.TOP_P, None),
    }

    reset_conversation_history = args.get(ArgAndParamNames.RESET_CONVERSATION_HISTORY, "") == "yes"
    chat_context = get_chat_context(reset_conversation_history, message)
    demisto.debug(
        f"[Command gpt-send-message] Prepared chat | messages_count={len(chat_context)} | "
        f"completion_params_keys={sorted(completion_params.keys())} | "
        f"reset_history={reset_conversation_history}"
    )

    response = client.get_chat_completions(chat_context=chat_context, completion_params=completion_params)
    demisto.debug(
        f"[Command gpt-send-message] Got response | "
        f"choices_count={len(response.get('choices') or [])} | "
        f"usage_keys={sorted((response.get('usage') or {}).keys())}"
    )

    assistant_message = extract_assistant_message(response)
    conversation_step = [{Roles.USER: message, Roles.ASSISTANT: assistant_message}]

    usage: dict[str, str] = response.get("usage", {})

    readable_output = (
        assistant_message
        + "\n"
        + tableToMarkdown(
            name=f'{response.get(ArgAndParamNames.MODEL, "")} response:',
            sort_headers=False,
            t={
                "Prompt tokens": usage.get("prompt_tokens", ""),
                "Completion tokens": usage.get("completion_tokens", ""),
                "Total tokens": usage.get("total_tokens", ""),
                "Context messages": str(len(chat_context)),
            },
        )
    )
    return CommandResults(
        outputs_prefix="OpenAiChatGPTV3.Conversation",
        outputs=conversation_step,
        replace_existing=reset_conversation_history,
        readable_output=readable_output,
    ), response


def check_email_headers_command(client: OpenAiClient, args: dict[str, Any]) -> CommandResults:
    return check_email_part(EmailParts.HEADERS, client, args)


def check_email_body_command(client: OpenAiClient, args: dict[str, Any]) -> CommandResults:
    return check_email_part(EmailParts.BODY, client, args)


def analyze_email_header_command(client: OpenAiClient, args: dict[str, Any], params: dict[str, Any]) -> CommandResults:
    """Analyze email headers using the OpenAI Responses API.

    This is the Responses-API counterpart of ``gpt-check-email-header`` (which uses
    Chat Completions). It parses the uploaded ``.eml`` file, builds the same prompt
    template, and sends it via ``POST /v1/responses``.

    Two war-room entries are produced (matching the ``gpt-check-email-header`` UX):
      1. A table of the parsed email headers.
      2. The AI verdict followed by a token-usage table (with a *Reasoning tokens*
         row when a reasoning model is used).

    Args:
        client: The configured ``OpenAiClient`` instance.
        args: Command arguments from ``demisto.args()``.
        params: Instance parameters from ``demisto.params()``.

    Returns:
        A ``CommandResults`` with the AI analysis and usage info.
    """
    # --- Parse the .eml file ---
    entry_id: str = args.get(ArgAndParamNames.ENTRY_ID, "")
    email_headers, _, _, file_name = get_email_parts(entry_id)

    if not email_headers:
        raise DemistoException("'parse_emails' did not extract any email headers from the provided file.")

    # --- Build the prompt (identical template to gpt-check-email-header) ---
    additional_instructions = args.get(ArgAndParamNames.ADDITIONAL_INSTRUCTIONS, "")

    email_headers_formatted = {
        header["name"]: header["value"] for header in email_headers if "name" in header and "value" in header
    }
    readable_input = tableToMarkdown(name=f"{file_name} headers:", t=email_headers_formatted, sort_headers=False)
    prompt = CHECK_EMAIL_HEADERS_PROMPT.format(additional_instructions, readable_input)

    demisto.debug(f"[Analyze Email Header] Built prompt | length={len(prompt)} chars")

    # --- Resolve model & build the Responses API request body ---
    model = _resolve_model(args, client)
    body = _build_responses_api_body(args, params, model, prompt)

    # --- Call the Responses API ---
    response = client.create_response(body)

    # --- Extract the assistant's verdict ---
    assistant_message = extract_response_output_text(response)
    readable_output = _build_response_readable_output(response, assistant_message, model)

    # --- Emit the headers table as a first war-room entry (same UX as gpt-check-email-header) ---
    return_results(
        CommandResults(
            readable_output=readable_input,
            outputs_prefix="OpenAiChatGPTV3.EmailHeaders",
            outputs={"EmailHeaders": readable_input, "Response": response},
            replace_existing=True,
        )
    )

    # --- Return the AI verdict + usage table as the second war-room entry ---
    return CommandResults(
        outputs_prefix="OpenAiChatGPTV3.Response",
        outputs=[
            {
                Roles.USER: prompt,
                Roles.ASSISTANT: assistant_message,
                "response_id": response.get("id", ""),
            }
        ],
        replace_existing=True,
        readable_output=readable_output,
        raw_response=response,
    )


def analyze_email_body_command(client: OpenAiClient, args: dict[str, Any], params: dict[str, Any]) -> CommandResults:
    """Analyze email body for security risks using the OpenAI Responses API.

    This is the Responses-API counterpart of ``gpt-check-email-body`` (which uses
    Chat Completions). It parses the uploaded ``.eml`` file, builds the same prompt
    template, and sends it via ``POST /v1/responses``.

    Two war-room entries are produced (matching the ``gpt-check-email-body`` UX):
      1. A table of the parsed email body (text and HTML).
      2. The AI verdict followed by a token-usage table (with a *Reasoning tokens*
         row when a reasoning model is used).

    Args:
        client: The configured ``OpenAiClient`` instance.
        args: Command arguments from ``demisto.args()``.
        params: Instance parameters from ``demisto.params()``.

    Returns:
        A ``CommandResults`` with the AI analysis and usage info.
    """
    # --- Parse the .eml file ---
    entry_id: str = args.get(ArgAndParamNames.ENTRY_ID, "")
    _, email_text_body, email_html_body, file_name = get_email_parts(entry_id)

    if not email_text_body and not email_html_body:
        raise DemistoException("'email_parser' did not extract any email body from the provided file.")

    # --- Build the prompt (identical template to gpt-check-email-body) ---
    additional_instructions = args.get(ArgAndParamNames.ADDITIONAL_INSTRUCTIONS, "")

    email_text_body = email_text_body if email_text_body else ""
    email_html_body = email_html_body if email_html_body else ""
    email_body = {"Body/Text": email_text_body, "HTML/Text": email_html_body}

    readable_input = tableToMarkdown(name=f"{file_name} body:", t=email_body, sort_headers=False)
    prompt = CHECK_EMAIL_BODY_PROMPT.format(additional_instructions, readable_input)

    demisto.debug(f"[Analyze Email Body] Built prompt | length={len(prompt)} chars")

    # --- Resolve model & build the Responses API request body ---
    model = _resolve_model(args, client)
    body = _build_responses_api_body(args, params, model, prompt)

    # --- Call the Responses API ---
    response = client.create_response(body)

    # --- Extract the assistant's verdict ---
    assistant_message = extract_response_output_text(response)
    readable_output = _build_response_readable_output(response, assistant_message, model)

    # --- Emit the body table as a first war-room entry (same UX as gpt-check-email-body) ---
    return_results(
        CommandResults(
            readable_output=readable_input,
            outputs_prefix="OpenAiChatGPTV3.EmailBody",
            outputs={"EmailBody": readable_input, "Response": response},
            replace_existing=True,
        )
    )

    # --- Return the AI verdict + usage table as the second war-room entry ---
    return CommandResults(
        outputs_prefix="OpenAiChatGPTV3.Response",
        outputs=[
            {
                Roles.USER: prompt,
                Roles.ASSISTANT: assistant_message,
                "response_id": response.get("id", ""),
            }
        ],
        replace_existing=True,
        readable_output=readable_output,
        raw_response=response,
    )


def create_soc_email_template_command(client: OpenAiClient, args: dict[str, Any]) -> CommandResults:
    additional_instructions = (
        f"Additional instructions: {args.get(ArgAndParamNames.ADDITIONAL_INSTRUCTIONS)}\n"
        if args.get(ArgAndParamNames.ADDITIONAL_INSTRUCTIONS, "")
        else ""
    )
    create_soc_email_template_message = CREATE_SOC_EMAIL_TEMPLATE_PROMPT.format(additional_instructions)
    args.update({ArgAndParamNames.MESSAGE: create_soc_email_template_message})
    send_message_command_results, response = send_message_command(client, args)
    # Setting the SOCEmailTemplate context prior to returning the 'send-message-command' results
    # and setting the entire conversation in the context.
    return_results(
        CommandResults(outputs_prefix="OpenAiChatGPTV3.SocEmailTemplate", outputs={"Response": response}, replace_existing=True)
    )
    return send_message_command_results


def draft_soc_email_command(client: OpenAiClient, args: dict[str, Any], params: dict[str, Any]) -> CommandResults:
    """Draft a SOC email template using the OpenAI Responses API.

    This is the Responses-API counterpart of ``gpt-create-soc-email-template``
    (which uses Chat Completions). It uses the same prompt template and is
    designed to consume prior conversation context (e.g. from a preceding
    ``gpt-analyze-email-body`` call).

    There is no ``reset_conversation_history`` argument — this command consumes
    prior conversation context by design.

    Two war-room entries are produced (matching the ``gpt-create-soc-email-template`` UX):
      1. The SOC email template context output.
      2. The AI-generated template followed by a token-usage table (with a
         *Reasoning tokens* row when a reasoning model is used).

    Args:
        client: The configured ``OpenAiClient`` instance.
        args: Command arguments from ``demisto.args()``.
        params: Instance parameters from ``demisto.params()``.

    Returns:
        A ``CommandResults`` with the AI-generated SOC email template and usage info.
    """
    # --- Build the prompt (identical template to gpt-create-soc-email-template) ---
    additional_instructions = (
        f"Additional instructions: {args.get(ArgAndParamNames.ADDITIONAL_INSTRUCTIONS)}\n"
        if args.get(ArgAndParamNames.ADDITIONAL_INSTRUCTIONS, "")
        else ""
    )
    prompt = CREATE_SOC_EMAIL_TEMPLATE_PROMPT.format(additional_instructions)

    demisto.debug(f"[Draft SOC Email] Built prompt | length={len(prompt)} chars")

    # --- Resolve model & build the Responses API request body ---
    model = _resolve_model(args, client)
    body = _build_responses_api_body(args, params, model, prompt)

    # --- Call the Responses API ---
    response = client.create_response(body)

    # --- Extract the assistant's draft ---
    assistant_message = extract_response_output_text(response)
    readable_output = _build_response_readable_output(response, assistant_message, model)

    # --- Emit the SOC email template context as a first war-room entry (same UX as gpt-create-soc-email-template) ---
    return_results(
        CommandResults(
            outputs_prefix="OpenAiChatGPTV3.SocEmailTemplate",
            outputs={"Response": response},
            replace_existing=True,
        )
    )

    # --- Return the AI draft + usage table as the second war-room entry ---
    return CommandResults(
        outputs_prefix="OpenAiChatGPTV3.Response",
        outputs=[
            {
                Roles.USER: prompt,
                Roles.ASSISTANT: assistant_message,
                "response_id": response.get("id", ""),
            }
        ],
        replace_existing=True,
        readable_output=readable_output,
        raw_response=response,
    )


def extract_response_output_text(response: dict[str, Any]) -> str:
    """Extract the assistant's text from a Responses API response.

    The Responses API returns output as a list of message objects, each containing
    a list of content blocks. This extracts the first ``output_text`` block.

    Args:
        response: The parsed JSON response from ``POST /v1/responses``.

    Returns:
        The assistant's text content.

    Raises:
        DemistoException: If the expected structure is missing or empty.
    """
    output: list = response.get("output") or []
    if not output:
        raise DemistoException("Could not retrieve output from Responses API: 'output' field is empty or missing.")

    content: list = output[0].get("content") or []
    if not content:
        raise DemistoException("Could not retrieve output from Responses API: 'output[0].content' is empty or missing.")

    text: str = content[0].get("text") or ""
    if not text:
        raise DemistoException("Could not retrieve output from Responses API: 'output[0].content[0].text' is empty.")

    demisto.debug(f"[Responses] Extracted assistant output | length={len(text)} chars")
    return text


def _is_reasoning_model(model: str) -> bool:
    """Return ``True`` if *model* belongs to a reasoning-capable family.

    Reasoning models (o1, o3, o4-mini, gpt-5*) expose extra usage details
    such as ``usage.output_tokens_details.reasoning_tokens``.
    """
    return model.lower().startswith(ResponsesConfig.REASONING_MODEL_PREFIXES)


def _build_response_readable_output(
    response: dict[str, Any],
    assistant_message: str,
    model: str,
) -> str:
    """Build the human-readable output for a completed Responses API call.

    Mirrors the ``gpt-send-message`` HR style using ``tableToMarkdown``.
    For reasoning models, an extra *Reasoning tokens* row is included.

    Args:
        response: The full API response dict.
        assistant_message: The extracted assistant text.
        model: The resolved model name.

    Returns:
        Markdown-formatted readable output string.
    """
    usage: dict[str, Any] = response.get("usage") or {}
    response_id = response.get("id", "")
    actual_model = response.get("model", model)

    usage_table: dict[str, Any] = {
        "Input tokens": usage.get("input_tokens", ""),
        "Output tokens": usage.get("output_tokens", ""),
    }

    # For reasoning models, include the reasoning tokens row.
    if _is_reasoning_model(actual_model):
        output_details = usage.get("output_tokens_details") or {}
        usage_table["Reasoning tokens"] = output_details.get("reasoning_tokens", "")

    usage_table["Total tokens"] = usage.get("total_tokens", "")
    usage_table["Response ID"] = response_id

    return (
        assistant_message
        + "\n"
        + tableToMarkdown(
            name=f"{actual_model} response:",
            sort_headers=False,
            t=usage_table,
        )
    )


@polling_function(
    name="gpt-create-response",
    interval=ResponsesConfig.DEFAULT_POLLING_INTERVAL_SECS,
    timeout=ResponsesConfig.DEFAULT_POLLING_TIMEOUT_SECS,
    polling_arg_name="background",
)
def create_response_command(args: dict[str, Any], client: OpenAiClient, params: dict[str, Any]) -> PollResult:
    """Execute the ``gpt-create-response`` command using the OpenAI Responses API.

    Builds the request body from command arguments (falling back to instance params
    where applicable), calls the Responses API, manages conversation context via
    ``previous_response_id`` for multi-turn conversations, and returns structured
    ``CommandResults``.

    Conversation continuity:
        The Responses API natively supports multi-turn conversations through the
        ``previous_response_id`` field. After each call, the response ``id`` is stored
        in the XSOAR context at ``OpenAiChatGPTV3.Response``. On subsequent
        calls (when ``reset_conversation_history`` is ``False``), this ID is sent as
        ``previous_response_id`` so the model retains full conversation context
        server-side — no need to re-send prior messages.

    Polling:
        When ``background=true`` is set, the initial API call returns immediately with
        a ``queued`` or ``in_progress`` status. The ``@polling_function`` decorator
        handles re-scheduling automatically based on the returned ``PollResult``.

    Note:
        This command uses a **separate** context key (``Response``) from
        the existing ``gpt-send-message`` command (``Conversation``) to avoid clashing.

    Args:
        args: Command arguments from ``demisto.args()``.
        client: The configured ``OpenAiClient`` instance.
        params: Instance parameters from ``demisto.params()``.

    Returns:
        A ``PollResult`` indicating whether polling should continue or the final result.
    """
    # --- Polling re-entry: if we already have a response_id, poll for completion ---
    polling_response_id = args.get("_polling_response_id")
    if polling_response_id:
        demisto.debug(f"[Responses] Polling re-entry | response_id={polling_response_id}")
        response = client.get_response(polling_response_id)
        status = response.get("status", "")
        demisto.debug(f"[Responses] Poll status={status}")

        if status in ResponsesConfig.BACKGROUND_PENDING_STATUSES:
            # Still running — continue polling
            return PollResult(
                response=None,
                continue_to_poll=True,
                args_for_next_run=args,
                partial_result=CommandResults(
                    readable_output=f"⏳ Response `{polling_response_id}` is still {status}. Polling...",
                ),
            )

        if status == "failed":
            error_info = response.get("error") or {}
            raise DemistoException(
                f"Background response failed. "
                f"Error: {error_info.get('message', 'N/A')} (code={error_info.get('code', 'N/A')})"
            )

        if status == "incomplete":
            incomplete_details = response.get("incomplete_details") or {}
            raise DemistoException(f"Background response is incomplete. " f"Reason: {incomplete_details.get('reason', 'N/A')}")

        if status == "cancelled":
            raise DemistoException("Background response was cancelled.")

        if status != "completed":
            raise DemistoException(f"Background response reached unexpected status '{status}'. " f"Full response: {response}")

        # Completed — return the final result
        return PollResult(
            response=_build_completed_response_result(response, args),
            continue_to_poll=False,
            partial_result=CommandResults(readable_output="Response completed successfully."),
        )

    # --- First call: build and send the request ---
    message: str = args.get(ArgAndParamNames.MESSAGE, "")
    if not message:
        raise DemistoException("The 'message' argument is required.")

    reset_conversation_history: bool = argToBoolean(args.get(ArgAndParamNames.RESET_CONVERSATION_HISTORY, "no"))

    # Resolve model & build the base Responses API request body
    model = _resolve_model(args, client)
    body = _build_responses_api_body(args, params, model, message)

    # Conversation continuity via previous_response_id.
    # The Responses API keeps the full conversation server-side; we only need to
    # pass the last response ID to continue the thread.
    # Uses a SEPARATE context key (Response) to avoid clashing with
    # the Chat Completions-based gpt-send-message command (Conversation).
    if not reset_conversation_history:
        conversation_ctx = demisto.context().get("OpenAiChatGPTV3", {}).get("Response")
        previous_response_id: str | None = None
        if isinstance(conversation_ctx, list) and conversation_ctx:
            previous_response_id = conversation_ctx[-1].get("response_id") if isinstance(conversation_ctx[-1], dict) else None
        elif isinstance(conversation_ctx, dict):
            previous_response_id = conversation_ctx.get("response_id")

        if previous_response_id:
            body["previous_response_id"] = previous_response_id
            demisto.debug(f"[Responses] Continuing conversation | previous_response_id={previous_response_id}")
        else:
            demisto.debug("[Responses] No previous response ID found - starting new conversation.")
    else:
        demisto.debug("[Responses] Conversation history reset requested - starting fresh.")

    # background
    background = args.get("background")
    if background is not None:
        body["background"] = argToBoolean(background)

    # compact_threshold - validated but not sent to API; for future context compaction logic
    compact_threshold = args.get("compact_threshold")
    if compact_threshold is not None:
        compact_threshold_val = int(compact_threshold)
        if compact_threshold_val < 1000:
            raise DemistoException("compact_threshold must be at least 1000.")

    # Call the Responses API
    response = client.create_response(body)

    # If background=true and the response is not yet completed, start polling
    response_status = response.get("status", "")
    if body.get("background") and response_status in ResponsesConfig.BACKGROUND_PENDING_STATUSES:
        response_id = response.get("id", "")
        demisto.debug(f"[Responses] Background response queued | id={response_id} status={response_status}")
        return PollResult(
            response=None,
            continue_to_poll=True,
            args_for_next_run={**args, "_polling_response_id": response_id},
            partial_result=CommandResults(
                readable_output=f"⏳ Background response `{response_id}` is {response_status}. Polling started...",
            ),
        )

    # Synchronous completion — return the final result
    return PollResult(
        response=_build_completed_response_result(response, args),
        continue_to_poll=False,
    )


def _build_completed_response_result(response: dict[str, Any], args: dict[str, Any]) -> CommandResults:
    """Build the final ``CommandResults`` for a completed Responses API response.

    Shared by both the synchronous path and the polling completion path.

    Args:
        response: The completed API response dict.
        args: The original command arguments.

    Returns:
        A ``CommandResults`` with context, HR, and raw response.
    """
    message: str = args.get(ArgAndParamNames.MESSAGE) or ""
    model: str = args.get(ArgAndParamNames.MODEL) or response.get("model") or ""

    assistant_message = extract_response_output_text(response)
    response_id = response.get("id", "")

    # Store conversation step with response_id for multi-turn continuity
    conversation_step = [
        {
            Roles.USER: message,
            Roles.ASSISTANT: assistant_message,
            "response_id": response_id,
        }
    ]

    readable_output = _build_response_readable_output(response, assistant_message, model)

    return CommandResults(
        outputs_prefix="OpenAiChatGPTV3.Response",
        outputs=conversation_step,
        replace_existing=True,
        readable_output=readable_output,
        raw_response=response,
    )


# region Helpers - JSON parsing
# =================================
# Parsers for non-standard response shapes (concatenated JSON / JSONL)
# =================================
def parse_concatenated_json(body: str) -> list[dict[str, Any]]:
    """Parse a stream of concatenated JSON objects (and/or JSONL lines) into a list of dicts.

    The OpenAI Compliance log-content endpoint returns a body that is NOT valid JSON -
    objects are concatenated with optional whitespace/newlines between them. Uses
    `json.JSONDecoder().raw_decode()` to walk the buffer object-by-object. Non-dict
    top-level values are dropped because they cannot represent an event record.

    Args:
        body: The raw response body text.

    Returns:
        A list of dicts parsed from the body, or `[]` for empty/undecodable input.
    """
    if not body:
        demisto.debug("[Parse] Empty body received - returning [].")
        return []

    body_length = len(body)
    decoder = json.JSONDecoder()
    records: list[dict[str, Any]] = []
    skipped_non_dict = 0
    buffer = body.lstrip()
    while buffer:
        try:
            obj, end = decoder.raw_decode(buffer)
        except json.JSONDecodeError as exc:
            # Stop on first decode failure rather than silently truncate.
            demisto.error(f"[Parse] Failed to decode concatenated-JSON: {exc.msg} (so far: {len(records)} records).")
            break
        if isinstance(obj, dict):
            records.append(obj)
        else:
            skipped_non_dict += 1
        buffer = buffer[end:].lstrip()

    demisto.debug(
        f"[Parse] Decoded concatenated JSON | body_size={body_length} bytes | "
        f"records={len(records)} | skipped_non_dict={skipped_non_dict}"
    )
    return records


def _parse_json_or_concatenated(body: Any, log_prefix: str) -> Any:
    """Best-effort JSON parser for OpenAI listing responses.

    Tries strict JSON first; on `json.JSONDecodeError` (the production "Extra data..." /
    "Expecting value..." failures), falls back to concatenated-JSON / JSONL parsing so a
    multi-document response degrades gracefully to a list of records instead of crashing
    the whole stream. Empty/whitespace-only bodies return `[]` so callers see "no events".

    Pass-through: if `body` is already a `dict` or `list` (e.g. when a test mocks
    `_http_request` to return a parsed object directly), it is returned as-is.

    Args:
        body: The raw response body (str expected, but tolerates dict/list/None).
        log_prefix: Per-call-site tag for the debug/error logs (e.g. "[API Audit]").

    Returns:
        - The decoded value when the body is a single valid JSON document (typically dict or list).
        - A list of records when the body is concatenated JSON / JSONL (the fallback path).
        - `[]` when the body is empty / whitespace / un-decodable.
    """
    # Pass-through for already-parsed responses (e.g. from a mocked `_http_request`).
    if isinstance(body, dict | list):
        return body
    if body is None:
        demisto.debug(f"{log_prefix} None response body - returning [].")
        return []
    if not isinstance(body, str):
        body = str(body)
    stripped = body.strip()
    if not stripped:
        demisto.debug(f"{log_prefix} Empty response body - returning [].")
        return []
    try:
        return json.loads(stripped)
    except json.JSONDecodeError as exc:
        demisto.debug(
            f"{log_prefix} Strict JSON parse failed ({exc.msg} at pos {exc.pos}); "
            f"falling back to concatenated-JSON parser (body_size={len(body)})."
        )
        return parse_concatenated_json(body)


# endregion


# region Helpers - Integration Params
# =================================
# Integration parameter parsing & validation
# =================================
def _extract_credential(raw: Any) -> str:
    """Extract a password string from a Cortex credentials dict or pass through a plain string."""
    if isinstance(raw, dict):
        return raw.get("password", "") or ""
    if raw is None:
        return ""
    return str(raw)


def parse_event_types_to_fetch(raw_event_types: Any) -> list[str]:
    """Validate and normalize the selected event-type labels.

    Shared by `parse_integration_params` (save-time validation) and `parse_collector_params`
    (runtime stream selection) so the parsing lives in one place.
    """
    valid_labels = {EventType.AUDIT, *EVENT_TYPE_LABEL_TO_API.keys()}
    event_types_to_fetch = argToList(raw_event_types or [])
    invalid = [t for t in event_types_to_fetch if t not in valid_labels]
    if invalid:
        raise DemistoException(f"Invalid event type(s) selected: {invalid}. Valid options: {sorted(valid_labels)}")
    return event_types_to_fetch


def parse_integration_params(params: dict[str, Any]) -> dict[str, Any]:
    """Parse the connection/credential settings into a config dict for the `OpenAiClient`.

    Collector run-time params (event types, max-fetch, first-fetch) are parsed by
    `parse_collector_params`. The cred-vs-events correlation check runs here so
    misconfigured instances fail at save-time.
    """
    base_url = (params.get("url") or "https://api.openai.com/").rstrip("/") + "/"
    api_key = _extract_credential(params.get("apikey"))
    admin_api_key = _extract_credential(params.get("admin_api_key"))
    compliance_api_key = _extract_credential(params.get("compliance_api_key"))
    chatgpt_base_url = params.get("chatgpt_api_url") or Config.DEFAULT_CHATGPT_URL
    workspace_id = params.get("workspace_id") or ""
    model = params.get("model-freetext") or params.get("model-select") or ""
    verify = not argToBoolean(params.get("insecure", False))
    proxy = argToBoolean(params.get("proxy", False))

    event_types_to_fetch = parse_event_types_to_fetch(params.get("event_types_to_fetch"))
    validate_event_types_credentials_correlation(
        event_types_to_fetch=event_types_to_fetch,
        admin_api_key=admin_api_key,
        compliance_api_key=compliance_api_key,
    )

    demisto.debug(f"[Config] URL: {base_url} | ChatGPT URL: {chatgpt_base_url}")
    demisto.debug(f"[Config] Model: {model or '<none>'} | verify={verify} | proxy={proxy}")
    demisto.debug(
        f"[Config] Credentials present: chat={bool(api_key)} | "
        f"admin={bool(admin_api_key)} | compliance={bool(compliance_api_key)}"
    )
    demisto.debug(f"[Config] event_types_to_fetch={event_types_to_fetch or '<none>'}")

    return {
        "base_url": base_url,
        "api_key": api_key,
        "model": model,
        "verify": verify,
        "proxy": proxy,
        "admin_api_key": admin_api_key,
        "compliance_api_key": compliance_api_key,
        "chatgpt_base_url": chatgpt_base_url,
        "workspace_id": workspace_id,
    }


def validate_event_types_credentials_correlation(
    event_types_to_fetch: list[str],
    admin_api_key: str,
    compliance_api_key: str,
) -> None:
    """Validate that each selected event-type group has its required API key.

    The OpenAI Audit Logs stream uses the Admin API key, while every Compliance
    stream uses the Compliance API key. If the user selects a stream without
    providing the matching key, raise an informative `DemistoException`
    naming the exact selected types and the missing parameter.

    Args:
        event_types_to_fetch: User-facing event-type labels selected in the integration parameters.
        admin_api_key: The Admin API Key (may be empty).
        compliance_api_key: The Compliance API Key (may be empty).

    Raises:
        DemistoException: If a selected stream is missing its required API key.
    """
    audit_selected = EventType.AUDIT in event_types_to_fetch
    selected_compliance = sorted(label for label in event_types_to_fetch if label in EVENT_TYPE_LABEL_TO_API)
    demisto.debug(
        f"[Validation] Cross-check credentials | audit_selected={audit_selected} | "
        f"compliance_selected_count={len(selected_compliance)} | "
        f"admin_key_present={bool(admin_api_key)} | compliance_key_present={bool(compliance_api_key)}"
    )

    if audit_selected and not admin_api_key:
        raise DemistoException(
            f"'{EventType.AUDIT}' is selected in 'Events types to fetch', "
            "but no 'Admin API Key' is provided. The Admin API Key is required to fetch the OpenAI Audit logs. "
            "Either provide the Admin API Key or remove 'OpenAI Audit logs' from the selected event types."
        )

    if selected_compliance and not compliance_api_key:
        raise DemistoException(
            f"Compliance event type(s) {selected_compliance} are selected in 'Events types to fetch', "
            "but no 'Compliance API Key' is provided. The Compliance API Key is required to fetch the "
            "ChatGPT Compliance logs. "
            "Either provide the Compliance API Key or remove the Compliance event types from the selection."
        )


# endregion


# region Event Collector - Helpers
# =================================
# Helpers for the audit + compliance event collector
# =================================
def parse_first_fetch_to_datetime(first_fetch: str) -> datetime:
    """Parse a first-fetch string (e.g., '1 minute ago', '2024-01-01T00:00:00Z') into a UTC datetime.

    On parse failure, falls back to `Config.DEFAULT_FIRST_FETCH` so a typo never silently
    widens the lookback window beyond what the integration parameter advertises. If even
    the configured default is unparseable, a final sentinel of `now() - 1 minute` is used
    to guarantee the function returns a valid datetime.

    Callers format the result themselves:
      - Audit stream wants Unix seconds: `int(dt.timestamp())`
      - Compliance stream wants ISO 8601: `dt.replace(microsecond=0).strftime(Config.DATE_FORMAT)`
    """
    parsed: datetime | None = None
    try:
        parsed = arg_to_datetime(first_fetch, is_utc=True)
    except ValueError:
        demisto.error(
            f"[First Fetch] Failed to parse first_fetch='{first_fetch}' as a datetime "
            f"or relative time expression - falling back to default='{Config.DEFAULT_FIRST_FETCH}'. "
            f"Please correct the integration parameter."
        )
    if not parsed:
        # Honor the documented default rather than silently widening the lookback window.
        parsed = arg_to_datetime(Config.DEFAULT_FIRST_FETCH, is_utc=True) or (datetime.now(UTC) - timedelta(minutes=1))
    # `arg_to_datetime(is_utc=True)` returns a NAIVE datetime; force UTC so callers' `.timestamp()`
    # does not silently interpret it in the system's local timezone (a latent portability bug).
    if parsed.tzinfo is None:
        parsed = parsed.replace(tzinfo=UTC)
    demisto.debug(f"[First Fetch] Resolved '{first_fetch}' to {parsed.isoformat()}.")
    return parsed


def event_id(event: dict[str, Any]) -> str | None:
    """Extract a stable event identifier from a record.

    Tries common identifier keys in order of preference: `id`, `log_id`, `event_id`, `uuid`.
    The first non-empty value found is coerced to a string and returned.

    Args:
        event: An event/listing dict from any of the OpenAI API surfaces.

    Returns:
        The first identifier found as a string, or None if none of the known keys are present.
    """
    for key in ("id", "log_id", "event_id", "uuid"):
        value = event.get(key)
        if value:
            return str(value)
    return None


def deduplicate_events(events: list[dict[str, Any]], previous_ids: list[str]) -> list[dict[str, Any]]:
    """Filter out events whose identifier was ingested in a previous fetch cycle.

    Used by the Compliance stream's tie-dedup at the persisted `last_end_time` cursor.
    Audit dedup uses the API's native cursor and does not need this helper.

    Args:
        events: Candidate events from the current fetch.
        previous_ids: Identifiers that were already ingested last run.

    Returns:
        A new list containing only events whose `event_id` is not in `previous_ids`.
        Returns the input unchanged when either list is empty (fast path).
    """
    if not events or not previous_ids:
        return events
    previous_set = set(previous_ids)
    new_events = [e for e in events if event_id(e) not in previous_set]
    skipped = len(events) - len(new_events)
    if skipped:
        demisto.debug(f"[Dedup] Skipped {skipped} previously-seen events; {len(new_events)} remaining.")
    return new_events


def enrich_audit_event(event: dict[str, Any]) -> dict[str, Any]:
    """Add `_time` (derived from `effective_at`) to an Audit Logs event."""
    effective_at = event.get("effective_at")
    if isinstance(effective_at, int | float):
        event["_time"] = datetime.fromtimestamp(effective_at, tz=UTC).strftime(Config.DATE_FORMAT)
    else:
        demisto.debug("[Enrich Audit] Event missing 'effective_at' - _time not set.")
    return event


def enrich_compliance_event(event: dict[str, Any], api_event_type: str, workspace_id: str) -> dict[str, Any]:
    """Add `_time` (from `timestamp`), `source_log_type` (per `event_type`), and `workspace_id` to a Compliance event."""
    timestamp = event.get("timestamp")
    if timestamp:
        event["_time"] = timestamp
    else:
        demisto.debug("[Enrich Compliance] Event missing 'timestamp' - _time not set.")

    if api_event_type in COMPLIANCE_EVENT_TYPE_TO_SOURCE_LOG_TYPE:
        event["source_log_type"] = COMPLIANCE_EVENT_TYPE_TO_SOURCE_LOG_TYPE[api_event_type]
    else:
        demisto.info(f"[Enrich Compliance] Unknown event_type='{api_event_type}' - using lowercase as source_log_type.")
        event["source_log_type"] = api_event_type.lower()
    event["_event_type"] = api_event_type
    event["workspace_id"] = workspace_id
    return event


def selected_audit_enabled(event_types_to_fetch: list[str]) -> bool:
    """Return True if the user selected the Audit Logs stream in `event_types_to_fetch`."""
    return EventType.AUDIT in event_types_to_fetch


def selected_compliance_event_types(event_types_to_fetch: list[str]) -> list[str]:
    """Return the upstream `event_type` values for compliance streams the user selected."""
    return [EVENT_TYPE_LABEL_TO_API[label] for label in event_types_to_fetch if label in EVENT_TYPE_LABEL_TO_API]


@dataclass
class CollectorParams:
    """Parsed collector run-time parameters - the single shape both commands consume."""

    event_types_to_fetch: list[str]
    audit_selected: bool
    api_event_types: list[str]
    audit_max_fetch: int
    compliance_max_fetch: int
    workspace_id: str
    first_fetch: str

    @property
    def compliance_selected(self) -> bool:
        return bool(self.api_event_types)

    def streams_to_run(self) -> list[str]:
        """Return the ordered list of `Stream.*` identifiers this run should fetch."""
        out: list[str] = []
        if self.audit_selected:
            out.append(Stream.AUDIT)
        if self.compliance_selected:
            out.append(Stream.COMPLIANCE)
        return out


def parse_collector_params(
    params: dict[str, Any],
    args: dict[str, Any] | None = None,
) -> CollectorParams:
    """Parse the collector run-time params (event types, max-fetch, first-fetch, workspace_id).

    When `args` is passed (manual `openai-get-events`), `event_type` / `limit` / `start_time`
    override the corresponding integration params; `limit` controls both streams since the
    manual command exposes a single knob.
    """
    args = args or {}

    raw_event_types = args.get("event_type") if args.get("event_type") else params.get("event_types_to_fetch")
    event_types_to_fetch = parse_event_types_to_fetch(raw_event_types)

    if args:
        unified_limit = arg_to_number(args.get("limit")) or Config.DEFAULT_GET_EVENTS_LIMIT
        audit_max_fetch = unified_limit
        compliance_max_fetch = unified_limit
    else:
        audit_max_fetch = arg_to_number(params.get("audit_max_fetch")) or Config.DEFAULT_AUDIT_MAX_FETCH
        compliance_max_fetch = arg_to_number(params.get("compliance_max_fetch")) or Config.DEFAULT_COMPLIANCE_MAX_FETCH

    first_fetch = args.get("start_time") or Config.DEFAULT_FIRST_FETCH
    workspace_id = params.get("workspace_id") or ""

    return CollectorParams(
        event_types_to_fetch=event_types_to_fetch,
        audit_selected=selected_audit_enabled(event_types_to_fetch),
        api_event_types=selected_compliance_event_types(event_types_to_fetch),
        audit_max_fetch=audit_max_fetch,
        compliance_max_fetch=compliance_max_fetch,
        workspace_id=workspace_id,
        first_fetch=first_fetch,
    )


@dataclass
class FetchResult:
    """Per-stream fetch output: events, last_run updates, and target XSIAM product."""

    stream: str
    events: list[dict[str, Any]] = field(default_factory=list)
    last_run_updates: dict[str, Any] = field(default_factory=dict)
    product: str = ""


def fetch_stream(
    client: OpenAiClient,
    stream: str,
    last_run: dict[str, Any],
    collector_params: CollectorParams,
) -> FetchResult:
    """Unified per-stream fetch entry point shared by both commands.

    Dispatches to `fetch_audit_logs` or `fetch_compliance_logs` based on `stream`.
    """
    demisto.debug(
        f"[Fetch Stream] Dispatching stream='{stream}' | "
        f"audit_max_fetch={collector_params.audit_max_fetch} | "
        f"compliance_max_fetch={collector_params.compliance_max_fetch} | "
        f"first_fetch='{collector_params.first_fetch}' | "
        f"last_run_keys={list(last_run.keys())}"
    )
    if stream == Stream.AUDIT:
        events, updates = fetch_audit_logs(
            client=client,
            last_run=last_run,
            max_fetch=collector_params.audit_max_fetch,
            first_fetch=collector_params.first_fetch,
        )
        return FetchResult(
            stream=Stream.AUDIT,
            events=events,
            last_run_updates=updates,
            product=Config.PRODUCT_AUDIT,
        )

    if stream == Stream.COMPLIANCE:
        if not collector_params.workspace_id:
            demisto.debug(
                "[Fetch Stream] Compliance stream skipped - no workspace_id configured " "(required by /v1/compliance/.../logs)."
            )
            return FetchResult(stream=Stream.COMPLIANCE, product=Config.PRODUCT_COMPLIANCE)
        events, updates = fetch_compliance_logs(
            client=client,
            workspace_id=collector_params.workspace_id,
            api_event_types=collector_params.api_event_types,
            last_run=last_run,
            max_fetch=collector_params.compliance_max_fetch,
            first_fetch=collector_params.first_fetch,
        )
        return FetchResult(
            stream=Stream.COMPLIANCE,
            events=events,
            last_run_updates=updates,
            product=Config.PRODUCT_COMPLIANCE,
        )

    raise DemistoException(f"Unknown stream identifier '{stream}'.")


# endregion


# region Event Collector - Fetch logic
# =================================
# Fetch logic for Audit & Compliance streams
# =================================
def fetch_audit_logs(
    client: OpenAiClient,
    last_run: dict[str, Any],
    max_fetch: int,
    first_fetch: str,
) -> tuple[list[dict[str, Any]], dict[str, Any]]:
    """Fetch the next batch of OpenAI Audit logs (cursor-based dedup).

    Audit dedup is delegated to the API's native cursor: every response carries `last_id`,
    which is persisted as `audit_after` and replayed verbatim as `after=` on the next run.
    No timestamp HWM, no per-id dedup list - the cursor guarantees no duplicates.

    The very first run (no stored cursor) is seeded by `first_fetch` (e.g. "1 day"), which
    is converted to a Unix-second `effective_at[gt]` lower bound for that single request.

    Returns:
        A tuple of (enriched_events, last_run_updates).
    """
    demisto.debug(f"[Audit Fetch] Starting | max_fetch={max_fetch} | first_fetch='{first_fetch}'")

    stored_cursor: str | None = last_run.get(LastRunKey.AUDIT_AFTER)
    stored_seed: int | None = last_run.get(LastRunKey.AUDIT_FIRST_FETCH_SEED)
    initial_effective_at_gt: int | None = None
    if stored_cursor:
        demisto.debug("[Audit Fetch] Resuming from stored cursor.")
    elif stored_seed is not None:
        # Replay seed persisted by a previous empty first-fetch (keeps the window anchored).
        initial_effective_at_gt = stored_seed
        demisto.debug(f"[Audit Fetch] Replaying persisted first-fetch seed effective_at>{initial_effective_at_gt}.")
    else:
        initial_effective_at_gt = int(parse_first_fetch_to_datetime(first_fetch).timestamp())
        demisto.debug(f"[Audit Fetch] No cursor in last_run - first fetch using effective_at>{initial_effective_at_gt}.")

    collected: list[dict[str, Any]] = []
    after: str | None = stored_cursor
    last_cursor: str | None = stored_cursor
    pages = 0
    while len(collected) < max_fetch and pages < Config.MAX_PAGES_PER_FETCH:
        # On every request after the first, `after=` is the cursor; `effective_at_gt` is only
        # honored on the very first ever request (when no cursor has been persisted yet).
        response = client.get_audit_logs(
            after=after,
            effective_at_gt=initial_effective_at_gt if after is None else None,
        )
        page = response.get("data") or []
        if not page:
            demisto.debug(f"[Audit Fetch] Page {pages + 1}: empty - stopping pagination.")
            break
        collected.extend(page)
        pages += 1
        page_last_id = response.get("last_id")
        if page_last_id:
            last_cursor = page_last_id
        demisto.debug(
            f"[Audit Fetch] Page {pages}: +{len(page)} events | total_collected={len(collected)} | "
            f"max_fetch={max_fetch} | has_more={response.get('has_more')} | new_cursor_set={bool(page_last_id)}"
        )
        if not response.get("has_more"):
            demisto.debug(f"[Audit Fetch] Page {pages}: has_more=false - stopping pagination.")
            break
        if not page_last_id:
            demisto.debug(f"[Audit Fetch] Page {pages}: no last_id cursor returned - stopping pagination.")
            break
        after = page_last_id

    # When over `max_fetch`, trim and advance the persisted cursor to the LAST kept event's id
    # so the next run resumes precisely after this batch (no gap, no overlap).
    if len(collected) > max_fetch:
        demisto.debug(f"[Audit Fetch] Trimming {len(collected)} collected events down to max_fetch={max_fetch}.")
        collected = collected[:max_fetch]
        last_cursor = event_id(collected[-1]) or last_cursor

    # Enrich with `_time` / `source_log_type` (cursor pagination guarantees no duplicates already).
    for event in collected:
        enrich_audit_event(event)

    # Persist the latest cursor so the next run picks up exactly after this batch.
    # On an empty first-ever fetch (no cursor + no stored seed), persist the resolved seed so
    # the next cycle replays the same lookback rather than sliding it forward with `now()`.
    last_run_updates: dict[str, Any] = {}
    if last_cursor:
        last_run_updates[LastRunKey.AUDIT_AFTER] = last_cursor
        last_run_updates[LastRunKey.AUDIT_FIRST_FETCH_SEED] = None  # cursor wins; drop stale seed.
        demisto.debug("[Audit Fetch] Persisting new cursor for next run.")
    elif initial_effective_at_gt is not None and stored_seed is None:
        last_run_updates[LastRunKey.AUDIT_FIRST_FETCH_SEED] = initial_effective_at_gt
        demisto.debug(f"[Audit Fetch] First fetch empty - persisting seed={initial_effective_at_gt}.")

    demisto.debug(
        f"[Audit Fetch] Done | new_events={len(collected)} | pages_fetched={pages} | " f"updates={list(last_run_updates.keys())}"
    )
    return collected, last_run_updates


def fetch_compliance_logs(
    client: OpenAiClient,
    workspace_id: str,
    api_event_types: list[str],
    last_run: dict[str, Any],
    max_fetch: int,
    first_fetch: str,
) -> tuple[list[dict[str, Any]], dict[str, Any]]:
    """Fetch the next batch of OpenAI Compliance logs (two-step list + content per id).

    Compliance dedup is timestamp-based with a per-timestamp ID set:
      - The listing API returns `last_end_time` directly in its response - we use that as the
        `after=` value on the next run (and persist it as `compliance_last_end_time`).
      - At any given `last_end_time`, multiple listings may share that exact timestamp.
        The IDs of those listings are persisted as `compliance_last_ids`, so on the next run
        any listing whose id is in that set is filtered out (preventing tie-duplicates).

    Returns:
        A tuple of (enriched_events, last_run_updates).
    """
    demisto.debug(
        f"[Compliance Fetch] Starting | event_types_count={len(api_event_types)} | "
        f"max_fetch={max_fetch} | first_fetch='{first_fetch}'"
    )

    cursor: str = last_run.get(LastRunKey.COMPLIANCE_LAST_END_TIME) or (
        parse_first_fetch_to_datetime(first_fetch).replace(microsecond=0).strftime(Config.DATE_FORMAT)
    )
    previous_ids: list[str] = list(last_run.get(LastRunKey.COMPLIANCE_LAST_IDS) or [])
    demisto.debug(f"[Compliance Fetch] Resolved cursor | cursor_set={bool(cursor)} | prev_ids_count={len(previous_ids)}")

    listing_response = client.list_compliance_logs(
        workspace_id=workspace_id,
        event_types=api_event_types,
        after=cursor,
        limit=max_fetch,
    )
    listings: list[dict[str, Any]] = listing_response.get("data") or []
    response_last_end_time: str | None = listing_response.get("last_end_time")

    # Empty listing - the API may still have advanced its cursor; persist it if so.
    if not listings:
        demisto.debug("[Compliance Fetch] No new compliance log entries returned.")
        updates: dict[str, Any] = {}
        if response_last_end_time and response_last_end_time != cursor:
            updates[LastRunKey.COMPLIANCE_LAST_END_TIME] = response_last_end_time
            updates[LastRunKey.COMPLIANCE_LAST_IDS] = []
        return [], updates

    if len(listings) > max_fetch:
        demisto.debug(f"[Compliance Fetch] Trimming {len(listings)} listings down to max_fetch={max_fetch}.")
        listings = listings[:max_fetch]

    # Dedupe against IDs that were already seen at the persisted `last_end_time` (tie-dedup).
    new_listings = deduplicate_events(listings, previous_ids)
    demisto.debug(f"[Compliance Fetch] Listings ready | listings_total={len(listings)} | new_listings={len(new_listings)}")

    # Step 2: for each new listing, fetch its content payload.
    events: list[dict[str, Any]] = []
    failed_content_fetches = 0
    for listing in new_listings:
        log_id = event_id(listing)
        api_event_type = listing.get("event_type", "")
        if not log_id:
            demisto.debug("[Compliance Fetch] Skipping listing entry with missing id.")
            continue
        try:
            content = client.get_compliance_log_content(workspace_id=workspace_id, log_id=log_id)
        except Exception as exc:
            failed_content_fetches += 1
            demisto.error(
                f"[Compliance Fetch] Failed to fetch content for log_id={log_id} " f"(api_event_type={api_event_type}): {exc}"
            )
            continue

        # The content endpoint returns a list of records (parsed from a concatenated-JSON / JSONL body).
        # Carry forward listing metadata so each event is self-describing downstream.
        for record in content:
            record.setdefault("id", log_id)
            record.setdefault("end_time", listing.get("end_time"))
            enrich_compliance_event(record, api_event_type, workspace_id)
            events.append(record)

    if failed_content_fetches:
        demisto.error(f"[Compliance Fetch] {failed_content_fetches} content fetch(es) failed and were skipped.")

    # Persist the API-reported `last_end_time` and the IDs of listings sharing that exact timestamp.
    # Falls back to the max `end_time` across listings if the API omits `last_end_time`.
    last_run_updates: dict[str, Any] = {}
    new_end_time: str | None = response_last_end_time or max(
        (et for listing in listings if (et := listing.get("end_time"))), default=None
    )
    if new_end_time:
        ids_at_end_time = [eid for listing in listings if listing.get("end_time") == new_end_time and (eid := event_id(listing))]
        # If the cursor didn't move, merge with previously-seen IDs to keep the dedup set complete.
        if new_end_time == cursor:
            ids_at_end_time = list(set(previous_ids) | set(ids_at_end_time))
        last_run_updates[LastRunKey.COMPLIANCE_LAST_END_TIME] = new_end_time
        last_run_updates[LastRunKey.COMPLIANCE_LAST_IDS] = ids_at_end_time
        demisto.debug(f"[Compliance Fetch] New cursor last_end_time advanced | ids_at_end_time_count={len(ids_at_end_time)}")

    demisto.debug(
        f"[Compliance Fetch] Done | events={len(events)} | listings_processed={len(new_listings)} | "
        f"updates={list(last_run_updates.keys())}"
    )
    return events, last_run_updates


# endregion


# region Event Collector - Commands
# =================================
# fetch-events / openai-get-events commands
# =================================
def _push_result(client: OpenAiClient, result: FetchResult, log_prefix: str) -> None:
    """Push a single stream's events to its dataset, isolating errors per stream."""
    if not result.events:
        demisto.debug(f"{log_prefix} No {result.stream} events to push.")
        return
    demisto.debug(f"{log_prefix} Pushing {len(result.events)} {result.stream} event(s) to product='{result.product}'.")
    try:
        client.send_events(result.events, product=result.product)
    except Exception as exc:
        demisto.error(f"{log_prefix} Failed to push {result.stream} events: {exc}")
        return
    demisto.debug(f"{log_prefix} Push of {len(result.events)} {result.stream} event(s) completed.")


def fetch_events_command(client: OpenAiClient, params: dict[str, Any]) -> None:
    """Scheduled fetch: pulls Audit + Compliance events in parallel and sends them to Cortex ingestion."""
    demisto.debug("[Command fetch-events] triggered")

    collector_params = parse_collector_params(params)
    last_run = demisto.getLastRun() or {}
    demisto.debug(f"[Command fetch-events] getLastRun returned successfully | last_run={last_run}")

    demisto.debug(
        f"[Command fetch-events] event_types_count={len(collector_params.event_types_to_fetch)} | "
        f"audit_max_fetch={collector_params.audit_max_fetch} | "
        f"compliance_max_fetch={collector_params.compliance_max_fetch} | "
        f"last_run_keys={list(last_run.keys())}"
    )

    streams_to_run = collector_params.streams_to_run()
    if not streams_to_run:
        demisto.debug("[Command fetch-events] No event-type group selected. Nothing to fetch.")
        demisto.setLastRun(last_run)
        return

    demisto.debug(f"[Command fetch-events] Launching {len(streams_to_run)} stream(s) in parallel: {streams_to_run}")

    # Each thread gets its own copy of last_run; main thread merges after `as_completed`.
    updated_last_run: dict[str, Any] = dict(last_run)
    results: list[FetchResult] = []
    stream_failures: dict[str, str] = {}  # surfaced to the UI after partial success is persisted.

    futures: dict[Future, str] = {}
    with ThreadPoolExecutor(max_workers=len(streams_to_run)) as executor:
        for stream in streams_to_run:
            demisto.debug(f"[Command fetch-events] Submitting {stream} stream to executor.")
            futures[
                executor.submit(
                    fetch_stream,
                    client=client,
                    stream=stream,
                    last_run=dict(last_run),
                    collector_params=collector_params,
                )
            ] = stream

        for future in as_completed(futures):
            stream_name = futures[future]
            demisto.debug(f"[Command fetch-events] Future completed for stream='{stream_name}' - collecting result.")
            try:
                result = future.result()
            except Exception as exc:
                # Keep iterating so successful streams still push events and persist their cursor.
                demisto.error(f"[Command fetch-events] {stream_name} stream failed: {exc}")
                stream_failures[stream_name] = str(exc)
                continue
            results.append(result)
            updated_last_run.update(result.last_run_updates)
            demisto.debug(
                f"[Command fetch-events] {stream_name} stream produced {len(result.events)} events | "
                f"last_run_updates={result.last_run_updates}"
            )

    demisto.debug(f"[Command fetch-events] All streams done. Pushing results for {len(results)} stream(s)...")
    for result in results:
        _push_result(client, result, log_prefix="[Command fetch-events]")

    demisto.debug(f"[Command fetch-events] Persisting merged last_run | last_run={updated_last_run}")
    demisto.setLastRun(updated_last_run)
    demisto.debug(f"[Command fetch-events] setLastRun persisted successfully | last_run={updated_last_run}")
    demisto.debug(
        f"[Command fetch-events] done | "
        f"streams={[r.stream for r in results]} | "
        f"events_per_stream={[len(r.events) for r in results]} | "
        f"last_run={updated_last_run} | "
        f"failed_streams={list(stream_failures.keys())}"
    )

    # Re-raise after partial progress is persisted so `main()` shows a visible UI error.
    if stream_failures:
        summary = "; ".join(f"{name}: {err}" for name, err in stream_failures.items())
        raise DemistoException(f"One or more event-collector streams failed: {summary}")


def get_events_command(client: OpenAiClient, args: dict[str, Any], params: dict[str, Any]) -> CommandResults:
    """Manual `openai-get-events` command. Runs the same fetch pipeline without persisting `last_run`."""
    demisto.debug("[Command openai-get-events] triggered")

    collector_params = parse_collector_params(params, args=args)
    should_push_events = argToBoolean(args.get("should_push_events", False))  # noqa: F405

    demisto.debug(
        f"[Command openai-get-events] event_type_count={len(collector_params.event_types_to_fetch)} | "
        f"limit={collector_params.audit_max_fetch} | should_push_events={should_push_events} | "
        f"first_fetch='{collector_params.first_fetch}'"
    )

    streams_to_run = collector_params.streams_to_run()
    if not streams_to_run:
        demisto.debug("[Command openai-get-events] No event-type group selected. Returning empty result.")
    else:
        demisto.debug(f"[Command openai-get-events] Fetching {len(streams_to_run)} stream(s) sequentially: {streams_to_run}")

    results: list[FetchResult] = []
    for stream in streams_to_run:
        demisto.debug(f"[Command openai-get-events] Fetching {stream} stream...")
        try:
            result = fetch_stream(client=client, stream=stream, last_run={}, collector_params=collector_params)
        except Exception as exc:
            # Isolate per-stream failures so one stream's exception does not abort the entire command.
            demisto.error(f"[Command openai-get-events] {stream} stream failed: {exc}")
            continue
        results.append(result)
        demisto.debug(f"[Command openai-get-events] {stream} stream returned {len(result.events)} event(s).")

    all_events: list[dict[str, Any]] = [event for result in results for event in result.events]
    demisto.debug(
        f"[Command openai-get-events] Returning {len(all_events)} event(s) "
        f"({', '.join(f'{r.stream}={len(r.events)}' for r in results) or 'no streams'}, "
        f"push={should_push_events})."
    )

    if should_push_events:
        for result in results:
            _push_result(client, result, log_prefix="[Command openai-get-events]")

    readable_output = tableToMarkdown(
        "OpenAI GPT Events",
        all_events,
        headers=["id", "_event_type", "source_log_type", "_time"],
        removeNull=True,
    )
    return CommandResults(
        readable_output=readable_output,
        outputs_prefix="OpenAI.Event",
        outputs_key_field="id",
        outputs=all_events,
    )


def validate_create_moderation_args(args: dict[str, Any]) -> None:
    """Validate that exactly one of text, entry_id, or image_url is provided.

    Raises:
        DemistoException: When zero or more than one input source is supplied.
    """
    text = argToList(args.get("text"))
    entry_id = args.get("entry_id")
    image_url = args.get("image_url")

    provided = sum(bool(v) for v in (text, entry_id, image_url))
    if provided == 0:
        raise DemistoException("Exactly one of 'text', 'entry_id', or 'image_url' must be provided.")
    if provided > 1:
        raise DemistoException("Only one of 'text', 'entry_id', or 'image_url' may be provided at a time.")


def _entry_id_to_data_url(entry_id: str) -> str:
    """Read a war-room image file and return a base64 data-URL string.

    Args:
        entry_id: The war-room entry ID referencing an uploaded image file.

    Returns:
        A ``data:<mime>;base64,<payload>`` string suitable for the Moderations API.

    Raises:
        DemistoException: When the file cannot be found or is not an image type.
    """
    file_result = demisto.getFilePath(entry_id)
    if not file_result:
        raise DemistoException(f"Could not find file for entry_id '{entry_id}'.")

    file_path: str = file_result.get("path", "")
    file_name: str = file_result.get("name", "")

    if not file_path:
        raise DemistoException(f"File path is empty for entry_id '{entry_id}'.")

    mime_type, _ = mimetypes.guess_type(file_name)
    if mime_type is None or not mime_type.startswith("image/"):
        raise DemistoException(
            f"Unsupported or unknown image type for file '{file_name}' "
            f"(detected MIME: {mime_type}). Only image files are supported."
        )

    with open(file_path, "rb") as f:
        image_b64 = base64.b64encode(f.read()).decode("utf-8")

    demisto.debug(
        f"[Moderation] Encoded image from entry_id '{entry_id}' | "
        f"file={file_name} | mime={mime_type} | b64_len={len(image_b64)}"
    )
    return f"data:{mime_type};base64,{image_b64}"


def create_moderation_command(client: OpenAiClient, args: dict[str, Any]) -> CommandResults:
    """Run content through the OpenAI Moderations API and return per-category results.

    Supports text (array / comma-separated), a war-room image entry, or a public image URL.
    When multiple texts are provided the API returns one result per text; this function
    iterates over all of them and produces a separate table and context entry for each.
    """
    validate_create_moderation_args(args)

    model = args.get("model", "omni-moderation-latest")
    text = argToList(args.get("text"))
    entry_id = args.get("entry_id")
    image_url = args.get("image_url")

    # Build the ``input`` field based on which argument was supplied and
    # track the input type/value for context output.
    if text:
        moderation_input: list[Any] = text
        input_type = "text"
    elif entry_id:
        data_url = _entry_id_to_data_url(entry_id)
        moderation_input = [{"type": "image_url", "image_url": {"url": data_url}}]
        input_type = "image"
    else:
        # image_url
        moderation_input = [{"type": "image_url", "image_url": {"url": image_url}}]
        input_type = "image_url"

    body: dict[str, Any] = {"model": model, "input": moderation_input}
    response = client.create_moderation(body)

    results_list: list[dict[str, Any]] = response.get("results", [])
    if not results_list:
        input_value = text[0] if text else (entry_id or image_url or "")
        return CommandResults(
            readable_output="No moderation results returned.",
            outputs_prefix="OpenAiChatGPTV3.Moderation",
            outputs=[
                {
                    "Flagged": False,
                    "Categories": {},
                    "CategoryScores": {},
                    "Input": {"input_type": input_type, "input_value": input_value},
                }
            ],
            replace_existing=True,
        )

    # Determine labels for each result.  For text inputs the label is the text
    # itself; for images we fall back to a generic "Image" label.
    input_labels: list[str] = text if text else ["Image"] * len(results_list)

    all_outputs: list[dict[str, Any]] = []
    readable_parts: list[str] = []

    for idx, result in enumerate(results_list):
        flagged: bool = result.get("flagged", False)
        categories: dict[str, bool] = result.get("categories", {})
        category_scores: dict[str, float] = result.get("category_scores", {})

        label = input_labels[idx] if idx < len(input_labels) else f"Input {idx + 1}"

        # Build human-readable table: one row per category.
        table_rows = [
            {
                "Category": cat,
                "Flagged": "" if categories.get(cat, False) else "",
                "Score": f"{category_scores.get(cat, 0.0):.4f}",
            }
            for cat in categories
        ]

        readable_parts.append(
            tableToMarkdown(
                f"Moderation Results for \"{label}\" (Flagged: {'Yes' if flagged else 'No'})",
                table_rows,
                headers=["Category", "Flagged", "Score"],
            )
        )

        # For text input each result corresponds to a text from the list;
        # for image/image_url there is always a single result.
        input_value = text[idx] if text and idx < len(text) else (entry_id or image_url or "")

        all_outputs.append(
            {
                "Input": {"input_type": input_type, "input_value": input_value},
                "Flagged": flagged,
                "Categories": categories,
                "CategoryScores": category_scores,
            }
        )

    readable_output = "\n".join(readable_parts)

    # When there is only a single result, unwrap the list for backward compatibility.
    outputs: list[dict[str, Any]] | dict[str, Any] = all_outputs[0] if len(all_outputs) == 1 else all_outputs

    return CommandResults(
        readable_output=readable_output, outputs_prefix="OpenAiChatGPTV3.Moderation", outputs=outputs, replace_existing=True
    )


def list_models_command(client: OpenAiClient) -> CommandResults:
    """List all models visible to the configured API key.

    Calls ``GET /v1/models`` and returns a table of Id / Created / OwnedBy.
    """
    response = client.list_models()
    models_data: list[dict[str, Any]] = response.get("data", [])

    outputs = [
        {
            "Id": m.get("id", ""),
            "Created": timestamp_to_datestring(m.get("created", 0) * 1000) if m.get("created") else "",
            "OwnedBy": m.get("owned_by", ""),
        }
        for m in models_data
    ]

    readable_output = tableToMarkdown(
        "OpenAI Models",
        outputs,
        headers=["Id", "Created", "OwnedBy"],
    )

    return CommandResults(
        readable_output=readable_output,
        outputs_prefix="OpenAiChatGPTV3.Model",
        outputs_key_field="Id",
        outputs=outputs,
    )


# endregion


# region Main
# =================================
# Main entry point
# =================================

# Maps command name -> handler with a uniform (client, args, params) signature.
COMMAND_MAP: dict[str, Callable[["OpenAiClient", dict[str, Any], dict[str, Any]], Any]] = {
    "test-module": lambda client, args, params: test_module(client=client, params=params),
    "gpt-send-message": lambda client, args, params: send_message_command(client=client, args=args)[0],
    "gpt-check-email-header": lambda client, args, params: check_email_headers_command(client=client, args=args),
    "gpt-check-email-body": lambda client, args, params: check_email_body_command(client=client, args=args),
    "gpt-create-soc-email-template": lambda client, args, params: create_soc_email_template_command(client=client, args=args),
    "gpt-draft-soc-email": lambda client, args, params: draft_soc_email_command(client=client, args=args, params=params),
    "gpt-analyze-email-header": lambda client, args, params: analyze_email_header_command(
        client=client, args=args, params=params
    ),
    "gpt-analyze-email-body": lambda client, args, params: analyze_email_body_command(client=client, args=args, params=params),
    "gpt-create-response": lambda client, args, params: create_response_command(args=args, client=client, params=params),
    "gpt-list-models": lambda client, args, params: list_models_command(client=client),
    "gpt-create-moderation": lambda client, args, params: create_moderation_command(client=client, args=args),
    "fetch-events": lambda client, args, params: fetch_events_command(client=client, params=params),
    "openai-get-events": lambda client, args, params: get_events_command(client=client, args=args, params=params),
}


def main() -> None:  # pragma: no cover
    """Main entry point.

    Parses integration params (via `parse_integration_params`), validates the requested
    command against the `COMMAND_MAP`, builds the `OpenAiClient`, and dispatches.
    Errors are logged with a full traceback and surfaced via `return_error`.
    """
    # Make `demisto.*` runtime-bridge calls thread-safe for `fetch_events_command`'s
    # ThreadPoolExecutor workers. Official CommonServerPython helper, idempotent.
    support_multithreading()  # noqa: F405
    demisto.debug(f"{INTEGRATION_NAME} integration started")

    try:
        params = demisto.params()
        args = demisto.args()
        command = demisto.command()

        handler = COMMAND_MAP.get(command)
        if handler is None:
            raise NotImplementedError(
                f"Command '{command}' is not implemented in the OpenAI GPT integration. "
                f"Available commands: {sorted(COMMAND_MAP.keys())}"
            )
        demisto.debug(f"[Main] Resolved handler for command '{command}'.")

        # Backfill GPT chat args from instance params when missing.
        setup_args(args, params)

        config = parse_integration_params(params)

        client = OpenAiClient(
            url=config["base_url"],
            api_key=config["api_key"],
            model=config["model"],
            verify=config["verify"],
            proxy=config["proxy"],
            admin_api_key=config["admin_api_key"],
            compliance_api_key=config["compliance_api_key"],
            chatgpt_base_url=config["chatgpt_base_url"],
        )

        demisto.debug("[Main] Client built. Dispatching command...")
        result = handler(client, args, params)
        if result is not None:
            return_results(result)

        demisto.debug(f"[Main] Command '{command}' completed successfully.")

    except Exception as error:
        error_msg = str(error)
        demisto.error(f"[Main] Command '{command}' failed: {error_msg}")
        demisto.error(traceback.format_exc())
        return_error(f"Failed to execute '{command}' command. Error: {error_msg}")

    demisto.debug(f"{INTEGRATION_NAME} integration finished")


# endregion


""" ENTRY POINT """
if __name__ in ("__main__", "__builtin__", "builtins"):
    main()