From bc3a4d9364cce8603ebaac924be439da24b4102d Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 9 Mar 2026 13:57:18 +0100 Subject: [PATCH 01/34] fix:pagination logic overwrite and unsafe error body parsing --- package/bin/genesyscloud_client.py | 27 ++++++++++++++++----------- 1 file changed, 16 insertions(+), 11 deletions(-) diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index f1e4317..8a015f7 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -12,8 +12,6 @@ class GenesysCloudClient: """ Interface with Genesys Cloud """ - client: ApiClient = None - def __init__(self, logger: logging.Logger, client_id: str, client_secret: str, aws_region: str): self.logger = logger if PureCloudPlatformClientV2.PureCloudRegionHosts.__members__.get(aws_region): @@ -40,6 +38,7 @@ 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) @@ -87,10 +86,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"): + 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 e.reason.contains("expir"): + 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() @@ -129,7 +128,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 +197,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 +212,19 @@ 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"): + 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 e.reason.contains("expir"): + 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}") + 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 From 8ae821c649c5b4283a0546a9bfd3855837ed786b Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 9 Mar 2026 14:09:28 +0100 Subject: [PATCH 02/34] fix:reduntant loop, indexing without safety --- package/bin/genesyscloud_models.py | 54 ++++++++++++++++-------------- 1 file changed, 28 insertions(+), 26 deletions(-) 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 From c841255b0b35e44be9eaa8eefe1f9fe7782c02a8 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 9 Mar 2026 14:13:57 +0100 Subject: [PATCH 03/34] fix:mixed checkpoint value formats --- package/bin/conversations_details_helper.py | 2 +- package/bin/edges_metrics_helper.py | 2 +- package/bin/edges_phones_helper.py | 2 +- package/bin/edges_trunks_metrics_helper.py | 2 +- package/bin/user_aggregates_helper.py | 2 +- package/bin/user_routing_status_helper.py | 2 +- 6 files changed, 6 insertions(+), 6 deletions(-) diff --git a/package/bin/conversations_details_helper.py b/package/bin/conversations_details_helper.py index c578c30..ddea3af 100644 --- a/package/bin/conversations_details_helper.py +++ b/package/bin/conversations_details_helper.py @@ -70,7 +70,7 @@ 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() + now = datetime.now(timezone.utc) fallback_start = (now - relativedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") start_date = input_item.get("start_date") if start_date is not None: diff --git a/package/bin/edges_metrics_helper.py b/package/bin/edges_metrics_helper.py index 10dc99c..699754e 100644 --- a/package/bin/edges_metrics_helper.py +++ b/package/bin/edges_metrics_helper.py @@ -117,7 +117,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) 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..073460d 100644 --- a/package/bin/edges_phones_helper.py +++ b/package/bin/edges_phones_helper.py @@ -107,7 +107,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) 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..242546e 100644 --- a/package/bin/edges_trunks_metrics_helper.py +++ b/package/bin/edges_trunks_metrics_helper.py @@ -117,7 +117,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) logger.debug(f"Updating checkpointer to {new_checkpoint}") kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) diff --git a/package/bin/user_aggregates_helper.py b/package/bin/user_aggregates_helper.py index faeae73..69f95f3 100644 --- a/package/bin/user_aggregates_helper.py +++ b/package/bin/user_aggregates_helper.py @@ -61,7 +61,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): # 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..555315e 100644 --- a/package/bin/user_routing_status_helper.py +++ b/package/bin/user_routing_status_helper.py @@ -98,7 +98,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) logger.debug(f"Updating checkpointer to {new_checkpoint}") kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) From 3c7682aa3e52db9855f11e56c7b3c513134a2076 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 11:27:09 +0100 Subject: [PATCH 04/34] style:removed duplicate option,renamed variable --- package/bin/conversations_metrics_helper.py | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) diff --git a/package/bin/conversations_metrics_helper.py b/package/bin/conversations_metrics_helper.py index 4ac3c2d..d9fd44a 100644 --- a/package/bin/conversations_metrics_helper.py +++ b/package/bin/conversations_metrics_helper.py @@ -73,14 +73,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 +122,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 From 929549ffd221ede9e4a151e0041bec09697e5492 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 11:29:07 +0100 Subject: [PATCH 05/34] chore:added missing import --- package/bin/status_page_metrics_helper.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/package/bin/status_page_metrics_helper.py b/package/bin/status_page_metrics_helper.py index 2bc2953..e468ffc 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" @@ -118,8 +120,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) From 98ca797966737bc4d841876311fcc52d6f86c110 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 11:46:33 +0100 Subject: [PATCH 06/34] syle:added default modular input log messages --- package/bin/audit_query_helper.py | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/package/bin/audit_query_helper.py b/package/bin/audit_query_helper.py index ba5b229..a042b4c 100644 --- a/package/bin/audit_query_helper.py +++ b/package/bin/audit_query_helper.py @@ -94,7 +94,7 @@ 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 @@ -116,12 +116,14 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): **params ) event_counter = 0 + sourcetype="genesyscloud:operational:audits" + for entity in results: event_writer.write_event( smi.Event( data=json.dumps(entity.to_dict(), ensure_ascii=False, default=str), index=input_item.get("index"), - sourcetype="genesyscloud:operational:audits", + sourcetype=sourcetype, time=entity.event_date.timestamp() ) ) @@ -133,6 +135,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, From 3badc0985f418a8fc8f9019ac30c3f0e0bd9d56b Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 11:47:50 +0100 Subject: [PATCH 07/34] fix:handled None returned values --- package/bin/actions_metrics_helper.py | 80 +++++++++++++-------------- package/bin/genesyscloud_client.py | 7 ++- 2 files changed, 44 insertions(+), 43 deletions(-) diff --git a/package/bin/actions_metrics_helper.py b/package/bin/actions_metrics_helper.py index 5fb4802..591b4bc 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" @@ -30,8 +30,6 @@ def validate_input(definition: smi.ValidationDefinition): 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) @@ -87,47 +85,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/genesyscloud_client.py b/package/bin/genesyscloud_client.py index 8a015f7..a327003 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -112,9 +112,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): From 2226ea8c89a796fa30bdf3f0a2d13a61bace04cd Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 11:53:29 +0100 Subject: [PATCH 08/34] fix:prevent IndexError in case of empty responses --- package/bin/user_routing_status_helper.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/package/bin/user_routing_status_helper.py b/package/bin/user_routing_status_helper.py index 555315e..129df96 100644 --- a/package/bin/user_routing_status_helper.py +++ b/package/bin/user_routing_status_helper.py @@ -72,13 +72,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: From b546be2eb44ffe33c7f3232d4b98220b342e24e2 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 11:59:41 +0100 Subject: [PATCH 09/34] refactor:aligned data dumping args among inputs --- package/bin/queue_observations_helper.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/package/bin/queue_observations_helper.py b/package/bin/queue_observations_helper.py index 1982132..a5b3468 100644 --- a/package/bin/queue_observations_helper.py +++ b/package/bin/queue_observations_helper.py @@ -102,7 +102,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 ) From a228482e7c293c7aab266dca10d4a63d721b1889 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 12:05:21 +0100 Subject: [PATCH 10/34] tests(deps):removed py2 support on import --- tests/integration/GenesysCloudTATest.py | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) 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 From 099f4336d4bf9b51e47393a5fe5ecda8a5ec9a89 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 12:05:55 +0100 Subject: [PATCH 11/34] tests(deps):removed py2 support on import --- tests/modinput_functional/BaseTATest.py | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/tests/modinput_functional/BaseTATest.py b/tests/modinput_functional/BaseTATest.py index 9b529ef..9ad4db7 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 From 235593290698127dbb5d8a7558a80ed9ddb361b1 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 12:06:58 +0100 Subject: [PATCH 12/34] tests:updated None assertion as per PEP8 --- tests/integration/test_genesyscloud_client.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 From ae20d461c77e9a9b4becf1f0db85a1a31d7a91ac Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 12:42:32 +0100 Subject: [PATCH 13/34] fix:added missing variable assignment --- package/bin/genesyscloud_client.py | 1 + 1 file changed, 1 insertion(+) diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index a327003..e86dd50 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -220,6 +220,7 @@ def post(self, api_instance_name: str, function_name: str, model_name: str, body # Haven't hit this yet. Message to be confirmed self.logger.warning("Token expired. Refreshing token.") self.client.handle_expired_access_token() + err_message = f"Exception when calling {api_instance_name}->{function_name}:" try: body = json.loads(e.body) message = body["message"] From a73d40f204039aaec0a3fea29569e9c21369a439 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 13:04:11 +0100 Subject: [PATCH 14/34] refactor(audit_query):removed undefined start_date input item --- package/bin/audit_query_helper.py | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/package/bin/audit_query_helper.py b/package/bin/audit_query_helper.py index a042b4c..a85ee83 100644 --- a/package/bin/audit_query_helper.py +++ b/package/bin/audit_query_helper.py @@ -66,16 +66,14 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): checkpointer_key_name = normalized_input_name now = datetime.now(timezone.utc) - fallback_start = (now - relativedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") + # Setting a fallback start date of 7 days ago from now + start_date = (now - relativedelta(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) or start_date end_time = now.strftime("%Y-%m-%dT%H:%M:%SZ") interval = f"{start_time}/{end_time}" - body = {"interval": interval} + + body = { "interval": interval } logger.debug(f"Request body: {body}") response = client.post( From 5890ecfc99349963d8e9d63d8c0f173a210cdb08 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 13:18:58 +0100 Subject: [PATCH 15/34] ci:updated gh actions to latest releases --- .github/workflows/build-appinspect.yml | 16 +++++++++------- .github/workflows/docs.yml | 5 +++-- .github/workflows/release.yml | 11 ++++++----- 3 files changed, 18 insertions(+), 14 deletions(-) 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 From 91afb97b2684f5ee7c377b95b19766aae5480df9 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 13:20:27 +0100 Subject: [PATCH 16/34] ci(tests):removed tests in 9.2.5,improved sleep with readiness poll --- .github/workflows/tests.yml | 27 ++++++++++++++++++--------- 1 file changed, 18 insertions(+), 9 deletions(-) diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 329c383..f100b0f 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,7 +140,7 @@ 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 @@ -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,7 +272,7 @@ 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 From 0fd706eca5f3d0e96dd97b36cb6c3905215e8d55 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 13:26:29 +0100 Subject: [PATCH 17/34] ops(makefile):prevented execution in a separate shell --- Makefile | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/Makefile b/Makefile index f9d8bc9..c688555 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: @@ -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; + cd tests && \ python -m pytest tests/modinput_functional/* From 4f1ab50aa151cd35917597a399eb534b942af61d Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 16 Mar 2026 14:07:03 +0100 Subject: [PATCH 18/34] ci(tests):updated gh action test-reporter to latest release --- .github/workflows/tests.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index f100b0f..de039ca 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -145,7 +145,7 @@ jobs: name: mockoon-logs-${{ matrix.version }} path: logs/mockoon-${{ matrix.version }}.log - - uses: dorny/test-reporter@v2 + - uses: dorny/test-reporter@v2.6.0 if: always() with: name: Tests Results @@ -277,7 +277,7 @@ jobs: name: mockoon-logs-${{ matrix.version }} path: logs/mockoon-${{ matrix.version }}.log - - uses: dorny/test-reporter@v2 + - uses: dorny/test-reporter@v2.6.0 if: always() with: name: Tests Results From 3cdbe384aabfe984f90a85f6b91723828e834f9a Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 17 Mar 2026 12:12:00 +0100 Subject: [PATCH 19/34] ops(makefile):prevented execution in a separate shell --- Makefile | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/Makefile b/Makefile index c688555..222a0ce 100644 --- a/Makefile +++ b/Makefile @@ -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 @@ -37,4 +37,4 @@ run-tests: install-tests run-functional-tests: install-tests cd tests && \ - python -m pytest tests/modinput_functional/* + python -m pytest modinput_functional/* From d9b1914d5704d365f935b8bc0b6e759e30ca621f Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 17 Mar 2026 12:13:59 +0100 Subject: [PATCH 20/34] fix:imported missing libraries, fixed checkpoint formatting --- package/bin/conversations_details_helper.py | 2 +- package/bin/edges_metrics_helper.py | 4 ++-- package/bin/edges_phones_helper.py | 4 ++-- package/bin/edges_trunks_metrics_helper.py | 4 ++-- package/bin/user_aggregates_helper.py | 2 +- package/bin/user_routing_status_helper.py | 4 ++-- 6 files changed, 10 insertions(+), 10 deletions(-) diff --git a/package/bin/conversations_details_helper.py b/package/bin/conversations_details_helper.py index ddea3af..80ebdb3 100644 --- a/package/bin/conversations_details_helper.py +++ b/package/bin/conversations_details_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 dateutil.relativedelta import relativedelta from genesyscloud_client import GenesysCloudClient diff --git a/package/bin/edges_metrics_helper.py b/package/bin/edges_metrics_helper.py index 699754e..908dd86 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 @@ -117,7 +117,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.now(timezone.utc) + 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 073460d..c258623 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 @@ -107,7 +107,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.now(timezone.utc) + 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 242546e..ce1006b 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 @@ -117,7 +117,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.now(timezone.utc) + 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/user_aggregates_helper.py b/package/bin/user_aggregates_helper.py index 69f95f3..3281407 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 diff --git a/package/bin/user_routing_status_helper.py b/package/bin/user_routing_status_helper.py index 129df96..194bd53 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 @@ -99,7 +99,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.now(timezone.utc) + new_checkpoint = datetime.now(timezone.utc).timestamp() logger.debug(f"Updating checkpointer to {new_checkpoint}") kvstore_checkpointer.update(checkpointer_key_name, new_checkpoint) From 0ff6652006399985739c6b1e1e07e4894b43769a Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 17 Mar 2026 14:11:00 +0100 Subject: [PATCH 21/34] tests:added scheme to splunk connection, skip index clean --- tests/modinput_functional/BaseTATest.py | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/tests/modinput_functional/BaseTATest.py b/tests/modinput_functional/BaseTATest.py index 9ad4db7..fd5919b 100644 --- a/tests/modinput_functional/BaseTATest.py +++ b/tests/modinput_functional/BaseTATest.py @@ -38,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, @@ -54,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 From aa35f8202e0beb8496af5a5039c5f4b3cbcc5b76 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 17 Mar 2026 14:11:24 +0100 Subject: [PATCH 22/34] tests:edited search params, removed timeout --- .../test_genesyscloud_ta.py | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 deletions(-) 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 From f3c7f4d9e0a85fb0781b2abf60b66b10420cc3d5 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 17 Mar 2026 15:05:02 +0100 Subject: [PATCH 23/34] refactor:removed usage of dateutil for calculating dates --- package/bin/audit_query_helper.py | 3 +-- package/bin/conversations_details_helper.py | 7 +++---- 2 files changed, 4 insertions(+), 6 deletions(-) diff --git a/package/bin/audit_query_helper.py b/package/bin/audit_query_helper.py index a85ee83..73621f5 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" @@ -67,7 +66,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): now = datetime.now(timezone.utc) # Setting a fallback start date of 7 days ago from now - start_date = (now - relativedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") + start_date = (now - timedelta(days=7)).strftime("%Y-%m-%dT%H:%M:%SZ") start_time = kvstore_checkpointer.get(checkpointer_key_name) or start_date end_time = now.strftime("%Y-%m-%dT%H:%M:%SZ") diff --git a/package/bin/conversations_details_helper.py b/package/bin/conversations_details_helper.py index 80ebdb3..3d0c66d 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, timezone -from dateutil.relativedelta import relativedelta +from datetime import datetime, timezone, timedelta from genesyscloud_client import GenesysCloudClient @@ -42,7 +41,7 @@ def validate_input(definition: smi.ValidationDefinition): start_date = definition.parameters.get("start_date") if start_date is not None: input_date = datetime.strptime(start_date, "%Y-%m-%d") - threshold_date = datetime.now() - relativedelta(days=7) + threshold_date = datetime.now() - timedelta(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.") return @@ -71,7 +70,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): 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(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") From a9964485dede85e095bb91b678089ef7b36b9d3c Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Thu, 9 Apr 2026 09:41:26 +0200 Subject: [PATCH 24/34] style:added ref for allowed statuses --- package/bin/audit_query_helper.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/package/bin/audit_query_helper.py b/package/bin/audit_query_helper.py index 73621f5..3645144 100644 --- a/package/bin/audit_query_helper.py +++ b/package/bin/audit_query_helper.py @@ -95,6 +95,8 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): 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"): break time.sleep(poll_sleep) From 4d2b5ee6086d5bea165a96f6031dde13820383f8 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Thu, 9 Apr 2026 09:42:28 +0200 Subject: [PATCH 25/34] fix:prevent TypeError: 'NoneType' --- package/bin/genesyscloud_client.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index e86dd50..1107ed7 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -46,9 +46,9 @@ def _fetch(self, api_instance, f_name: str, *args, **kwargs): # 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: From 1f3afc110e00fa498f8b02ff67b456fc9468826d Mon Sep 17 00:00:00 2001 From: ahoang-splunk <106760044+ahoang-splunk@users.noreply.github.com> Date: Mon, 8 Jun 2026 04:07:18 -0400 Subject: [PATCH 26/34] [FDSE-3233]: Add proxy support (#43) * Aded proxy support * Remove client * Implemented requested changes: change version to 0.3.1 and improved code readability --- globalConfig.json | 94 ++++++++++++++++++++- package/bin/actions_metrics_helper.py | 39 ++++++++- package/bin/audit_query_helper.py | 39 ++++++++- package/bin/conversations_details_helper.py | 39 ++++++++- package/bin/conversations_metrics_helper.py | 40 ++++++++- package/bin/edges_metrics_helper.py | 38 ++++++++- package/bin/edges_phones_helper.py | 38 ++++++++- package/bin/edges_trunks_metrics_helper.py | 38 ++++++++- package/bin/genesyscloud_client.py | 15 +++- package/bin/queue_observations_helper.py | 39 ++++++++- package/bin/status_page_metrics_helper.py | 51 ++++++++++- package/bin/user_aggregates_helper.py | 39 ++++++++- package/bin/user_routing_status_helper.py | 39 ++++++++- 13 files changed, 530 insertions(+), 18 deletions(-) diff --git a/globalConfig.json b/globalConfig.json index 4a28bc8..da907a1 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" }, @@ -1074,7 +1166,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/package/bin/actions_metrics_helper.py b/package/bin/actions_metrics_helper.py index 591b4bc..14e4faf 100644 --- a/package/bin/actions_metrics_helper.py +++ b/package/bin/actions_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 @@ -52,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) diff --git a/package/bin/audit_query_helper.py b/package/bin/audit_query_helper.py index 3645144..f4ec6f4 100644 --- a/package/bin/audit_query_helper.py +++ b/package/bin/audit_query_helper.py @@ -26,6 +26,42 @@ 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 validate_input(definition: smi.ValidationDefinition): pass @@ -60,7 +96,8 @@ 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 diff --git a/package/bin/conversations_details_helper.py b/package/bin/conversations_details_helper.py index 3d0c66d..99dd1c7 100644 --- a/package/bin/conversations_details_helper.py +++ b/package/bin/conversations_details_helper.py @@ -24,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. @@ -75,8 +111,9 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): 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] diff --git a/package/bin/conversations_metrics_helper.py b/package/bin/conversations_metrics_helper.py index d9fd44a..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 = ( diff --git a/package/bin/edges_metrics_helper.py b/package/bin/edges_metrics_helper.py index 908dd86..a704673 100644 --- a/package/bin/edges_metrics_helper.py +++ b/package/bin/edges_metrics_helper.py @@ -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] diff --git a/package/bin/edges_phones_helper.py b/package/bin/edges_phones_helper.py index c258623..e0981b6 100644 --- a/package/bin/edges_phones_helper.py +++ b/package/bin/edges_phones_helper.py @@ -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] diff --git a/package/bin/edges_trunks_metrics_helper.py b/package/bin/edges_trunks_metrics_helper.py index ce1006b..eb09266 100644 --- a/package/bin/edges_trunks_metrics_helper.py +++ b/package/bin/edges_trunks_metrics_helper.py @@ -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] diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index 1107ed7..9aabf38 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -6,13 +6,14 @@ from typing import List from PureCloudPlatformClientV2.rest import ApiException from PureCloudPlatformClientV2.api_client import ApiClient - +from PureCloudPlatformClientV2.configuration import Configuration class GenesysCloudClient: """ Interface with Genesys Cloud """ - 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] @@ -21,7 +22,15 @@ def __init__(self, logger: logging.Logger, client_id: str, client_secret: str, a self.logger.warning(f"Region {aws_region} not found: searching 'GENESYSCLOUD_HOST' env variable") self.host = os.environ.get("GENESYSCLOUD_HOST", None) # If host is none, default value will be "https://api.mypurecloud.com" - + + #Proxy + config = Configuration() + config.host = self.host + if proxy_url: + self.logger.info(f"Using proxy: {proxy_url}") + config.proxy = proxy_url + config.proxy_username = proxy_username + config.proxy_password = proxy_password self.client = ApiClient(self.host).get_client_credentials_token( client_id, client_secret ) diff --git a/package/bin/queue_observations_helper.py b/package/bin/queue_observations_helper.py index a5b3468..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 diff --git a/package/bin/status_page_metrics_helper.py b/package/bin/status_page_metrics_helper.py index e468ffc..fe83119 100644 --- a/package/bin/status_page_metrics_helper.py +++ b/package/bin/status_page_metrics_helper.py @@ -28,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() @@ -74,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" diff --git a/package/bin/user_aggregates_helper.py b/package/bin/user_aggregates_helper.py index 3281407..f114d66 100644 --- a/package/bin/user_aggregates_helper.py +++ b/package/bin/user_aggregates_helper.py @@ -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,9 +90,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 diff --git a/package/bin/user_routing_status_helper.py b/package/bin/user_routing_status_helper.py index 194bd53..34a0156 100644 --- a/package/bin/user_routing_status_helper.py +++ b/package/bin/user_routing_status_helper.py @@ -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 From 4da41b70c3f08f12e3cd2e60757f1a551748a56a Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 8 Jun 2026 10:37:30 +0200 Subject: [PATCH 27/34] fix:proxy configuration reset --- package/bin/genesyscloud_client.py | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index 9aabf38..f954fc1 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -8,11 +8,11 @@ from PureCloudPlatformClientV2.api_client import ApiClient from PureCloudPlatformClientV2.configuration import Configuration + class GenesysCloudClient: """ Interface with Genesys Cloud """ - 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): @@ -22,15 +22,18 @@ def __init__(self, logger: logging.Logger, client_id: str, client_secret: str, a self.logger.warning(f"Region {aws_region} not found: searching 'GENESYSCLOUD_HOST' env variable") self.host = os.environ.get("GENESYSCLOUD_HOST", None) # If host is none, default value will be "https://api.mypurecloud.com" - - #Proxy + + # 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}") - config.proxy = proxy_url - config.proxy_username = proxy_username - config.proxy_password = proxy_password + + # 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 ) From 8d07ae3a7240d50cdc90713453588aea3f82dd81 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Mon, 8 Jun 2026 10:39:10 +0200 Subject: [PATCH 28/34] docs:added proxy configuration info --- docs/ProxyConfiguration/index.md | 12 ++++++++++++ mkdocs.yml | 1 + 2 files changed, 13 insertions(+) create mode 100644 docs/ProxyConfiguration/index.md 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/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' From 29f588244951c6ded291e892ea39133635d3d894 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 9 Jun 2026 12:09:38 +0200 Subject: [PATCH 29/34] feat:added download events from url when available --- globalConfig.json | 45 ++++++++++++++-------- package/bin/audit_query_helper.py | 58 ++++++++++++++++++++++++---- package/bin/genesyscloud_client.py | 62 +++++++++++++++++++++++++++--- 3 files changed, 138 insertions(+), 27 deletions(-) diff --git a/globalConfig.json b/globalConfig.json index da907a1..c6d2187 100644 --- a/globalConfig.json +++ b/globalConfig.json @@ -357,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", @@ -374,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", diff --git a/package/bin/audit_query_helper.py b/package/bin/audit_query_helper.py index f4ec6f4..8ca5ea2 100644 --- a/package/bin/audit_query_helper.py +++ b/package/bin/audit_query_helper.py @@ -62,10 +62,30 @@ def get_account_proxy(logger, session_key: str): 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(): @@ -101,12 +121,27 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): checkpointer_key_name = normalized_input_name + # Setting a default start date of 7 days ago from now now = datetime.now(timezone.utc) - # Setting a fallback start date of 7 days ago from now - start_date = (now - timedelta(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 start_date + 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 } @@ -135,6 +170,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): # 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) @@ -143,6 +179,7 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): params = { "transaction_id": transaction_id, + "allow_redirect": True, "page_size": 500 } @@ -155,12 +192,19 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): 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=sourcetype, - time=entity.event_date.timestamp() + time=time_value ) ) event_counter += 1 diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index f954fc1..2b58bc4 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -2,8 +2,10 @@ 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 @@ -52,7 +54,28 @@ def _fetch(self, api_instance, f_name: str, *args, **kwargs): # 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.warn(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 @@ -99,8 +122,8 @@ def get(self, api_instance_name: str, function_name: str, *args, **kwargs): self.logger.error(f"Error: {e}") except ApiException as e: 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() + 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.") @@ -116,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. @@ -226,8 +278,8 @@ def post(self, api_instance_name: str, function_name: str, model_name: str, body except ApiException as e: 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() + 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.") From e5520a92b2bbe295367265b68a0a4f299023d25d Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 9 Jun 2026 12:10:48 +0200 Subject: [PATCH 30/34] feat:handled the exceed 31 days interval API error --- package/bin/conversations_details_helper.py | 46 ++++++++++++++++----- 1 file changed, 35 insertions(+), 11 deletions(-) diff --git a/package/bin/conversations_details_helper.py b/package/bin/conversations_details_helper.py index 99dd1c7..5525d2e 100644 --- a/package/bin/conversations_details_helper.py +++ b/package/bin/conversations_details_helper.py @@ -72,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() - timedelta(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): @@ -115,14 +130,23 @@ def stream_events(inputs: smi.InputDefinition, event_writer: smi.EventWriter): client = GenesysCloudClient( 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 = { From aafb7323a8e3a75dbd1c9f072905025e9750de0c Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 9 Jun 2026 12:11:14 +0200 Subject: [PATCH 31/34] docs:updated documentation accordingly --- docs/ConfigureAnalyticsInputs/index.md | 2 +- docs/ConfigureOperationalInputs/index.md | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) 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`. From eecb90a32546f5f205b979b3e64da54188cd4ffd Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 9 Jun 2026 12:29:23 +0200 Subject: [PATCH 32/34] ci(tests):bumped action test-reporter to latest version --- .github/workflows/tests.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index de039ca..6e9f8b7 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -277,7 +277,7 @@ jobs: name: mockoon-logs-${{ matrix.version }} path: logs/mockoon-${{ matrix.version }}.log - - uses: dorny/test-reporter@v2.6.0 + - uses: dorny/test-reporter@v3 if: always() with: name: Tests Results From e6299563f5bd8032d93f009b8b4eb6e1ec379e83 Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 9 Jun 2026 13:34:35 +0200 Subject: [PATCH 33/34] ci(tests):bumped action test-reporter to latest version --- .github/workflows/tests.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 6e9f8b7..3810427 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -145,7 +145,7 @@ jobs: name: mockoon-logs-${{ matrix.version }} path: logs/mockoon-${{ matrix.version }}.log - - uses: dorny/test-reporter@v2.6.0 + - uses: dorny/test-reporter@v3 if: always() with: name: Tests Results From bf7d4dd84aaa446e0b248081820225f466e2859f Mon Sep 17 00:00:00 2001 From: Erica Pescio Date: Tue, 9 Jun 2026 13:43:07 +0200 Subject: [PATCH 34/34] style:enhanced log messages --- package/bin/genesyscloud_client.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/package/bin/genesyscloud_client.py b/package/bin/genesyscloud_client.py index 2b58bc4..99c0dd7 100644 --- a/package/bin/genesyscloud_client.py +++ b/package/bin/genesyscloud_client.py @@ -61,7 +61,7 @@ def _fetch(self, api_instance, f_name: str, *args, **kwargs): # When getting audit query results API could return a redirect with downloadUrl. body = json.loads(e.body) if "downloadUrl" in body.keys(): - self.logger.warn(f"Got URL to download events from.") + 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) @@ -109,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)