diff --git a/orchestrator/activities/email.py b/orchestrator/activities/email.py index 35e5951..8bd6718 100644 --- a/orchestrator/activities/email.py +++ b/orchestrator/activities/email.py @@ -197,8 +197,8 @@ class Email(BaseActivity): email_group=group_name ).inc() - self.info(f"Email sent to {group_name}: {receivers}", - metadata=metadata) + self.info(f"Email sent to {group_name}: {receivers}", + metadata=metadata) self.info(f"Email sent for {mail_type} mail type.", metadata=metadata) diff --git a/orchestrator/activities/formatters.py b/orchestrator/activities/formatters.py index 0bfbc15..34732be 100644 --- a/orchestrator/activities/formatters.py +++ b/orchestrator/activities/formatters.py @@ -176,6 +176,10 @@ class Formatters(BaseActivity): Returns: dict[str, Any]: The formatted schedule config. """ + metadata = input_data["metadata"] + + self.info("Formatting schedule config...", metadata=metadata) + schedule_config = input_data['schedule_config'] config = {} @@ -189,6 +193,10 @@ class Formatters(BaseActivity): config[namespace][schedule_name] = updated_at + self.info("Formatted schedule config", metadata=metadata) + self.debug(json.dumps( + config, indent=4, sort_keys=True), metadata=metadata) + return config def compare_config_timestamps(self, diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index be0d40f..e7a3ac1 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -219,6 +219,10 @@ class MongoDB(BaseActivity): now = datetime.now() collection = self.database["orchestrated_schedules"] + self.info("Updating pipelines timestamps...", metadata=metadata) + + success_count = 0 + argument = [ {"schedule_name": pipeline["schedule_name"], "namespace": pipeline["namespace"]} @@ -231,6 +235,7 @@ class MongoDB(BaseActivity): data_filter, {"$set": {"updated_at": now}} ) + success_count += 1 except Exception as e: trace = traceback.format_exc() self.send_notification( @@ -244,6 +249,9 @@ class MongoDB(BaseActivity): self.error(trace, metadata=metadata) raise e + self.info( + f"Updated {success_count} of {len(updated_pipelines)} pipelines timestamps", metadata=metadata) + @activity.defn(name="create_pipelines_timestamps") async def create_pipelines_timestamps(self, input_data: dict[str, Any]) -> None: """ @@ -255,6 +263,10 @@ class MongoDB(BaseActivity): metadata = input_data.get("metadata", {}) collection = self.database["orchestrated_schedules"] + self.info("Creating pipelines timestamps...", metadata=metadata) + + success_count = 0 + now = datetime.now() argument = [ @@ -268,6 +280,7 @@ class MongoDB(BaseActivity): try: if data_filter: collection.insert_many(data_filter) + success_count += 1 except Exception as e: trace = traceback.format_exc() self.send_notification( @@ -281,6 +294,9 @@ class MongoDB(BaseActivity): self.error(trace, metadata=metadata) raise e + self.info( + f"Created {success_count} of {len(created_pipelines)} pipelines timestamps", metadata=metadata) + @activity.defn(name="delete_pipelines_timestamps") async def delete_pipelines_timestamps(self, input_data: dict[str, Any]) -> None: """ @@ -292,6 +308,10 @@ class MongoDB(BaseActivity): metadata = input_data.get("metadata", {}) collection = self.database["orchestrated_schedules"] + self.info("Deleting pipelines timestamps...", metadata=metadata) + + success_count = 0 + argument = [ {"schedule_name": pipeline["schedule_name"], "namespace": pipeline["namespace"]} @@ -301,6 +321,7 @@ class MongoDB(BaseActivity): try: collection.delete_many(data_filter) + success_count += 1 except Exception as e: trace = traceback.format_exc() self.send_notification( @@ -314,6 +335,9 @@ class MongoDB(BaseActivity): self.error(trace, metadata=metadata) raise e + self.info( + f"Deleted {success_count} of {len(deleted_pipelines)} pipelines timestamps", metadata=metadata) + @activity.defn(name="create_collection_with_ttl_index") async def create_collection_with_ttl_index(self, input_data: dict[str, Any]) -> None: """ @@ -331,6 +355,9 @@ class MongoDB(BaseActivity): collection_names = self.database.list_collection_names() + created_collections = [] + created_indexes = [] + for pipeline_name, pipeline_config in pipelines.items(): collection = pipeline_config["topic"] @@ -338,6 +365,7 @@ class MongoDB(BaseActivity): # Check if collection exists if collection not in collection_names: self.database.create_collection(collection) + created_collections.append(collection) collection = self.database[collection] # Check if TTL index exists @@ -355,6 +383,7 @@ class MongoDB(BaseActivity): expireAfterSeconds=self.ttl_index_seconds, background=True ) + created_indexes.append(collection) except Exception as e: trace = traceback.format_exc() @@ -369,6 +398,20 @@ class MongoDB(BaseActivity): self.error(trace, metadata=metadata) raise e + self.info( + f"Created {len(created_collections)} collections and {len(created_indexes)} indexes", + metadata=metadata + ) + + self.debug( + f"Created collections: {created_collections}", + metadata=metadata + ) + self.debug( + f"Created indexes: {created_indexes}", + metadata=metadata + ) + @activity.defn(name="load_latest_data") async def load_latest_data(self, input_data: dict[str, Any]) -> list[dict[str, Any]]: """ diff --git a/orchestrator/activities/slot_manager.py b/orchestrator/activities/slot_manager.py index 14aa70c..a7dab02 100644 --- a/orchestrator/activities/slot_manager.py +++ b/orchestrator/activities/slot_manager.py @@ -134,6 +134,8 @@ class SlotManager(Redis): report = {} + success_count = 0 + for slot in to_insert: try: self.set(f"slot:opc_tags:{slot}", @@ -142,6 +144,7 @@ class SlotManager(Redis): "success": True, "message": "Slot updated successfully" } + success_count += 1 except Exception as e: self.error( f"Failed to update slot {slot}: {str(e)}", metadata=metadata) @@ -150,7 +153,8 @@ class SlotManager(Redis): "message": str(e) } - self.info(f"Updated {len(to_insert)} OPC slots", metadata=metadata) + self.info( + f"Updated {success_count} of {len(to_insert)} OPC slots", metadata=metadata) self.debug( f"Report: \n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata) @@ -174,8 +178,12 @@ class SlotManager(Redis): to_delete = input_data['to_delete'] metadata = input_data.get("metadata", {}) + self.info("Deleting OPC slots...", metadata=metadata) + report = {} + success_count = 0 + for slot in to_delete: try: self.redis_client.delete(f"slot:opc_tags:{slot}") @@ -183,6 +191,7 @@ class SlotManager(Redis): "success": True, "message": "Slot deleted successfully" } + success_count += 1 except Exception as e: self.error( f"Failed to delete slot {slot}: {str(e)}", metadata=metadata) @@ -191,7 +200,8 @@ class SlotManager(Redis): "message": str(e) } - self.info(f"Deleted {len(to_delete)} OPC slots", metadata=metadata) + self.info( + f"Deleted {success_count} of {len(to_delete)} OPC slots", metadata=metadata) self.debug( f"Report: \n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata) diff --git a/orchestrator/activities/temporal_manager.py b/orchestrator/activities/temporal_manager.py index d6f64bf..3c9e93a 100644 --- a/orchestrator/activities/temporal_manager.py +++ b/orchestrator/activities/temporal_manager.py @@ -74,7 +74,9 @@ class TemporalManager(BaseActivity): input_data: - orchestrated_schedules (dict[str, Any]): The orchestrated schedules to compare. """ - metadata = input_data.get("metadata", {}) + metadata = input_data["metadata"] + + remove_count = 0 self.info("Getting orchestrated schedules...", metadata=metadata) @@ -101,6 +103,8 @@ class TemporalManager(BaseActivity): await handle.delete() + remove_count += 1 + except Exception as e: trace = traceback.format_exc() self.send_notification( @@ -114,6 +118,9 @@ class TemporalManager(BaseActivity): self.error(trace, metadata=metadata) raise e + self.info( + f"Removed {remove_count} schedules", metadata=metadata) + @activity.defn(name="create_schedules") async def create_schedules(self, input_data: dict[str, Any]) -> list[dict[str, Any]]: """ @@ -133,6 +140,10 @@ class TemporalManager(BaseActivity): report = [] + success_count = 0 + + self.info("Creating schedules...", metadata=metadata) + for namespace, schedules in schedules_to_create.items(): client = self.temporal_clients.get(namespace) @@ -192,6 +203,8 @@ class TemporalManager(BaseActivity): "success": True, "message": "Schedule created successfully" }) + + success_count += 1 except Exception as e: self.error( f"Failed to create schedule {schedule_name}: {str(e)}", metadata=metadata) @@ -203,7 +216,7 @@ class TemporalManager(BaseActivity): }) self.info( - f"Processed {len(schedules_to_create)} schedules", metadata=metadata) + f"Created {success_count} of {len(schedules_to_create)} schedules", metadata=metadata) self.debug( f"\n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata) @@ -228,6 +241,10 @@ class TemporalManager(BaseActivity): metadata = input_data.get("metadata", {}) report = [] + success_count = 0 + + self.info("Updating schedules...", metadata=metadata) + for namespace, schedules in schedules_to_update.items(): client = self.temporal_clients.get(namespace) @@ -276,6 +293,8 @@ class TemporalManager(BaseActivity): "success": True, "message": "Schedule updated successfully" }) + + success_count += 1 except Exception as e: self.error( f"Failed to update schedule {schedule_name}: {str(e)}", metadata=metadata) @@ -287,7 +306,7 @@ class TemporalManager(BaseActivity): }) self.info( - f"Processed {len(schedules_to_update)} schedules", metadata=metadata) + f"Updated {success_count} of {len(schedules_to_update)} schedules", metadata=metadata) self.debug( f"\n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata) @@ -312,6 +331,10 @@ class TemporalManager(BaseActivity): metadata = input_data.get("metadata", {}) report = [] + success_count = 0 + + self.info("Deleting schedules...", metadata=metadata) + for namespace, schedules in schedules_to_delete.items(): client = self.temporal_clients.get(namespace) @@ -337,6 +360,8 @@ class TemporalManager(BaseActivity): "success": True, "message": "Schedule deleted successfully" }) + + success_count += 1 except Exception as e: trace = traceback.format_exc() self.error( @@ -350,7 +375,7 @@ class TemporalManager(BaseActivity): }) self.info( - f"Processed {len(schedules_to_delete)} schedules", metadata=metadata) + f"Deleted {success_count} of {len(schedules_to_delete)} schedules", metadata=metadata) self.debug( f"\n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata) diff --git a/orchestrator/worker/worker.py b/orchestrator/worker/worker.py index d5d967e..9439a38 100644 --- a/orchestrator/worker/worker.py +++ b/orchestrator/worker/worker.py @@ -41,12 +41,21 @@ async def main(): namespace = os.getenv('TEMPORAL_NAMESPACE', 'default') logger = get_logger(__name__) - logger.info(f'Starting Worker with POD_ID: {POD_ID}') + metadata = { + 'pod_id': POD_ID, + 'model_name': '-', + 'model_id': '-', + 'workflow_name': '-', + 'schedule_name': '-', + } - logger.info("Starting prometheus client...") + logger.custom_info( + f'Starting Worker with POD_ID: {POD_ID}', metadata=metadata) + + logger.custom_info("Starting prometheus client...", metadata=metadata) start_prometheus_server() - logger.info('Starting Notification Handler...') + logger.custom_info('Starting Notification Handler...', metadata=metadata) mongo_config = build_mongodb_config() notification_handler = NotificationHandler( @@ -56,7 +65,8 @@ async def main(): project_name=os.getenv('PROJECT_NAME', 'orchestrator'), ) - logger.info(f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...') + logger.custom_info( + f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...', metadata=metadata) new_runtime = Runtime( telemetry=TelemetryConfig( @@ -65,7 +75,8 @@ async def main(): ) ) - logger.info(f'Starting Temporal Client at {host}:{namespace}') + logger.custom_info( + f'Starting Temporal Client at {host}:{namespace}', metadata=metadata) temporal_client = await client.Client.connect( target_host=host, @@ -73,7 +84,7 @@ async def main(): runtime=new_runtime ) - logger.info('Starting Activities...') + logger.custom_info('Starting Activities...', metadata=metadata) activities = Activities( temporal_config=build_temporal_config(), @@ -87,7 +98,7 @@ async def main(): await activities.connect_to_temporal() - logger.info('Starting Workers...') + logger.custom_info('Starting Workers...', metadata=metadata) workers = [ Worker( @@ -175,7 +186,7 @@ async def main(): for w in workers: handlers.append(w.run()) - logger.info('Workers started successfully') + logger.custom_info('Workers started successfully', metadata=metadata) try: # This will run the workers and wait for them to complete.