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
| ID | OpenAi ChatGPT v3 |
|---|---|
| Provider | OpenAI |
| Category | Messaging and Conferencing |
| From Version | 6.0.0 |
| Docker Image | demisto/parse-emails:0.1.48.10120494 |
| Supported Modules | Agentix XSIAM |
README
OpenAI GPT
Instance Configuration
-
Generate an API Key
- Sign up or log in to OpenAI developer platform.
- Generate a new API key at OpenAI developer platform - api-keys.
-
Choose a GPT model to interact with
-
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).
-
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)
- max-tokens: The maximum number of tokens that can be generated for the response. (Allows controlling tokens’ consumption). Default: unset.
- 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.
- 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
- Admin API Key (required for
OpenAI Audit logs): generate from the OpenAI Platform admin console. Used to call/v1/organization/audit_logs. - Compliance API Key (required for any Compliance event type): generate from the ChatGPT Platform. Used to call
/v1/compliance/workspaces/{workspace_id}/.... - Workspace ID (required for any Compliance event type): the UUID of the compliance workspace whose events you want to collect.
- Admin API Key (required for
-
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 openaichatgpt_auditopenai_chatgpt_audit_rawCompliance logs (all) openaichatgpt_complianceopenai_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.comBase 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:
- A table of the parsed email headers.
- 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:
- A table of the parsed email body (text and HTML).
- 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:
- The SOC email template context output (
replace_existing=True— running twice overwrites the previous draft). - 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 URLapikey—admin_api_key— Admin API Keycompliance_api_key— Compliance API Keyworkspace_id— Workspace IDmodel-select— Modelmodel-freetext— Model (Optional - overrides selected choice)max_tokens— Max tokenstemperature— Temperaturetop_p— Top Pinsecure— Trust any certificate (not secure)proxy— Use system proxy settingsisFetchEvents— Fetch eventsevent_types_to_fetch— Events types to fetchaudit_max_fetch— Maximum number of OpenAI Audit events per fetchcompliance_max_fetch— Maximum number of Compliance events per fetcheventFetchInterval— Events Fetch Interval
Commands (11)
-
gpt-analyze-email-bodyAnalyzes 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-headerAnalyzes 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-bodyChecks 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-headerChecking 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-moderationRuns 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-responseSends 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-templateCreate an email template out of the conversation context to be sent from the SOC.
-
gpt-draft-soc-emailDrafts 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-modelsLists 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-messageSend a plain message to the selected GPT model and receive the generated response.
-
openai-get-eventsManually 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()