Merge pull request #16 from Aignosi/SIENTIAPDE-1199-revisar-e-testar-observabilidade
Sientiapde 1199 revisar e testar observabilidade
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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]]:
|
||||
"""
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -27,11 +27,6 @@ class TemporalManager(BaseActivity):
|
||||
self.laborious_namespace = laborious_namespace
|
||||
self.temporal_clients = {}
|
||||
|
||||
self.schedule_handles = {
|
||||
self.scouter_namespace: {},
|
||||
self.laborious_namespace: {}
|
||||
}
|
||||
|
||||
self.model_id_id_key = SearchAttributeKey.for_keyword("model_id")
|
||||
self.model_name_id_key = SearchAttributeKey.for_keyword("model_name")
|
||||
self.orchestrated_id_key = SearchAttributeKey.for_keyword(
|
||||
@@ -74,7 +69,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 +98,8 @@ class TemporalManager(BaseActivity):
|
||||
|
||||
await handle.delete()
|
||||
|
||||
remove_count += 1
|
||||
|
||||
except Exception as e:
|
||||
trace = traceback.format_exc()
|
||||
self.send_notification(
|
||||
@@ -114,6 +113,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 +135,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 +198,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 +211,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 +236,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 +288,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 +301,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 +326,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)
|
||||
|
||||
@@ -329,14 +347,14 @@ class TemporalManager(BaseActivity):
|
||||
|
||||
await handler.delete()
|
||||
|
||||
del self.schedule_handles[namespace][schedule_name]
|
||||
|
||||
report.append({
|
||||
"namespace": namespace,
|
||||
"schedule_name": schedule_name,
|
||||
"success": True,
|
||||
"message": "Schedule deleted successfully"
|
||||
})
|
||||
|
||||
success_count += 1
|
||||
except Exception as e:
|
||||
trace = traceback.format_exc()
|
||||
self.error(
|
||||
@@ -350,7 +368,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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -354,7 +354,8 @@ async def test_format_schedule_config(formatters):
|
||||
"updated_at": "2021-01-01"},
|
||||
{"namespace": "test_namespace2", "schedule_name": "test2",
|
||||
"updated_at": "2021-01-02"}
|
||||
]
|
||||
],
|
||||
**metadata
|
||||
}
|
||||
|
||||
result = await formatters.format_schedule_config(input_data)
|
||||
|
||||
@@ -76,7 +76,8 @@ async def test_normalize_schedules(temporal_manager):
|
||||
"orchestrated_schedules": {
|
||||
"scouter": {"test-scouter": "2021-01-01"},
|
||||
"laborious": {"test-schedule-id1": "2021-01-01"}
|
||||
}
|
||||
},
|
||||
**metadata
|
||||
}
|
||||
|
||||
temporal_manager.temporal_clients['scouter'].list_schedules = AsyncMock(
|
||||
|
||||
Reference in New Issue
Block a user