diff --git a/.github/workflows/build-appinspect.yml b/.github/workflows/build-appinspect.yml index 2dea3c7..2bd7771 100644 --- a/.github/workflows/build-appinspect.yml +++ b/.github/workflows/build-appinspect.yml @@ -17,10 +17,11 @@ jobs: app_name: ${{ steps.app.outputs.name }} app_version: ${{ steps.app.outputs.version }} steps: - - uses: actions/checkout@v4 - - uses: actions/setup-python@v5 + - uses: actions/checkout@v6 + - uses: actions/setup-python@v6 with: python-version: 3.9 + cache: 'pip' - name: Install dependencies run: | python -m pip install --upgrade pip @@ -55,7 +56,7 @@ jobs: ucc-gen package --path output/${{ steps.app.outputs.name }} shell: bash - - uses: actions/upload-artifact@v4 + - uses: actions/upload-artifact@v6 with: name: ${{ steps.app.outputs.name }}-${{ steps.app.outputs.version }} path: ${{ steps.app.outputs.name }}*.tar.gz @@ -73,11 +74,12 @@ jobs: - "private_classic" - "private_victoria" steps: - - uses: actions/checkout@v4 - - uses: actions/setup-python@v5 + - uses: actions/checkout@v6 + - uses: actions/setup-python@v6 with: python-version: 3.9 - - uses: actions/download-artifact@v4 + cache: 'pip' + - uses: actions/download-artifact@v5 with: name: ${{ needs.build.outputs.app_name }}-${{ needs.build.outputs.app_version }} path: dist @@ -85,7 +87,7 @@ jobs: with: app_path: dist/${{ needs.build.outputs.app_name }}-${{ needs.build.outputs.app_version }}.tar.gz included_tags: ${{ matrix.tags }} - - uses: actions/upload-artifact@v4 + - uses: actions/upload-artifact@v6 with: name: appinspect_result_${{ matrix.tags }} path: appinspect_result.json diff --git a/.github/workflows/docs.yml b/.github/workflows/docs.yml index c5c1ba1..695c11e 100644 --- a/.github/workflows/docs.yml +++ b/.github/workflows/docs.yml @@ -15,9 +15,10 @@ jobs: contents: write pages: write steps: - - uses: actions/checkout@v4 - - uses: actions/setup-python@v5 + - uses: actions/checkout@v6 + - uses: actions/setup-python@v6 with: python-version: 3.12 + cache: 'pip' - run: pip install mkdocs==1.6.0 mkdocs-material==9.5.32 mkdocs-print-site-plugin==2.6.0 - run: mkdocs gh-deploy --force \ No newline at end of file diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 63c8aa9..ddb2d25 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -22,10 +22,11 @@ jobs: artifact_name: ${{ steps.artifact.outputs.name }} ta_version: ${{ steps.app.outputs.version }} steps: - - uses: actions/checkout@v4 - - uses: actions/setup-python@v5 + - uses: actions/checkout@v6 + - uses: actions/setup-python@v6 with: python-version: ${{ matrix.python-version }} + cache: 'pip' - name: Install dependencies run: | python -m pip install --upgrade pip @@ -64,7 +65,7 @@ jobs: run: | echo "name=${{ steps.app.outputs.name }}-${{ steps.app.outputs.version }}-py${{ matrix.python-version }}" >> $GITHUB_OUTPUT - - uses: actions/upload-artifact@v4 + - uses: actions/upload-artifact@v6 with: name: ${{ steps.artifact.outputs.name }} path: ${{ steps.app.outputs.name }}*.tar.gz @@ -74,8 +75,8 @@ jobs: runs-on: ubuntu-latest needs: build steps: - - uses: actions/checkout@v4 - - uses: actions/download-artifact@v4 + - uses: actions/checkout@v6 + - uses: actions/download-artifact@v5 with: name: ${{ needs.build.outputs.artifact_name }} path: dist diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 329c383..3810427 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -8,6 +8,7 @@ on: paths: - package/bin/** - package/lib/requirements.txt + - tests/** concurrency: group: ${{ github.workflow }}-${{ github.ref }} @@ -22,10 +23,11 @@ jobs: matrix: version: [10.2.0] steps: - - uses: actions/checkout@v4 - - uses: actions/setup-python@v5 + - uses: actions/checkout@v6 + - uses: actions/setup-python@v6 with: python-version: 3.9 + cache: 'pip' - name: Install dependencies run: | python -m pip install --upgrade pip @@ -79,7 +81,10 @@ jobs: -p 8089:8089 \ splunk/splunk:${{ matrix.version }} # Wait for splunk to be up and running (~2min) - sleep 150 + for i in {1..60}; do + docker logs splunk 2>&1 | grep -q "Ansible playbook complete" && break + sleep 5 + done value=$(docker exec splunk printenv GENESYSCLOUD_HOST) echo "GENESYSCLOUD_HOST=$value" >> $GITHUB_ENV @@ -135,12 +140,12 @@ jobs: - name: Upload Mockoon logs if: always() - uses: actions/upload-artifact@v4 + uses: actions/upload-artifact@v6 with: name: mockoon-logs-${{ matrix.version }} path: logs/mockoon-${{ matrix.version }}.log - - uses: dorny/test-reporter@v2 + - uses: dorny/test-reporter@v3 if: always() with: name: Tests Results @@ -154,12 +159,13 @@ jobs: strategy: fail-fast: false matrix: - version: [9.4.1, 9.3.3, 9.2.5] + version: [9.4.1, 9.3.3] steps: - - uses: actions/checkout@v4 - - uses: actions/setup-python@v5 + - uses: actions/checkout@v6 + - uses: actions/setup-python@v6 with: python-version: 3.9 + cache: 'pip' - name: Install dependencies run: | python -m pip install --upgrade pip @@ -212,7 +218,10 @@ jobs: -p 8089:8089 \ splunk/splunk:${{ matrix.version }} # Wait for splunk to be up and running (~2min) - sleep 150 + for i in {1..60}; do + docker logs splunk 2>&1 | grep -q "Ansible playbook complete" && break + sleep 5 + done value=$(docker exec splunk printenv GENESYSCLOUD_HOST) echo "GENESYSCLOUD_HOST=$value" >> $GITHUB_ENV @@ -263,12 +272,12 @@ jobs: - name: Upload Mockoon logs if: always() - uses: actions/upload-artifact@v4 + uses: actions/upload-artifact@v6 with: name: mockoon-logs-${{ matrix.version }} path: logs/mockoon-${{ matrix.version }}.log - - uses: dorny/test-reporter@v2 + - uses: dorny/test-reporter@v3 if: always() with: name: Tests Results diff --git a/Makefile b/Makefile index f9d8bc9..222a0ce 100644 --- a/Makefile +++ b/Makefile @@ -7,7 +7,7 @@ venv: python3 -m venv .venv build: venv - source .venv/bin/activate; + source .venv/bin/activate && \ ucc-gen build --ta-version=$(APP_VERSION) run: @@ -21,14 +21,14 @@ package: build ucc-gen package --path output/$(APP_NAME) -o dist install-docs: venv - source .venv/bin/activate; + source .venv/bin/activate && \ pip install mkdocs==1.6.0 mkdocs-material==9.5.32 mkdocs-print-site-plugin==2.6.0 run-docs: install-docs mkdocs serve install-tests: venv - source .venv/bin/activate; + source .venv/bin/activate && \ pip install pytest==6.2.4 splunk-sdk run-tests: install-tests @@ -36,5 +36,5 @@ run-tests: install-tests export GENESYSCLOUD_HOST="http://localhost:3004" && python -m pytest integration/* run-functional-tests: install-tests - cd tests; - python -m pytest tests/modinput_functional/* + cd tests && \ + python -m pytest modinput_functional/* diff --git a/docs/ConfigureAnalyticsInputs/index.md b/docs/ConfigureAnalyticsInputs/index.md index dc9b950..567a649 100644 --- a/docs/ConfigureAnalyticsInputs/index.md +++ b/docs/ConfigureAnalyticsInputs/index.md @@ -91,6 +91,6 @@ Each attribute in the following table corresponds to a field in Splunk Web. |`interval` |Interval (seconds) |Rerun the input after the defined value, in seconds. The default value is 300.| | `direction` |Direction |The direction of the communication. | `media_types` |Media Type(s) |The session media type(s). -| `start_date` |Start Date |Date from which start collecting data. The default value is 7 days ago from now. Format: `YYYY-MM-DD`. +| `start_date` |Start Date |Date from which start collecting data. The default value is 7 days from now. Format: `YYYY-MM-DD`. Direction and Media Type(s) possible values are taken from [Genesys Cloud Specs](https://developer.genesys.cloud/analyticsdatamanagement/analytics/aggregate/conversation-query#dimensions). \ No newline at end of file diff --git a/docs/ConfigureOperationalInputs/index.md b/docs/ConfigureOperationalInputs/index.md index 2f34bcd..f324424 100644 --- a/docs/ConfigureOperationalInputs/index.md +++ b/docs/ConfigureOperationalInputs/index.md @@ -51,6 +51,7 @@ index = interval = poll_interval_seconds = max_poll_attempts = +start_date = ``` @@ -75,3 +76,4 @@ Each attribute in the following table corresponds to a field in Splunk Web. |`interval` |Interval (seconds) |Rerun the input after the defined value, in seconds. The default value is 300.| |`max_poll_attempts` |Max Poll Attempts |Maximum number of status checks (polls) performed for an audit query transaction before giving up. The default value is 10, to be increased for long-running queries.| |`poll_interval_seconds` |Poll Interval (seconds) |Seconds to wait between each status check of the audit query `transaction_id`. Lower values give quicker response but perform more API calls; higher values reduce rate-limit pressure. The default value is 2.| +| `start_date` |Start Date |Date from which start collecting data. The default value is 7 days from now. Format: `YYYY-MM-DD`. diff --git a/docs/ProxyConfiguration/index.md b/docs/ProxyConfiguration/index.md new file mode 100644 index 0000000..e989100 --- /dev/null +++ b/docs/ProxyConfiguration/index.md @@ -0,0 +1,12 @@ +# Configure Proxy for the Genesys Cloud Add-on for Splunk + +Perform the following steps to configure a Proxy for the Genesys Cloud Add-on for Splunk. + +1. On the Splunk Web home page, click **Genesys Cloud Add-on for Splunk** in the left navigation bar. +2. Click **Configuration** in the app navigation bar. +3. Click the **Proxy** tab. +4. Select the **Enable** box to enable the proxy connection and fill in the fields required for your proxy. +5. Click **Save**. + +To disable your proxy, uncheck the +**Enable** box. diff --git a/globalConfig.json b/globalConfig.json index 4a28bc8..c6d2187 100644 --- a/globalConfig.json +++ b/globalConfig.json @@ -141,6 +141,98 @@ ], "title": "Accounts" }, + { + "name": "proxy", + "title": "Proxy", + "entity": [ + { + "field": "proxy_enabled", + "label": "Enable", + "type": "checkbox" + }, + { + "field": "proxy_type", + "label": "Proxy Type", + "type": "singleSelect", + "options": { + "disableSearch": true, + "autoCompleteFields": [ + { + "label": "http", + "value": "http" + }, + { + "value": "https", + "label": "https" + }, + { + "label": "socks4", + "value": "socks4" + }, + { + "label": "socks5", + "value": "socks5" + } + ] + }, + "defaultValue": "http" + }, + { + "field": "proxy_url", + "type": "text", + "label": "Host", + "validators": [ + { + "type": "string", + "errorMsg": "Max host length is 4096", + "minLength": 1, + "maxLength": 4096 + } + ] + }, + { + "field": "proxy_port", + "label": "Port", + "type": "text", + "validators": [ + { + "type": "number", + "range": [ + 1, + 65535 + ] + } + ] + }, + { + "field": "proxy_username", + "label": "Username", + "type": "text", + "validators": [ + { + "type": "string", + "minLength": 0, + "maxLength": 50, + "errorMsg": "Max length of username is 50" + } + ] + }, + { + "field": "proxy_password", + "label": "Password", + "type": "text", + "encrypted": true, + "validators": [ + { + "type": "string", + "minLength": 0, + "maxLength": 8192, + "errorMsg": "Max length of password is 8192" + } + ] + } + ] + }, { "type": "loggingTab" }, @@ -265,6 +357,26 @@ "referenceName": "account" } }, + { + "type": "singleSelect", + "field": "index", + "label": "Index", + "defaultValue": "default", + "required": true, + "options": { + "endpointUrl": "data/indexes", + "denyList": "^_.*$", + "createSearchChoice": true + }, + "validators": [ + { + "type": "string", + "minLength": 1, + "maxLength": 80, + "errorMsg": "Length of index name should be between 1 and 80." + } + ] + }, { "type": "text", "field": "max_poll_attempts", @@ -282,24 +394,19 @@ "help": "Seconds to wait between status checks (polls) of the audit query transaction_id. Lower values respond faster but use more API calls; higher values reduce rate-limit pressure." }, { - "type": "singleSelect", - "field": "index", - "label": "Index", - "defaultValue": "default", - "required": true, - "options": { - "endpointUrl": "data/indexes", - "denyList": "^_.*$", - "createSearchChoice": true - }, + "field": "start_date", + "label": "Start Date", + "help": "Start date for data collection. Default: 7 days ago. Format: YYYY-MM-DD", + "required": false, + "type": "text", "validators": [ { - "type": "string", - "minLength": 1, - "maxLength": 80, - "errorMsg": "Length of index name should be between 1 and 80." + "type": "date" } - ] + ], + "options": { + "disableonEdit": true + } } ], "inputHelperModule": "audit_query_helper", @@ -1074,7 +1181,7 @@ "restRoot": "genesys_cloud_ta", "version": "0.3.1", "displayName": "Genesys Cloud Add-on for Splunk", - "schemaVersion": "0.0.9", + "schemaVersion": "0.0.10", "supportedThemes": [ "light", "dark" diff --git a/mkdocs.yml b/mkdocs.yml index 6fecc5a..11fe454 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -61,6 +61,7 @@ nav: - Configure Operational Inputs: 'ConfigureOperationalInputs/index.md' - Configure Telephony Providers Edge Inputs: 'ConfigureTelephonyProvidersEdgeInputs/index.md' - Configure Users Inputs: 'ConfigureUsersInputs/index.md' + - Proxy Configuration: 'ProxyConfiguration/index.md' - Amazon EventBridge Integration: 'IntegrateEventBridge/index.md' - Troubleshoot: - Troubleshoot: 'Troubleshooting/index.md' diff --git a/package/bin/actions_metrics_helper.py b/package/bin/actions_metrics_helper.py index 5fb4802..14e4faf 100644 --- a/package/bin/actions_metrics_helper.py +++ b/package/bin/actions_metrics_helper.py @@ -6,7 +6,7 @@ from solnlib.modular_input import checkpointer from splunklib import modularinput as smi -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from genesyscloud_client import GenesysCloudClient ADDON_NAME = "genesys_cloud_ta" @@ -25,13 +25,47 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password + def validate_input(definition: smi.ValidationDefinition): """Validation function for the modular input (currently unused).""" return def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): - from datetime import datetime, timedelta, timezone - for input_name, input_item in inputs.inputs.items(): normalized_input_name = input_name.split("/")[-1] logger = logger_for_input(normalized_input_name) @@ -54,11 +88,12 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): account = input_item.get("account") logger.info(f"Retrieving credentials for account: {account}") + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) account_region = get_account_property(session_key, account, "region") client_id = get_account_property(session_key, account, "client_id") client_secret = get_account_property(session_key, account, "client_secret") - client = GenesysCloudClient(logger, client_id, client_secret, account_region) + client = GenesysCloudClient(logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password) checkpointer_key_name = normalized_input_name now = datetime.now(timezone.utc) @@ -87,47 +122,49 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): "ActionAggregationQuery", body ) + sourcetype="genesyscloud:analytics:actions:metrics" + + if response: + event_counter = 0 + res_dict = response.to_dict() or {} + to_process_data = res_dict.get("results") or [] + + for event in to_process_data: + try: + group_info = event.get("group", {}) or {} + for data_entry in event.get("data", []): + interval_str = data_entry.get("interval") + interval_start_time = ( + datetime.strptime(interval_str.split("/")[0], "%Y-%m-%dT%H:%M:%S.%fZ").timestamp() + if interval_str else round(start_time.timestamp(), 3) + ) - event_counter = 0 - res_dict = response.to_dict() or {} - to_process_data = res_dict.get("results") or [] - - for event in to_process_data: - try: - group_info = event.get("group", {}) or {} - for data_entry in event.get("data", []): - interval_str = data_entry.get("interval") - interval_start_time = ( - datetime.strptime(interval_str.split("/")[0], "%Y-%m-%dT%H:%M:%S.%fZ").timestamp() - if interval_str else round(start_time.timestamp(), 3) - ) - - for metric in data_entry.get("metrics", []): - enriched_metric = { - "metric": metric.get("metric"), - "qualifier": metric.get("qualifier"), - "stats": metric.get("stats", {}), - "group": group_info, - "interval": interval_str - } - - event_writer.write_event( - smi.Event( - data=json.dumps(enriched_metric, ensure_ascii=False, default=str), - index=input_item.get("index"), - sourcetype="genesyscloud:analytics:actions:metrics", - time=interval_start_time + for metric in data_entry.get("metrics", []): + enriched_metric = { + "metric": metric.get("metric"), + "qualifier": metric.get("qualifier"), + "stats": metric.get("stats", {}), + "group": group_info, + "interval": interval_str + } + + event_writer.write_event( + smi.Event( + data=json.dumps(enriched_metric, ensure_ascii=False, default=str), + index=input_item.get("index"), + sourcetype=sourcetype, + time=interval_start_time + ) ) - ) - event_counter += 1 - except Exception as e: - logger.error(f"Failed to write event: {e}") + event_counter += 1 + except Exception as e: + logger.error(f"Failed to write event: {e}") - if event_counter > 0: - new_checkpoint = end_time.timestamp() - kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) + if event_counter > 0: + new_checkpoint = end_time.timestamp() + kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) - log.events_ingested(logger, input_name, "genesyscloud:analytics:actions:metrics", event_counter, input_item.get("index"), account=account) + log.events_ingested(logger, input_name, sourcetype, event_counter, input_item.get("index"), account=account) log.modular_input_end(logger, normalized_input_name) except Exception as e: diff --git a/package/bin/audit_query_helper.py b/package/bin/audit_query_helper.py index ba5b229..8ca5ea2 100644 --- a/package/bin/audit_query_helper.py +++ b/package/bin/audit_query_helper.py @@ -8,7 +8,6 @@ from splunklib import modularinput as smi from datetime import datetime, timedelta, timezone, time as dt_time -from dateutil.relativedelta import relativedelta from genesyscloud_client import GenesysCloudClient ADDON_NAME = "genesys_cloud_ta" @@ -27,10 +26,66 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf = cfm.get_conf("genesys_cloud_ta_account") return account_conf.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password + +def exceed_range(start: str, end: str, max_days: int = 31, + fmt: str = "%Y-%m-%dT%H:%M:%SZ") -> bool: + """Verify interval range does not exceed threshold. + :param start: Interval start string reference. + :param end: Interval end string reference. + :param max_days: Threshold in number of days. + :param fmt: Datetime string format. + :return: True if the range between start and end exceeds max_days. + """ + start_dt = datetime.strptime(start, fmt) + end_dt = datetime.strptime(end, fmt) + + if end_dt <= start_dt: + raise Exception(f"Invalid start date: {start} cannot be later or equal to {end}") + + return (end_dt - start_dt) >= timedelta(days=max_days) def validate_input(definition: smi.ValidationDefinition): - pass - + start_date = definition.parameters.get("start_date") + fmt_str = "%Y-%m-%d" + if start_date is not None: + if exceed_range(start_date, datetime.now().strftime(fmt_str), fmt=fmt_str): + raise Exception(f"Invalid start date {start_date}. Start date for data collection cannot exceed 31 days from now.") + return def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): for input_name, input_item in inputs.inputs.items(): @@ -61,21 +116,35 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_id = get_account_property(session_key, input_item.get("account"), "client_id") client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") - client = GenesysCloudClient(logger, client_id, client_secret, account_region) + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) + client = GenesysCloudClient(logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password) checkpointer_key_name = normalized_input_name + # Setting a default start date of 7 days ago from now now = datetime.now(timezone.utc) - fallback_start = (now - relativedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") - + fallback_start = (now - timedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") start_date = input_item.get("start_date") if start_date is not None: fallback_start = datetime.strptime(start_date, "%Y-%m-%d").strftime("%Y-%m-%dT%H:%M:%SZ") - start_time = kvstore_checkpointer.get(checkpointer_key_name) or fallback_start + start_time = kvstore_checkpointer.get(checkpointer_key_name) end_time = now.strftime("%Y-%m-%dT%H:%M:%SZ") + + # Evaluating the interval. It can be no greater than 31 days. + # [400] BadRequest - You must specify a search interval as part of your query that does not exceed 31 days. + if start_time is None or (start_time and exceed_range(start_time, end_time)): + logger.info(f"Checkpoint not found or exceeds interval range of 31 days [{start_time}]. Using the configured start_date as fallback [{fallback_start}].") + start_time = fallback_start + if exceed_range(fallback_start, end_time): + # This case keeps the system running in case of too far away in time checkpoint. + reset_start = (now - timedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") + logger.warn(f"Fallback start_date exceeds interval range of 31 days. Resetting it to {reset_start}.") + start_time = reset_start + interval = f"{start_time}/{end_time}" - body = {"interval": interval} + + body = { "interval": interval } logger.debug(f"Request body: {body}") response = client.post( @@ -94,11 +163,14 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): status = "" for _ in range(max_polls): - state_resp = client.get("AuditApi","get_audits_query_transaction_id", transaction_id) + state_resp = client.get("AuditApi", "get_audits_query_transaction_id", transaction_id) if not state_resp: raise Exception(f"Failed to get status for transaction {transaction_id}") status = state_resp[0].state + # Allowed statuses: + # https://github.com/MyPureCloud/platform-client-sdk-python/blob/master/build/PureCloudPlatformClientV2/models/audit_query_execution_status_response.py#L126 if status in ("Succeeded", "Failed", "Cancelled"): + logger.debug(f"Audit complete with status '{status}'") break time.sleep(poll_sleep) @@ -107,6 +179,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): params = { "transaction_id": transaction_id, + "allow_redirect": True, "page_size": 500 } @@ -116,13 +189,22 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): **params ) event_counter = 0 + sourcetype="genesyscloud:operational:audits" + for entity in results: + if not isinstance(entity, dict): + value = entity.to_dict() + time_value = entity.event_date.timestamp() + else: + # Retrieved via download URL file + value = entity + time_value = datetime.strptime(entity["eventTime"], "%Y-%m-%dT%H:%M:%SZ").timestamp() event_writer.write_event( smi.Event( - data=json.dumps(entity.to_dict(), ensure_ascii=False, default=str), + data=json.dumps(value, ensure_ascii=False, default=str), index=input_item.get("index"), - sourcetype="genesyscloud:operational:audits", - time=entity.event_date.timestamp() + sourcetype=sourcetype, + time=time_value ) ) event_counter += 1 @@ -133,6 +215,16 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): logger.debug(f"Updating checkpointer to {end_time}") kvstore_checkpointer.update(checkpointer_key_name, end_time) + log.events_ingested( + logger, + input_name, + sourcetype, + event_counter, + input_item.get("index"), + account=input_item.get("account"), + ) + + log.modular_input_end(logger, normalized_input_name) except Exception as e: log.log_exception( logger, diff --git a/package/bin/conversations_details_helper.py b/package/bin/conversations_details_helper.py index c578c30..5525d2e 100644 --- a/package/bin/conversations_details_helper.py +++ b/package/bin/conversations_details_helper.py @@ -6,8 +6,7 @@ from solnlib.modular_input import checkpointer from splunklib import modularinput as smi -from datetime import datetime -from dateutil.relativedelta import relativedelta +from datetime import datetime, timezone, timedelta from genesyscloud_client import GenesysCloudClient @@ -25,6 +24,42 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password + def get_conversation_duration(start: datetime, end: datetime) -> int: """ Calculate conversation duration. @@ -37,14 +72,29 @@ def get_conversation_duration(start: datetime, end: datetime) -> int: duration = end - start return int(duration.total_seconds() * 1000) +def exceed_range(start: str, end: str, max_days: int = 31, + fmt: str = "%Y-%m-%dT%H:%M:%SZ") -> bool: + """Verify interval range does not exceed threshold. + :param start: Interval start string reference. + :param end: Interval end string reference. + :param max_days: Threshold in number of days. + :param fmt: Datetime string format. + :return: True if the range between start and end exceeds max_days. + """ + start_dt = datetime.strptime(start, fmt) + end_dt = datetime.strptime(end, fmt) + + if end_dt <= start_dt: + raise Exception(f"Invalid start date: {start} cannot be later or equal to {end}") + + return (end_dt - start_dt) >= timedelta(days=max_days) + def validate_input(definition: smi.ValidationDefinition): - # Interval can be no greater than 7 days. start_date = definition.parameters.get("start_date") + fmt_str = "%Y-%m-%d" if start_date is not None: - input_date = datetime.strptime(start_date, "%Y-%m-%d") - threshold_date = datetime.now() - relativedelta(days=7) - if input_date < threshold_date: - raise Exception(f"Invalid start date {input_date}. Start date for data collection can be no greater than 7 days ago from now.") + if exceed_range(start_date, datetime.now().strftime(fmt_str), fmt=fmt_str): + raise Exception(f"Invalid start date {start_date}. Start date for data collection cannot exceed 31 days from now.") return def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): @@ -70,23 +120,33 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_id = get_account_property(session_key, input_item.get("account"), "client_id") client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") # Setting a default start date of 7 days ago from now - now = datetime.now() - fallback_start = (now - relativedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") + now = datetime.now(timezone.utc) + fallback_start = (now - timedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") start_date = input_item.get("start_date") if start_date is not None: fallback_start = datetime.strptime(start_date, "%Y-%m-%d").strftime("%Y-%m-%dT%H:%M:%SZ") + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) client = GenesysCloudClient( - logger, client_id, client_secret, account_region + logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password ) - checkpointer_key_name = input_name.split("/")[-1] + checkpointer_key_name = normalized_input_name - # Retrieve the last checkpoint or set it to the fallback start date. - start_time = ( - kvstore_checkpointer.get(checkpointer_key_name) - or fallback_start - ) + # Retrieve the last checkpoint + start_time = kvstore_checkpointer.get(checkpointer_key_name) end_time = now.strftime("%Y-%m-%dT%H:%M:%SZ") + + # Evaluating the interval. It can be no greater than 31 days. + # [400] BadRequest - You must specify a search interval as part of your query that does not exceed 31 days. + if start_time is None or (start_time and exceed_range(start_time, end_time)): + logger.info(f"Checkpoint not found or exceeds interval range of 31 days [{start_time}]. Using the configured start_date as fallback [{fallback_start}].") + start_time = fallback_start + if exceed_range(fallback_start, end_time): + # This case keeps the system running in case of too far away in time checkpoint. + reset_start = (now - timedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") + logger.warn(f"Fallback start_date exceeds interval range of 31 days. Resetting it to {reset_start}.") + start_time = reset_start + interval = f"{start_time}/{end_time}" body = { diff --git a/package/bin/conversations_metrics_helper.py b/package/bin/conversations_metrics_helper.py index 4ac3c2d..575ce9c 100644 --- a/package/bin/conversations_metrics_helper.py +++ b/package/bin/conversations_metrics_helper.py @@ -25,6 +25,42 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password + def validate_input(definition: smi.ValidationDefinition): """Validation function for the modular input (currently unused).""" return @@ -55,8 +91,8 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): account_region = get_account_property(session_key, input_item.get("account"), "region") client_id = get_account_property(session_key, input_item.get("account"), "client_id") client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") - - client = GenesysCloudClient(logger, client_id, client_secret, account_region) + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) + client = GenesysCloudClient(logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password) checkpointer_key_name = normalized_input_name current_checkpoint = ( @@ -73,14 +109,13 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): "nConsult", "nConsultTransferred", "nError", "nOffered", "nOutbound", "nOutboundAbandoned", "nOutboundAttempted", "nOutboundConnected", "nOverSla", "nStateTransitionError", "nTransferred", "oExternalMediaCount", "oMediaCount", - "oMessageTurn", "oServiceLevel", - "oServiceTarget", "tAbandon", "tAcd", "tActiveCallback", "tActiveCallbackComplete", + "oMessageTurn", "oServiceLevel", "oServiceTarget", + "tWait", "tAbandon", "tAcd", "tActiveCallback", "tActiveCallbackComplete", "tAcw", "tAgentResponseTime", "tAlert", "tAnswered", "tBarging", "tCoaching", "tCoachingComplete", "tConnected", "tContacting", "tDialing", "tFirstConnect", "tFirstDial", "tFlowOut", "tHandle", "tHeld", "tHeldComplete", "tIvr", "tMonitoring", "tMonitoringComplete", "tNotResponding", "tPark", "tParkComplete", - "tShortAbandon", "tTalk", "tTalkComplete", "tUserResponseTime", "tVoicemail", - "tWait", "nOffered" + "tShortAbandon", "tTalk", "tTalkComplete", "tUserResponseTime", "tVoicemail" ] group_by = ["queueId"] @@ -123,12 +158,12 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): datetime.strptime(data_entry["interval"].split("/")[0], "%Y-%m-%dT%H:%M:%S.%fZ").timestamp() if event.get("data") else round(start_time.timestamp(), 3) ) - for metrics in data_entry["metrics"]: - metrics["group"] = event["group"] - metrics["interval"] = data_entry["interval"] + for metric in data_entry["metrics"]: + metric["group"] = event["group"] + metric["interval"] = data_entry["interval"] event_writer.write_event( smi.Event( - data=json.dumps(metrics, ensure_ascii=False, default=str), + data=json.dumps(metric, ensure_ascii=False, default=str), index=input_item.get("index"), sourcetype=sourcetype, time=interval_start_time diff --git a/package/bin/edges_metrics_helper.py b/package/bin/edges_metrics_helper.py index 10dc99c..a704673 100644 --- a/package/bin/edges_metrics_helper.py +++ b/package/bin/edges_metrics_helper.py @@ -7,7 +7,7 @@ from solnlib.modular_input import checkpointer from splunklib import modularinput as smi -from datetime import datetime +from datetime import datetime, timezone from genesyscloud_client import GenesysCloudClient from genesyscloud_models import EdgeModel @@ -27,6 +27,41 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password def validate_input(definition: smi.ValidationDefinition): return @@ -67,8 +102,9 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") # aws_region = input_item.get('region') + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) client = GenesysCloudClient( - logger, client_id, client_secret, account_region + logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password ) checkpointer_key_name = input_name.split("/")[-1] @@ -117,7 +153,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): # Updating checkpoint if data was indexed to avoid losing info if event_counter > 0: logger.debug(f"Indexed '{event_counter}' events") - new_checkpoint = datetime.utcnow().timestamp() + new_checkpoint = datetime.now(timezone.utc).timestamp() logger.debug(f"Updating checkpointer to {new_checkpoint}") kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) diff --git a/package/bin/edges_phones_helper.py b/package/bin/edges_phones_helper.py index 54700fc..e0981b6 100644 --- a/package/bin/edges_phones_helper.py +++ b/package/bin/edges_phones_helper.py @@ -7,7 +7,7 @@ from solnlib.modular_input import checkpointer from splunklib import modularinput as smi -from datetime import datetime +from datetime import datetime, timezone from genesyscloud_client import GenesysCloudClient from genesyscloud_models import PhoneModel @@ -27,6 +27,41 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password def validate_input(definition: smi.ValidationDefinition): return @@ -66,8 +101,9 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_id = get_account_property(session_key, input_item.get("account"), "client_id") client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) client = GenesysCloudClient( - logger, client_id, client_secret, account_region + logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password ) checkpointer_key_name = input_name.split("/")[-1] @@ -107,7 +143,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): # Updating checkpoint if data was indexed to avoid losing info if event_counter > 0: logger.debug(f"Indexed '{event_counter}' events") - new_checkpoint = datetime.utcnow().timestamp() + new_checkpoint = datetime.now(timezone.utc).timestamp() logger.debug(f"Updating checkpointer to {new_checkpoint}") kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) diff --git a/package/bin/edges_trunks_metrics_helper.py b/package/bin/edges_trunks_metrics_helper.py index 454fef7..eb09266 100644 --- a/package/bin/edges_trunks_metrics_helper.py +++ b/package/bin/edges_trunks_metrics_helper.py @@ -7,7 +7,7 @@ from solnlib.modular_input import checkpointer from splunklib import modularinput as smi -from datetime import datetime +from datetime import datetime, timezone from genesyscloud_client import GenesysCloudClient from genesyscloud_models import TrunkModel @@ -27,6 +27,41 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password def validate_input(definition: smi.ValidationDefinition): return @@ -67,8 +102,9 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") # aws_region = input_item.get('region') + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) client = GenesysCloudClient( - logger, client_id, client_secret, account_region + logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password ) checkpointer_key_name = input_name.split("/")[-1] @@ -117,7 +153,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): # Updating checkpoint if data was indexed to avoid losing info if event_counter > 0: logger.debug(f"Indexed '{event_counter}' events") - new_checkpoint = datetime.utcnow().timestamp() + new_checkpoint = datetime.now(timezone.utc).timestamp() logger.debug(f"Updating checkpointer to {new_checkpoint}") kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index f1e4317..99c0dd7 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -2,19 +2,20 @@ import logging import json import os +import urllib3 from typing import List +from io import BytesIO from PureCloudPlatformClientV2.rest import ApiException from PureCloudPlatformClientV2.api_client import ApiClient +from PureCloudPlatformClientV2.configuration import Configuration class GenesysCloudClient: """ Interface with Genesys Cloud """ - client: ApiClient = None - - def __init__(self, logger: logging.Logger, client_id: str, client_secret: str, aws_region: str): + def __init__(self, logger: logging.Logger, client_id: str, client_secret: str, aws_region: str, proxy_url: str = None, proxy_username: str = None, proxy_password: str = None): self.logger = logger if PureCloudPlatformClientV2.PureCloudRegionHosts.__members__.get(aws_region): region = PureCloudPlatformClientV2.PureCloudRegionHosts[aws_region] @@ -24,6 +25,17 @@ def __init__(self, logger: logging.Logger, client_id: str, client_secret: str, a self.host = os.environ.get("GENESYSCLOUD_HOST", None) # If host is none, default value will be "https://api.mypurecloud.com" + # Singleton pattern. Configuration() is a globally shared object. + # Always set values to avoid data persistance from previous execution. + config = Configuration() + config.host = self.host + config.proxy = proxy_url + config.proxy_username = proxy_username + config.proxy_password = proxy_password + if proxy_url: + self.logger.info(f"Using proxy: {proxy_url}") + + # Note that passing self.host to the client can be removed as per singleton behavior. self.client = ApiClient(self.host).get_client_credentials_token( client_id, client_secret ) @@ -40,16 +52,38 @@ def _fetch(self, api_instance, f_name: str, *args, **kwargs): if not callable(function): raise AttributeError(f"{f_name} is not a callable function of the API instance") + # FIXME add a max_pages safeguard to avoid unbounded pagination while True: - api_response = function(*args, **kwargs) + try: + api_response = function(*args, **kwargs) + except ApiException as e: + if e.status in (301, 302, 303, 307, 308): + # When getting audit query results API could return a redirect with downloadUrl. + body = json.loads(e.body) + if "downloadUrl" in body.keys(): + self.logger.info(f"Got URL to download events from.") + buf = self.download(body["downloadUrl"]) + for item in json.loads(buf.read()): + items.append(item) + + if "cursor" in body.keys(): + cursor = body["cursor"] + kwargs["cursor"] = cursor + continue + else: + # Nothing else to be fetched + break + raise + else: + raise if isinstance(api_response, list): # A simple list (of strings) is returned as response items.extend(api_response) else: - # An object such as EdgeEntityListing|RoutingStatus|etc is returned + # An object such as EdgeEntityListing|RoutingStatus|AuditQueryExecutionResultsResponse|etc is returned enable_pagination = any(key in api_response.attribute_map for key in pagination_params) - if "entities" in api_response.attribute_map: + if hasattr(api_response, "entities") and api_response.entities: for item in api_response.entities: items.append(item) else: @@ -75,7 +109,7 @@ def get(self, api_instance_name: str, function_name: str, *args, **kwargs): :param api_instance_name: Name of the API instance e.g. TelephonyProvidersEdgeApi, RoutingApi, etc :param function_name: Name of the function to call in the API instance """ - self.logger.info(f"Getting data from {api_instance_name}") + self.logger.info(f"Getting data from {api_instance_name}->{function_name}") # Get the API class dynamically api_class = getattr(PureCloudPlatformClientV2, api_instance_name) @@ -87,10 +121,10 @@ def get(self, api_instance_name: str, function_name: str, *args, **kwargs): except AttributeError as e: self.logger.error(f"Error: {e}") except ApiException as e: - if e.status == 429 and e.reason.contains("Rate limit exceeded the maximum"): - self.logger.warning("Rate limit exceeded. Refreshing token.") - self.client.handle_expired_access_token() - if e.status == 401 and e.reason.contains("expir"): + if e.status == 429 and "Rate limit exceeded the maximum" in e.reason: + self.logger.warning("Rate limit exceeded. Refreshing token.") + self.client.handle_expired_access_token() + if e.status == 401 and "expir" in e.reason: # Haven't hit this yet. Message to be confirmed self.logger.warning("Token expired. Refreshing token.") self.client.handle_expired_access_token() @@ -105,6 +139,35 @@ def get(self, api_instance_name: str, function_name: str, *args, **kwargs): return [] + def download(self, url: str, chunk_size: int = 8192) -> BytesIO: + """ + Download a URL in chunks into an in-memory buffer. + :param url: URL to download events from. + :param chunk_size: Number of bytes to read per iteration. + :return: BytesIO buffer positioned at the start, containing the downloaded bytes. + """ + config = Configuration() + if config.proxy: + headers = None + if config.proxy_username and config.proxy_password: + headers = urllib3.make_headers( + proxy_basic_auth=f"{config.proxy_username}:{config.proxy_password}" + ) + http = urllib3.ProxyManager(config.proxy, proxy_headers=headers) + else: + http = urllib3.PoolManager() + + buffer = BytesIO() + with http.request("GET", url, preload_content=False) as response: + if response.status != 200: + raise urllib3.exceptions.HTTPError(f"Unexpected status: {response.status}") + for chunk in response.stream(chunk_size): + if chunk: + buffer.write(chunk) + + buffer.seek(0) # rewind so the buffer can be read from the start + return buffer + def convert_response(self, response: list, key: str) -> list: """ Convert data returned from paginating POST API. @@ -113,9 +176,10 @@ def convert_response(self, response: list, key: str) -> list: :return: List of items to be ingested. """ total_items = [] - for obj in response: - res_dict = obj.to_dict() or {} - total_items.extend(res_dict.get(key, []) or []) + if response is not None: + for obj in response: + res_dict = obj.to_dict() or {} + total_items.extend(res_dict.get(key, []) or []) return total_items def post(self, api_instance_name: str, function_name: str, model_name: str, body: dict, *args, **kwargs): @@ -129,7 +193,7 @@ def post(self, api_instance_name: str, function_name: str, model_name: str, body """ enable_pagination = False api_responses = [] - # Tipically 100 items per page is the max accepted + # Typically 100 items per page is the max accepted page_size = 100 page_number = 1 total_hits = 0 @@ -198,7 +262,8 @@ def post(self, api_instance_name: str, function_name: str, model_name: str, body except Exception as e: self.logger.error(f"Pagination enabled but neither 'total_hits' nor 'total' is returned: {api_response.attribute_map}") return None - total_hits = api_response.total_hits + else: + total_hits = api_response.total_hits if total_hits - (page_size * page_number) > 0: page_number += 1 @@ -212,14 +277,20 @@ def post(self, api_instance_name: str, function_name: str, model_name: str, body return api_responses except ApiException as e: - if e.status == 429 and e.reason.contains("Rate limit exceeded the maximum"): - self.logger.warning("Rate limit exceeded. Refreshing token.") - self.client.handle_expired_access_token() - if e.status == 401 and e.reason.contains("expir"): + if e.status == 429 and "Rate limit exceeded the maximum" in e.reason: + self.logger.warning("Rate limit exceeded. Refreshing token.") + self.client.handle_expired_access_token() + if e.status == 401 and "expir" in e.reason: # Haven't hit this yet. Message to be confirmed self.logger.warning("Token expired. Refreshing token.") self.client.handle_expired_access_token() - body = json.loads(e.body) - message = body["message"] - self.logger.error(f"Exception when calling {api_instance_name}->{function_name}: [{e.status}] {e.reason} - {message}") + err_message = f"Exception when calling {api_instance_name}->{function_name}:" + try: + body = json.loads(e.body) + message = body["message"] + self.logger.error(f"{err_message} [{e.status}] {e.reason} - {message}") + except ValueError as ve: + self.logger.warning(f"{err_message} {ve}") + self.logger.error(f"{err_message} [{e.status}] {e.reason} - {e.body}") + return None diff --git a/package/bin/genesyscloud_models.py b/package/bin/genesyscloud_models.py index 556930a..9ae96ee 100644 --- a/package/bin/genesyscloud_models.py +++ b/package/bin/genesyscloud_models.py @@ -12,8 +12,6 @@ ) class GCBaseModel: - data: List[dict] = [] - def __init__(self, data: List[dict]) -> None: self.data = data @@ -21,12 +19,12 @@ def to_camelcase(self, s: str) -> str: return re.sub(r'(?!^)_([a-zA-Z])', lambda m: m.group(1).upper(), s) def to_string(self, dt: datetime) -> str: - format = "%d-%m-%YT%H:%M:%S.%f%z" - return dt.strftime(format) + formatting_str = "%d-%m-%YT%H:%M:%S.%f%z" + return dt.strftime(formatting_str) def to_datetime(self, dt_string: str) -> datetime: - format = "%Y-%m-%dT%H:%M:%S.%fZ" - return datetime.datetime.strptime(dt_string, format) + formatting_str = "%Y-%m-%dT%H:%M:%S.%fZ" + return datetime.datetime.strptime(dt_string, formatting_str) def extract(self, idx: int, sub_key: str, keys_to_extract: list, enable_camelcase: bool = False) -> dict: """ @@ -62,10 +60,10 @@ def trunk_ids(self) -> List[str]: def get_trunk_ids(self, batch: int = 0) -> Tuple[List[str], bool]: factor = self.MAX_TRUNK_IDS * batch - slice = self.MAX_TRUNK_IDS + factor + slice_limit = self.MAX_TRUNK_IDS + factor remaining_trunks = abs(len(self.data) - factor) has_next_batch = remaining_trunks > self.MAX_TRUNK_IDS - return [trunk["id"] for trunk in self.data[factor:slice]], has_next_batch + return [trunk["id"] for trunk in self.data[factor:slice_limit]], has_next_batch def get_trunk(self, tid: str) -> dict: ret_trunk = {} @@ -75,9 +73,10 @@ def get_trunk(self, tid: str) -> dict: "connected_status", "ip_status" ] - trunk = [t for t in self.data if t["id"] == tid][0] - for key, value in trunk.items(): - ret_trunk.update({k: trunk[k] for k in required_keys}) + trunk = next((t for t in self.data if t["id"] == tid), None) + if trunk is None: + raise ValueError(f"Trunk {tid} not found") + ret_trunk.update({k: trunk[k] for k in required_keys if k in trunk}) return ret_trunk @@ -92,10 +91,10 @@ def __init__(self, edges: List[Edge]): def get_edge_ids(self, batch: int = 0) -> Tuple[List[str], bool]: factor = self.MAX_EDGE_IDS*batch - slice = self.MAX_EDGE_IDS + factor + slice_limit = self.MAX_EDGE_IDS + factor remaining_edges = abs(len(self.data) - factor) has_next_batch = remaining_edges > self.MAX_EDGE_IDS - return [edge["id"] for edge in self.data[factor:slice]], has_next_batch + return [edge["id"] for edge in self.data[factor:slice_limit]], has_next_batch def get_edge(self, eid: str) -> dict: ret_edge = { "site": {} } @@ -106,9 +105,10 @@ def get_edge(self, eid: str) -> dict: "conversation_count", "os_name" ] - edge = [e for e in self.data if e["id"] == eid][0] - for key, value in edge.items(): - ret_edge.update({k: edge[k] for k in required_keys}) + edge = next((e for e in self.data if e["id"] == eid), None) + if edge is None: + raise ValueError(f"Edge {eid} not found") + ret_edge.update({k: edge[k] for k in required_keys if k in edge}) # Avoid indexing a lot of "null" values added # by the to_dict() SDK function for "site" data for key, value in edge["site"].items(): @@ -161,18 +161,19 @@ def queue_ids(self) -> List[str]: def get_queue_ids(self, batch: int = 0) -> Tuple[List[str], bool]: factor = self.MAX_QUEUE_IDS*batch - slice = self.MAX_QUEUE_IDS + factor + slice_limit = self.MAX_QUEUE_IDS + factor remaining_queues = abs(len(self.data) - factor) has_next_batch = remaining_queues > self.MAX_QUEUE_IDS - return [queue["id"] for queue in self.data[factor:slice]], has_next_batch + return [queue["id"] for queue in self.data[factor:slice_limit]], has_next_batch def get_queue(self, qid: str) -> dict: ret_queue = {} required_keys = ["id", "name"] - queue = [q for q in self.data if q["id"] == qid][0] - for key, value in queue.items(): - ret_queue.update({k: queue[k] for k in required_keys}) + queue = next((q for q in self.data if q["id"] == qid), None) + if queue is None: + raise ValueError(f"Queue {qid} not found") + ret_queue.update({k: queue[k] for k in required_keys if k in queue}) return ret_queue @@ -191,16 +192,17 @@ def user_ids(self) -> List[str]: def get_user_ids(self, batch: int = 0) -> Tuple[List[str], bool]: factor = self.MAX_USER_IDS*batch - slice = self.MAX_USER_IDS + factor + slice_limit = self.MAX_USER_IDS + factor remaining_users = abs(len(self.data) - factor) has_next_batch = remaining_users > self.MAX_USER_IDS - return [user["id"] for user in self.data[factor:slice]], has_next_batch + return [user["id"] for user in self.data[factor:slice_limit]], has_next_batch def get_user(self, uid: str) -> dict: ret_user = {} required_keys = ["id", "name", "chat", "email", "division"] - user = [u for u in self.data if u["id"] == uid][0] - for key, value in user.items(): - ret_user.update({k: user[k] for k in required_keys}) + user = next((u for u in self.data if u["id"] == uid), None) + if user is None: + raise ValueError(f"User {uid} not found") + ret_user.update({k: user[k] for k in required_keys if k in user}) return ret_user \ No newline at end of file diff --git a/package/bin/queue_observations_helper.py b/package/bin/queue_observations_helper.py index 1982132..59d953b 100644 --- a/package/bin/queue_observations_helper.py +++ b/package/bin/queue_observations_helper.py @@ -22,6 +22,42 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password + def validate_input(definition: smi.ValidationDefinition): return @@ -43,9 +79,10 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_id = get_account_property(session_key, input_item.get("account"), "client_id") client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") account_region = get_account_property(session_key, input_item.get("account"), "region") + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) # Initialize Genesys Cloud client client = GenesysCloudClient( - logger, client_id, client_secret, account_region + logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password ) # Getting data from API @@ -102,7 +139,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): event_writer.write_event( smi.Event( # Index time not needed? - data=json.dumps(data_entry), + data=json.dumps(data_entry, ensure_ascii=False, default=str), index=input_item.get("index"), sourcetype=sourcetype ) diff --git a/package/bin/status_page_metrics_helper.py b/package/bin/status_page_metrics_helper.py index 2bc2953..fe83119 100644 --- a/package/bin/status_page_metrics_helper.py +++ b/package/bin/status_page_metrics_helper.py @@ -1,12 +1,14 @@ import json import logging import requests -from datetime import datetime, timezone +import import_declare_test from solnlib import conf_manager, log from splunklib import modularinput as smi from solnlib.modular_input import checkpointer +from datetime import datetime, timezone + ADDON_NAME = "genesys_cloud_ta" STATUS_PAGE_API_URL = "https://status.mypurecloud.com/api" @@ -26,11 +28,55 @@ def get_account_property(session_key: str, account_name: str, property_name: str account_conf_file = cfm.get_conf("genesys_cloud_ta_account") return account_conf_file.get(account_name).get(property_name) -def fetch_status_page_data(logger: logging.Logger): +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + proxy_url = f"{proxy_type}://{url}:{port}" + else: + proxy_url = f"{proxy_type}://{user}:{password}@{url}:{port}" + + return proxy_url + +def fetch_status_page_data(logger: logging.Logger, proxy_url: str = None): """Fetch status data from the Genesys Cloud Status Page API""" try: # Get status summary instead of incidents - summary_response = requests.get(f"{STATUS_PAGE_API_URL}/v2/summary.json") + proxies = None + if proxy_url: + proxies = { + "http": proxy_url, + "https": proxy_url, + } + logger.info(f"Using proxy: {proxy_url}") + summary_response = requests.get(f"{STATUS_PAGE_API_URL}/v2/summary.json", proxies=proxies) summary_response.raise_for_status() summary_data = summary_response.json() @@ -72,7 +118,8 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): logger.debug(f"Current status page checkpoint: {status_page_checkpoint}") # Fetch data from Status Page API - summary = fetch_status_page_data(logger) + proxy = get_account_proxy(logger=logger, session_key=session_key) + summary = fetch_status_page_data(logger, proxy_url=proxy) # Process summary data sourcetype = "genesyscloud:operational:system" @@ -118,8 +165,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): input_name, sourcetype, events_count, - input_item.get("index"), - account=input_item.get("account") + input_item.get("index") ) log.modular_input_end(logger, normalized_input_name) diff --git a/package/bin/user_aggregates_helper.py b/package/bin/user_aggregates_helper.py index faeae73..f114d66 100644 --- a/package/bin/user_aggregates_helper.py +++ b/package/bin/user_aggregates_helper.py @@ -5,7 +5,7 @@ from solnlib import conf_manager, log from solnlib.modular_input import checkpointer from splunklib import modularinput as smi -from datetime import datetime, timedelta +from datetime import datetime, timezone from dateutil.relativedelta import relativedelta from genesyscloud_client import GenesysCloudClient @@ -30,6 +30,42 @@ def get_account_property(session_key: str, account_name: str, property_name: str def validate_input(definition: smi.ValidationDefinition): return +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password + def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): for input_name, input_item in inputs.inputs.items(): normalized_input_name = input_name.split("/")[-1] @@ -54,14 +90,15 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") account_region = get_account_property(session_key, input_item.get("account"), "region") + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) # Initialize Genesys Cloud client client = GenesysCloudClient( - logger, client_id, client_secret, account_region + logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password ) # Initialize checkpointing checkpointer_key_name = input_name.split("/")[-1] - now = datetime.now() + now = datetime.now(timezone.utc) # No checkpoint? Default it to four years ago per API docs last_checkpoint = ( kvstore_checkpointer.get(checkpointer_key_name) diff --git a/package/bin/user_routing_status_helper.py b/package/bin/user_routing_status_helper.py index 3c16686..34a0156 100644 --- a/package/bin/user_routing_status_helper.py +++ b/package/bin/user_routing_status_helper.py @@ -6,7 +6,7 @@ from solnlib.modular_input import checkpointer from splunklib import modularinput as smi -from datetime import datetime +from datetime import datetime, timezone from genesyscloud_client import GenesysCloudClient from genesyscloud_models import UserModel @@ -29,6 +29,42 @@ def get_account_property(session_key: str, account_name: str, property_name: str def validate_input(definition: smi.ValidationDefinition): return +def get_account_proxy(logger, session_key: str): + try: + proxy_config = conf_manager.get_proxy_dict( + logger=logger, + session_key=session_key, + app_name=ADDON_NAME, + conf_name="genesys_cloud_ta_settings", + ) + # Handle invalid port case + except InvalidPortError as e: + logger.error(f"Proxy configuration error: {e}") + + # Handle invalid hostname case + except InvalidHostnameError as e: + logger.error(f"Proxy configuration error: {e}") + + if not proxy_config or not proxy_config.get('proxy_enabled'): + logger.info('Proxy is not enabled') + return None, None, None + + url = proxy_config.get('proxy_url') + port = proxy_config.get('proxy_port') + user = proxy_config.get('proxy_username') + password = proxy_config.get('proxy_password') + + if not all((user, password)): + logger.info('Proxy has no credentials found') + user, password = None, None + + proxy_type = proxy_config.get('proxy_type') + proxy_type = proxy_type.lower() if proxy_type else 'http' + + proxy_url = f"{proxy_type}://{url}:{port}" + + return proxy_url, user, password + def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): for input_name, input_item in inputs.inputs.items(): normalized_input_name = input_name.split("/")[-1] @@ -53,9 +89,10 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client_secret = get_account_property(session_key, input_item.get("account"), "client_secret") account_region = get_account_property(session_key, input_item.get("account"), "region") + proxy_url, proxy_username, proxy_password = get_account_proxy(logger=logger, session_key=session_key) # Initialize Genesys Cloud client client = GenesysCloudClient( - logger, client_id, client_secret, account_region + logger, client_id, client_secret, account_region, proxy_url=proxy_url, proxy_username=proxy_username, proxy_password=proxy_password ) # Initialize checkpointing @@ -72,13 +109,14 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client.get("UsersApi", "get_users") ) - #Getting user routing status + # Getting user routing status sourcetype = "genesyscloud:users:users:routingstatus" rcounter = 0 + # CAREFUL! 300+ users = 300+ API calls. Most likely hit rate limits. for uid in user_model.user_ids: response = client.get("UsersApi", "get_user_routingstatus", uid) - if (response[0].start_time): + if response and response[0].start_time: event_time_epoch = response[0].start_time.timestamp() if event_time_epoch > current_checkpoint: @@ -98,7 +136,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): # Updating checkpoint if data was indexed to avoid losing info if rcounter > 0: logger.debug(f"Indexed '{rcounter}' events") - new_checkpoint = datetime.utcnow().timestamp() + new_checkpoint = datetime.now(timezone.utc).timestamp() logger.debug(f"Updating checkpointer to {new_checkpoint}") kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) diff --git a/tests/integration/GenesysCloudTATest.py b/tests/integration/GenesysCloudTATest.py index 749a87c..deba5e5 100644 --- a/tests/integration/GenesysCloudTATest.py +++ b/tests/integration/GenesysCloudTATest.py @@ -3,10 +3,7 @@ import pytest import logging -try: - from StringIO import StringIO -except ImportError: - from io import StringIO +from io import StringIO sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "../../package/bin"))) @@ -77,5 +74,5 @@ def get_genesyscloud_accounts_configuration(cls): "client_secret": cls.CLIENT_SECRET, "aws_region": cls.AWS_REGION } - cls.logger.info(f"{configs}") + # cls.logger.info(f"{configs}") return configs diff --git a/tests/integration/test_genesyscloud_client.py b/tests/integration/test_genesyscloud_client.py index ebd7d18..8defb63 100644 --- a/tests/integration/test_genesyscloud_client.py +++ b/tests/integration/test_genesyscloud_client.py @@ -177,4 +177,4 @@ def test_POST_wrong_body(self, body_routing_queues, body_users, api_name, func_n model_name, body ) - assert response == None + assert response is None diff --git a/tests/modinput_functional/BaseTATest.py b/tests/modinput_functional/BaseTATest.py index 9b529ef..fd5919b 100644 --- a/tests/modinput_functional/BaseTATest.py +++ b/tests/modinput_functional/BaseTATest.py @@ -1,10 +1,7 @@ import pytest import logging -try: - from StringIO import StringIO -except ImportError: - from io import StringIO +from io import StringIO import splunklib.client as client @@ -41,6 +38,7 @@ def setup_class(cls): cls.splunk_client = client.connect( host=cls.splunk_url, port=cls.splunkd_port, + scheme="https", username=cls.username, password=cls.password, verify=False, @@ -57,9 +55,12 @@ def setup_class(cls): def teardown_class(cls): cls.delete_genesyscloud_accounts(f"{cls.TA_APP_NAME}_account") index = cls.splunk_client.indexes[cls.INDEX] - # Clean index before deleting (removes all events) - index.clean(timeout=300) - index.delete() + try: + # Clean index before deleting (removes all events). + # index.clean(timeout=300) + index.delete() + except Exception as e: + cls.logger.warning(f"Teardown failed: {e}") cls.splunk_client = None diff --git a/tests/modinput_functional/test_genesyscloud_ta.py b/tests/modinput_functional/test_genesyscloud_ta.py index a238750..86097ac 100644 --- a/tests/modinput_functional/test_genesyscloud_ta.py +++ b/tests/modinput_functional/test_genesyscloud_ta.py @@ -28,9 +28,10 @@ def _search(self, search_query: str, run_counter: int=0, timeout: int=40, sleep_ lst_results = [] elapsed_time = 0 kwargs = { - "earliest_time": "0", + "earliest_time": 0, "latest_time": "now", - "output_mode": "json" + "output_mode": "json", + "count": 0 } # +1min timeout each retry tot_timeout = timeout + (run_counter * 60) @@ -39,6 +40,9 @@ def _search(self, search_query: str, run_counter: int=0, timeout: int=40, sleep_ oneshot = self.splunk_client.jobs.export(search_query, **kwargs) reader = results.JSONResultsReader(oneshot) for result in reader: + # if isinstance(result, results.Message): + # # ⚠️ Don't ignore these — they may explain why results are empty + # self.logger.warning(f"[{result.type}] {result.message}") if not isinstance(result, dict): # Diagnostic messages may be returned in the results continue @@ -109,7 +113,7 @@ def test_input_user_routing_status(self): sourcetype = "genesyscloud:users:users:routingstatus" spl = f"search index={self.INDEX} sourcetype={sourcetype}" for attempt in range(self.RETRY): - results = self._search(search_query=spl, run_counter=attempt, timeout=100) + results = self._search(search_query=spl, run_counter=attempt) #, timeout=100) if len(results) > 0 or attempt == (self.RETRY - 1): self.logger.debug(f"Results {len(results)} and attempt: {attempt}") break @@ -124,7 +128,7 @@ def test_input_edges_metrics(self): sourcetype = "genesyscloud:telephonyprovidersedge:edges:metrics" spl = f"search index={self.INDEX} sourcetype={sourcetype}" for attempt in range(self.RETRY): - results = self._search(search_query=spl, run_counter=attempt, timeout=100) + results = self._search(search_query=spl, run_counter=attempt) #, timeout=100) if len(results) > 0 or attempt == (self.RETRY - 1): self.logger.debug(f"Results {len(results)} and attempt: {attempt}") break @@ -134,6 +138,7 @@ def test_input_edges_metrics(self): assert len(results) > 0 and len(results) <= 157 assert results[0]["source"] == "edges_metrics://edges_metrics" + # @pytest.mark.skip(reason="already tested") def test_input_edges_phones(self): """ This test will check whether data was successfully indexed @@ -159,7 +164,7 @@ def test_input_edges_trunks_metrics(self): sourcetype = "genesyscloud:telephonyprovidersedge:trunks:metrics" spl = f"search index={self.INDEX} sourcetype={sourcetype}" for attempt in range(self.RETRY): - results = self._search(search_query=spl, run_counter=attempt, timeout=100) + results = self._search(search_query=spl, run_counter=attempt) #, timeout=100) if len(results) > 0 or attempt == (self.RETRY - 1): self.logger.debug(f"Results {len(results)} and attempt: {attempt}") break @@ -174,7 +179,7 @@ def test_input_queue_observations(self): sourcetype = "genesyscloud:analytics:queues:observations" spl = f"search index={self.INDEX} sourcetype={sourcetype}" for attempt in range(self.RETRY): - results = self._search(search_query=spl, run_counter=attempt, timeout=100) + results = self._search(search_query=spl, run_counter=attempt) #, timeout=100) if len(results) > 0 or attempt == (self.RETRY - 1): self.logger.debug(f"Results {len(results)} and attempt: {attempt}") break @@ -206,7 +211,7 @@ def test_input_actions_metrics(self): source = "actions_metrics://actions_metrics" spl = f"search index={self.INDEX} sourcetype={sourcetype} source={source}" for attempt in range(self.RETRY): - results = self._search(search_query=spl, run_counter=attempt, timeout=90) + results = self._search(search_query=spl, run_counter=attempt) #, timeout=90) if results or attempt == (self.RETRY - 1): self.logger.debug(f"Results {len(results)} and attempt: {attempt}") break @@ -214,7 +219,6 @@ def test_input_actions_metrics(self): assert len(results) > 0 and len(results) <= 5 assert results[0]["source"] == source - def test_input_audit_query(self): """ This test will check whether data was successfully indexed @@ -222,7 +226,7 @@ def test_input_audit_query(self): sourcetype = "genesyscloud:operational:audits" spl = f"search index={self.INDEX} sourcetype={sourcetype}" for attempt in range(self.RETRY): - results = self._search(search_query=spl, run_counter=attempt, timeout=120) + results = self._search(search_query=spl, run_counter=attempt) #, timeout=120) if len(results) > 0 or attempt == (self.RETRY - 1): self.logger.debug(f"Results {len(results)} and attempt: {attempt}") break