diff --git a/orchestrator/activities/formatters.py b/orchestrator/activities/formatters.py index e04434b..a5c8053 100644 --- a/orchestrator/activities/formatters.py +++ b/orchestrator/activities/formatters.py @@ -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'] diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index fadc2de..1c1510f 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -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 diff --git a/orchestrator/activities/slot_manager.py b/orchestrator/activities/slot_manager.py index d55990b..404c8d8 100644 --- a/orchestrator/activities/slot_manager.py +++ b/orchestrator/activities/slot_manager.py @@ -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 diff --git a/orchestrator/activities/temporal_manager.py b/orchestrator/activities/temporal_manager.py index bdca7dd..f8461eb 100644 --- a/orchestrator/activities/temporal_manager.py +++ b/orchestrator/activities/temporal_manager.py @@ -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