SIENTIAPDE-1163
refactor: enhance logging in activities to include metadata for improved context in notifications and error handling
This commit is contained in:
@@ -42,7 +42,9 @@ class Formatters(BaseActivity):
|
||||
- dict[str, Any]: The schedule config dictionary
|
||||
"""
|
||||
|
||||
self.logger.info("Processing schedules...")
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.info("Processing schedules...", metadata=metadata)
|
||||
|
||||
pipelines = input_data['pipelines']
|
||||
|
||||
@@ -66,10 +68,9 @@ class Formatters(BaseActivity):
|
||||
"updated_at", datetime.now().strftime(DEFAULT_DATE_FORMAT))
|
||||
}
|
||||
|
||||
self.logger.info("Processed schedules")
|
||||
|
||||
self.logger.debug(json.dumps(
|
||||
schedule_config, indent=4, sort_keys=True))
|
||||
self.info("Processed schedules", metadata=metadata)
|
||||
self.debug(json.dumps(
|
||||
schedule_config, indent=4, sort_keys=True), metadata=metadata)
|
||||
|
||||
return schedule_config
|
||||
|
||||
@@ -91,7 +92,9 @@ class Formatters(BaseActivity):
|
||||
- dict[str, Any]: The slot config dictionary
|
||||
"""
|
||||
|
||||
self.logger.info("Processing slots...")
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.info("Processing slots...", metadata=metadata)
|
||||
|
||||
pipelines = input_data['pipelines']
|
||||
opc_servers_list = input_data['opc_servers']
|
||||
@@ -124,9 +127,9 @@ class Formatters(BaseActivity):
|
||||
slot_config = build_tag_config(
|
||||
tag, slot_config.copy(), opc_servers, number_of_slots)
|
||||
|
||||
self.logger.info("Processed slots")
|
||||
self.logger.debug(json.dumps(
|
||||
slot_config, indent=4, sort_keys=True))
|
||||
self.info("Processed slots", metadata=metadata)
|
||||
self.debug(json.dumps(
|
||||
slot_config, indent=4, sort_keys=True), metadata=metadata)
|
||||
|
||||
return slot_config
|
||||
|
||||
@@ -158,7 +161,7 @@ class Formatters(BaseActivity):
|
||||
def compare_config_timestamps(self,
|
||||
schedules: dict[str, Any], current_schedules: dict[str, Any],
|
||||
to_update: dict[str, Any], to_create: dict[str, Any],
|
||||
namespace: str):
|
||||
namespace: str, metadata: dict[str, Any]):
|
||||
"""
|
||||
Compares the timestamps of the schedule and the current schedule.
|
||||
"""
|
||||
@@ -170,8 +173,8 @@ class Formatters(BaseActivity):
|
||||
|
||||
old_timestamp = current_schedules[schedule_name]
|
||||
|
||||
self.logger.debug(
|
||||
f"Comparing schedule {schedule_name}:{update_timestamp} vs {old_timestamp}")
|
||||
self.debug(
|
||||
f"Comparing schedule {schedule_name}:{update_timestamp} vs {old_timestamp}", metadata=metadata)
|
||||
|
||||
if update_timestamp > old_timestamp:
|
||||
to_update[namespace][schedule_name] = schedule
|
||||
@@ -197,7 +200,9 @@ class Formatters(BaseActivity):
|
||||
- dict[str, Any]: The schedule config dictionary
|
||||
"""
|
||||
|
||||
self.logger.info("Creating schedule config...")
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.info("Creating schedule config...", metadata=metadata)
|
||||
|
||||
current_schedule_config = input_data['current_schedule_config']
|
||||
schedule_config = input_data['schedule_config']
|
||||
@@ -218,7 +223,7 @@ class Formatters(BaseActivity):
|
||||
for namespace, schedules in schedule_config.items():
|
||||
current_schedules = current_schedule_config.get(namespace, {})
|
||||
self.compare_config_timestamps(
|
||||
schedules, current_schedules, to_update, to_create, namespace)
|
||||
schedules, current_schedules, to_update, to_create, namespace, metadata)
|
||||
|
||||
for namespace, schedules in current_schedule_config.items():
|
||||
for schedule_name in schedules:
|
||||
@@ -231,9 +236,9 @@ class Formatters(BaseActivity):
|
||||
"to_delete": to_delete
|
||||
}
|
||||
|
||||
self.logger.info("Created schedule config")
|
||||
self.logger.debug(json.dumps(
|
||||
output, indent=4, sort_keys=True))
|
||||
self.info("Created schedule config", metadata=metadata)
|
||||
self.debug(json.dumps(
|
||||
output, indent=4, sort_keys=True), metadata=metadata)
|
||||
|
||||
return output
|
||||
|
||||
@@ -256,10 +261,10 @@ class Formatters(BaseActivity):
|
||||
- dict[str, Any]: The slot config dictionary
|
||||
"""
|
||||
|
||||
self.logger.info("Creating slot config...")
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.info("Creating slot config...", metadata=metadata)
|
||||
|
||||
current_slot_config = input_data['current_slot_config']
|
||||
slot_config = input_data['slot_config']
|
||||
to_delete = []
|
||||
@@ -276,9 +281,9 @@ class Formatters(BaseActivity):
|
||||
"to_insert": slot_config
|
||||
}
|
||||
|
||||
self.logger.info("Created slot config")
|
||||
self.logger.debug(json.dumps(
|
||||
output, indent=4, sort_keys=True))
|
||||
self.info("Created slot config", metadata=metadata)
|
||||
self.debug(json.dumps(
|
||||
output, indent=4, sort_keys=True), metadata=metadata)
|
||||
|
||||
return output
|
||||
|
||||
@@ -333,11 +338,10 @@ class Formatters(BaseActivity):
|
||||
- updated_schedules (dict[str, Any]): The updated schedules.
|
||||
- deleted_schedules (list[str]): The deleted schedules.
|
||||
"""
|
||||
|
||||
self.logger.info("Reporting orchestration...")
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.info("Reporting orchestration...", metadata=metadata)
|
||||
|
||||
created_schedules = input_data['created_schedules']
|
||||
updated_schedules = input_data['updated_schedules']
|
||||
deleted_schedules = input_data['deleted_schedules']
|
||||
@@ -414,10 +418,10 @@ class Formatters(BaseActivity):
|
||||
- deleted_slots (list[str]): The deleted slots.
|
||||
"""
|
||||
|
||||
self.logger.info("Reporting orchestration...")
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.info("Reporting orchestration...", metadata=metadata)
|
||||
|
||||
inserted_slots = input_data['inserted_slots']
|
||||
deleted_slots = input_data['deleted_slots']
|
||||
|
||||
|
||||
@@ -100,8 +100,8 @@ class MongoDB(BaseActivity):
|
||||
|
||||
filters = query.get("filters", {})
|
||||
|
||||
self.logger.info(
|
||||
f"Loading documents from collection '{collection_name}' with filters: {filters}")
|
||||
self.info(
|
||||
f"Loading documents from collection '{collection_name}' with filters: {filters}", metadata=metadata)
|
||||
|
||||
try:
|
||||
collection = self.database[collection_name]
|
||||
@@ -110,12 +110,11 @@ class MongoDB(BaseActivity):
|
||||
|
||||
documents = clear_mongo_id(documents)
|
||||
|
||||
self.logger.info(
|
||||
f"Loaded {len(documents)} documents from collection '{collection_name}'")
|
||||
self.info(
|
||||
f"Loaded {len(documents)} documents from collection '{collection_name}'", metadata=metadata)
|
||||
|
||||
self.logger.debug(
|
||||
f"Documents loaded: {documents}"
|
||||
)
|
||||
self.debug(
|
||||
f"Documents loaded: {documents}", metadata=metadata)
|
||||
|
||||
return documents
|
||||
|
||||
@@ -129,7 +128,7 @@ class MongoDB(BaseActivity):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
|
||||
raise e
|
||||
|
||||
@@ -158,8 +157,8 @@ class MongoDB(BaseActivity):
|
||||
raise ValueError("Aggregation must be provided.")
|
||||
aggregation.append({"$project": {"_id": 0}})
|
||||
|
||||
self.logger.info(
|
||||
f"Aggregating documents from collection '{collection_name}' with aggregation: {aggregation}")
|
||||
self.info(
|
||||
f"Aggregating documents from collection '{collection_name}' with aggregation: {aggregation}", metadata=metadata)
|
||||
|
||||
try:
|
||||
collection = self.database[collection_name]
|
||||
@@ -169,12 +168,11 @@ class MongoDB(BaseActivity):
|
||||
|
||||
aggregated_documents = clear_mongo_id(aggregated_documents)
|
||||
|
||||
self.logger.info(
|
||||
f"Aggregated {len(aggregated_documents)} documents from collection '{collection_name}'")
|
||||
self.info(
|
||||
f"Aggregated {len(aggregated_documents)} documents from collection '{collection_name}'", metadata=metadata)
|
||||
|
||||
self.logger.debug(
|
||||
f"Aggregation result: {aggregated_documents}"
|
||||
)
|
||||
self.debug(
|
||||
f"Aggregation result: {aggregated_documents}", metadata=metadata)
|
||||
|
||||
return aggregated_documents
|
||||
|
||||
@@ -188,7 +186,7 @@ class MongoDB(BaseActivity):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
|
||||
raise e
|
||||
|
||||
@@ -226,7 +224,7 @@ class MongoDB(BaseActivity):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
raise e
|
||||
|
||||
@activity.defn(name="create_pipelines_timestamps")
|
||||
@@ -262,7 +260,7 @@ class MongoDB(BaseActivity):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
raise e
|
||||
|
||||
@activity.defn(name="delete_pipelines_timestamps")
|
||||
@@ -295,5 +293,5 @@ class MongoDB(BaseActivity):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
raise e
|
||||
|
||||
@@ -30,7 +30,7 @@ class SlotManager(Redis):
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.logger.info("Loading OPC slots...")
|
||||
self.info("Loading OPC slots...", metadata=metadata)
|
||||
|
||||
opc_slots = {}
|
||||
|
||||
@@ -38,7 +38,7 @@ class SlotManager(Redis):
|
||||
|
||||
slot_keys = self.redis_client.keys("slot:opc_tags:*")
|
||||
|
||||
self.logger.debug(f"Slot keys: {slot_keys}")
|
||||
self.debug(f"Slot keys: {slot_keys}", metadata=metadata)
|
||||
|
||||
if slot_keys:
|
||||
if isinstance(slot_keys[0], bytes):
|
||||
@@ -58,10 +58,10 @@ class SlotManager(Redis):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
raise e
|
||||
|
||||
self.logger.info(f"Loaded {len(opc_slots)} OPC slots")
|
||||
self.info(f"Loaded {len(opc_slots)} OPC slots", metadata=metadata)
|
||||
|
||||
return opc_slots
|
||||
|
||||
@@ -76,16 +76,17 @@ class SlotManager(Redis):
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.logger.info("Loading active ingestors...")
|
||||
self.info("Loading active ingestors...", metadata=metadata)
|
||||
|
||||
try:
|
||||
|
||||
active_ingestors = self.redis_client.keys("heartbeat:ingestor:*")
|
||||
|
||||
self.logger.info(
|
||||
f"Loaded {len(active_ingestors)} active ingestors")
|
||||
self.info(
|
||||
f"Loaded {len(active_ingestors)} active ingestors", metadata=metadata)
|
||||
|
||||
self.logger.debug(f"Active ingestors: \n {active_ingestors}")
|
||||
self.debug(
|
||||
f"Active ingestors: \n {active_ingestors}", metadata=metadata)
|
||||
|
||||
ingestors = []
|
||||
|
||||
@@ -106,7 +107,7 @@ class SlotManager(Redis):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
raise e
|
||||
|
||||
@activity.defn(name="update_slots")
|
||||
@@ -124,7 +125,8 @@ class SlotManager(Redis):
|
||||
"""
|
||||
|
||||
to_insert = input_data['to_insert']
|
||||
self.logger.info("Updating OPC slots...")
|
||||
metadata = input_data.get("metadata", {})
|
||||
self.info("Updating OPC slots...", metadata=metadata)
|
||||
|
||||
report = {}
|
||||
|
||||
@@ -137,16 +139,17 @@ class SlotManager(Redis):
|
||||
"message": "Slot updated successfully"
|
||||
}
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to update slot {slot}: {str(e)}")
|
||||
self.error(
|
||||
f"Failed to update slot {slot}: {str(e)}", metadata=metadata)
|
||||
report[slot] = {
|
||||
"success": False,
|
||||
"message": str(e)
|
||||
}
|
||||
|
||||
self.logger.info(f"Updated {len(to_insert)} OPC slots")
|
||||
self.info(f"Updated {len(to_insert)} OPC slots", metadata=metadata)
|
||||
|
||||
self.logger.debug(
|
||||
f"Report: \n {json.dumps(report, indent=4, sort_keys=True)}")
|
||||
self.debug(
|
||||
f"Report: \n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata)
|
||||
|
||||
return report
|
||||
|
||||
@@ -165,7 +168,7 @@ class SlotManager(Redis):
|
||||
"""
|
||||
|
||||
to_delete = input_data['to_delete']
|
||||
self.logger.info("Deleting OPC slots...")
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
report = {}
|
||||
|
||||
@@ -177,15 +180,16 @@ class SlotManager(Redis):
|
||||
"message": "Slot deleted successfully"
|
||||
}
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to delete slot {slot}: {str(e)}")
|
||||
self.error(
|
||||
f"Failed to delete slot {slot}: {str(e)}", metadata=metadata)
|
||||
report[slot] = {
|
||||
"success": False,
|
||||
"message": str(e)
|
||||
}
|
||||
|
||||
self.logger.info(f"Deleted {len(to_delete)} OPC slots")
|
||||
self.info(f"Deleted {len(to_delete)} OPC slots", metadata=metadata)
|
||||
|
||||
self.logger.debug(
|
||||
f"Report: \n {json.dumps(report, indent=4, sort_keys=True)}")
|
||||
self.debug(
|
||||
f"Report: \n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata)
|
||||
|
||||
return report
|
||||
|
||||
@@ -70,18 +70,18 @@ class TemporalManager(BaseActivity):
|
||||
input_data:
|
||||
- orchestrated_schedules (dict[str, Any]): The orchestrated schedules to compare.
|
||||
"""
|
||||
self.logger.info("Getting orchestrated schedules...")
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
self.info("Getting orchestrated schedules...", metadata=metadata)
|
||||
|
||||
orchestrated_schedules = input_data.get('orchestrated_schedules', {})
|
||||
|
||||
for namespace, client in self.temporal_clients.items():
|
||||
try:
|
||||
schedules = orchestrated_schedules.get(namespace, {})
|
||||
|
||||
self.logger.info(
|
||||
f"Getting orchestrated schedules for {namespace}")
|
||||
self.info(
|
||||
f"Getting orchestrated schedules for {namespace}", metadata=metadata)
|
||||
|
||||
async for schedule in await client.list_schedules():
|
||||
search_attrs = getattr(schedule, "search_attributes", {})
|
||||
@@ -89,8 +89,8 @@ class TemporalManager(BaseActivity):
|
||||
schedule_id = schedule.id
|
||||
|
||||
if schedule_id not in schedules:
|
||||
self.logger.info(
|
||||
f"Schedule {schedule_id} not found in mongo db, cleaning up")
|
||||
self.info(
|
||||
f"Schedule {schedule_id} not found in mongo db, cleaning up", metadata=metadata)
|
||||
|
||||
handle = client.get_schedule_handle(
|
||||
schedule_id)
|
||||
@@ -107,7 +107,7 @@ class TemporalManager(BaseActivity):
|
||||
level=NotificationLevel.ERROR,
|
||||
attachment_content=trace
|
||||
)
|
||||
self.logger.error(trace)
|
||||
self.error(trace, metadata=metadata)
|
||||
raise e
|
||||
|
||||
@activity.defn(name="create_schedules")
|
||||
@@ -125,6 +125,7 @@ class TemporalManager(BaseActivity):
|
||||
"""
|
||||
|
||||
schedules_to_create = input_data['schedules']
|
||||
metadata = input_data.get("metadata", {})
|
||||
|
||||
report = []
|
||||
|
||||
@@ -153,9 +154,10 @@ class TemporalManager(BaseActivity):
|
||||
workflow_type = schedule['workflow_type']
|
||||
|
||||
try:
|
||||
self.logger.debug(f"Creating schedule {schedule_name}:")
|
||||
self.logger.debug(
|
||||
f"{json.dumps(schedule, indent=4, sort_keys=True)}")
|
||||
self.debug(
|
||||
f"Creating schedule {schedule_name}:", metadata=metadata)
|
||||
self.debug(
|
||||
f"{json.dumps(schedule, indent=4, sort_keys=True)}", metadata=metadata)
|
||||
|
||||
await client.create_schedule(
|
||||
schedule_name,
|
||||
@@ -186,8 +188,8 @@ class TemporalManager(BaseActivity):
|
||||
"message": "Schedule created successfully"
|
||||
})
|
||||
except Exception as e:
|
||||
self.logger.error(
|
||||
f"Failed to create schedule {schedule_name}: {str(e)}")
|
||||
self.error(
|
||||
f"Failed to create schedule {schedule_name}: {str(e)}", metadata=metadata)
|
||||
report.append({
|
||||
"namespace": namespace,
|
||||
"schedule_name": schedule_name,
|
||||
@@ -195,10 +197,11 @@ class TemporalManager(BaseActivity):
|
||||
"message": str(e)
|
||||
})
|
||||
|
||||
self.logger.info(f"Processed {len(schedules_to_create)} schedules")
|
||||
self.info(
|
||||
f"Processed {len(schedules_to_create)} schedules", metadata=metadata)
|
||||
|
||||
self.logger.debug(
|
||||
f"\n {json.dumps(report, indent=4, sort_keys=True)}")
|
||||
self.debug(
|
||||
f"\n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata)
|
||||
|
||||
return report
|
||||
|
||||
@@ -217,7 +220,7 @@ class TemporalManager(BaseActivity):
|
||||
"""
|
||||
|
||||
schedules_to_update = input_data['schedules']
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
report = []
|
||||
|
||||
for namespace, schedules in schedules_to_update.items():
|
||||
@@ -239,12 +242,12 @@ class TemporalManager(BaseActivity):
|
||||
async def update_schedule(input_data: ScheduleUpdateInput) -> ScheduleUpdate: # NOSONAR
|
||||
schedule_action = input_data.description.schedule.action
|
||||
|
||||
self.logger.debug("Updating schedule:")
|
||||
self.debug("Updating schedule:", metadata=metadata)
|
||||
|
||||
if hasattr(schedule_action, "args"):
|
||||
self.logger.debug("New schedule:")
|
||||
self.logger.debug(
|
||||
f"{json.dumps(schedule, indent=4, sort_keys=True)}") # NOSONAR
|
||||
self.debug("New schedule:", metadata=metadata)
|
||||
self.debug(
|
||||
f"{json.dumps(schedule, indent=4, sort_keys=True)}", metadata=metadata) # NOSONAR
|
||||
|
||||
schedule_action.args = [schedule]
|
||||
|
||||
@@ -269,8 +272,8 @@ class TemporalManager(BaseActivity):
|
||||
"message": "Schedule updated successfully"
|
||||
})
|
||||
except Exception as e:
|
||||
self.logger.error(
|
||||
f"Failed to update schedule {schedule_name}: {str(e)}")
|
||||
self.error(
|
||||
f"Failed to update schedule {schedule_name}: {str(e)}", metadata=metadata)
|
||||
report.append({
|
||||
"namespace": namespace,
|
||||
"schedule_name": schedule_name,
|
||||
@@ -278,10 +281,11 @@ class TemporalManager(BaseActivity):
|
||||
"message": str(e)
|
||||
})
|
||||
|
||||
self.logger.info(f"Processed {len(schedules_to_update)} schedules")
|
||||
self.info(
|
||||
f"Processed {len(schedules_to_update)} schedules", metadata=metadata)
|
||||
|
||||
self.logger.debug(
|
||||
f"\n {json.dumps(report, indent=4, sort_keys=True)}")
|
||||
self.debug(
|
||||
f"\n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata)
|
||||
|
||||
return report
|
||||
|
||||
@@ -300,7 +304,7 @@ class TemporalManager(BaseActivity):
|
||||
"""
|
||||
|
||||
schedules_to_delete = input_data['schedules']
|
||||
|
||||
metadata = input_data.get("metadata", {})
|
||||
report = []
|
||||
|
||||
for namespace, schedules in schedules_to_delete.items():
|
||||
@@ -329,8 +333,8 @@ class TemporalManager(BaseActivity):
|
||||
"message": "Schedule deleted successfully"
|
||||
})
|
||||
except Exception as e:
|
||||
self.logger.error(
|
||||
f"Failed to delete schedule {schedule_name}: {str(e)}")
|
||||
self.error(
|
||||
f"Failed to delete schedule {schedule_name}: {str(e)}", metadata=metadata)
|
||||
report.append({
|
||||
"namespace": namespace,
|
||||
"schedule_name": schedule_name,
|
||||
@@ -338,9 +342,10 @@ class TemporalManager(BaseActivity):
|
||||
"message": str(e)
|
||||
})
|
||||
|
||||
self.logger.info(f"Processed {len(schedules_to_delete)} schedules")
|
||||
self.info(
|
||||
f"Processed {len(schedules_to_delete)} schedules", metadata=metadata)
|
||||
|
||||
self.logger.debug(
|
||||
f"\n {json.dumps(report, indent=4, sort_keys=True)}")
|
||||
self.debug(
|
||||
f"\n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata)
|
||||
|
||||
return report
|
||||
|
||||
Reference in New Issue
Block a user