Merge pull request #13 from Aignosi/SIENTIAPDE-1174-mapear-e-implementar-metricas-a-serem-criadas

Sientiapde 1174 mapear e implementar metricas a serem criadas
This commit is contained in:
Bruno Domingues
2025-08-13 16:34:35 -03:00
committed by GitHub
12 changed files with 131 additions and 23 deletions

View File

@@ -2,6 +2,8 @@
from smtplib import SMTPServerDisconnected from smtplib import SMTPServerDisconnected
from temporalio import workflow, activity from temporalio import workflow, activity
from orchestrator import metrics
with workflow.unsafe.imports_passed_through(): with workflow.unsafe.imports_passed_through():
import traceback import traceback
import smtplib import smtplib
@@ -147,7 +149,6 @@ class Email(BaseActivity):
for group_name, group_config in receiver_groups.items(): for group_name, group_config in receiver_groups.items():
try: try:
receivers = ", ".join(group_config['members']) receivers = ", ".join(group_config['members'])
self.info(f"Sending email to {group_name}: {receivers}", self.info(f"Sending email to {group_name}: {receivers}",
@@ -177,6 +178,12 @@ class Email(BaseActivity):
group_config['status'] = 'failed' group_config['status'] = 'failed'
else: else:
group_config['status'] = 'sent' group_config['status'] = 'sent'
metrics.EMAIL_SENT_COUNT.labels(
pod_id=self.pod_id,
model_name=metadata['model_name'],
pipeline_name=metadata['workflow_name'],
email_group=group_name
).inc()
self.info(f"Email sent to {group_name}: {receivers}", self.info(f"Email sent to {group_name}: {receivers}",
metadata=metadata) metadata=metadata)

View File

@@ -574,4 +574,11 @@ class Formatters(BaseActivity):
notification) notification)
already_added_keys.append(key) already_added_keys.append(key)
# Remove groups with no notifications
receiver_groups = {
group_name: group
for group_name, group in receiver_groups.items()
if group['notifications']
}
return receiver_groups return receiver_groups

View File

@@ -318,6 +318,13 @@ class SlotManager(Redis):
notification) notification)
already_added_keys.append(key) already_added_keys.append(key)
# Remove groups with no notifications
receiver_groups = {
group_name: group
for group_name, group in receiver_groups.items()
if group['notifications']
}
return receiver_groups return receiver_groups
@activity.defn(name="store_notification_cache") @activity.defn(name="store_notification_cache")

16
orchestrator/metrics.py Normal file
View File

@@ -0,0 +1,16 @@
from prometheus_client import Gauge, Counter
APP_UP = Gauge(
"app_up",
"Indicates if the application is running (1) or shutting down (0)",
["pod_id"],
)
CORE_LABELS = ["pod_id", "model_name", "pipeline_name"]
EMAIL_SENT_COUNT = Counter(
"email_sent_count",
"Number of emails sent",
[*CORE_LABELS, "email_group"],
)

View File

@@ -22,6 +22,10 @@ with workflow.unsafe.imports_passed_through():
) )
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
from sientia_do.temporal.utils.logger import get_logger from sientia_do.temporal.utils.logger import get_logger
from prometheus_client import start_http_server
from orchestrator import metrics
POD_ID = os.getenv("POD_ID")
async def main(): async def main():
@@ -29,7 +33,10 @@ async def main():
namespace = os.getenv('TEMPORAL_NAMESPACE', 'default') namespace = os.getenv('TEMPORAL_NAMESPACE', 'default')
logger = get_logger(__name__) logger = get_logger(__name__)
logger.info('Starting Worker...') logger.info(f'Starting Worker with POD_ID: {POD_ID}')
logger.info("Starting prometheus client...")
start_prometheus_server()
logger.info('Starting Notification Handler...') logger.info('Starting Notification Handler...')
@@ -164,7 +171,20 @@ async def main():
if activities: if activities:
activities.shutdown() activities.shutdown()
# Exit with a non-zero status code to indicate failure to Kubernetes # Exit with a non-zero status code to indicate failure to Kubernetes
metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN
sys.exit(1) sys.exit(1)
def start_prometheus_server():
try:
port = int(os.getenv("HTTP_METRICS_PORT", 9090))
start_http_server(port)
print(f"Prometheus server started on port {port}.")
metrics.APP_UP.labels(pod_id=POD_ID).set(1)
except Exception as e:
print(f"Failed to start Prometheus server: {e}")
os._exit(1)
if __name__ == '__main__': if __name__ == '__main__':
asyncio.run(main()) asyncio.run(main())

View File

@@ -82,6 +82,9 @@ class Alerts:
} }
) )
if not log_report:
return
# Store the notification_id sendings to avoid sending them again # Store the notification_id sendings to avoid sending them again
await workflow.execute_activity_method( await workflow.execute_activity_method(
Activities.store_notification_cache, Activities.store_notification_cache,

View File

@@ -65,6 +65,9 @@ class Reports:
retry_policy=retry_policy retry_policy=retry_policy
) )
if not receiver_groups:
return
# Call subworkflow "process_notifications" passing the notification package # Call subworkflow "process_notifications" passing the notification package
await workflow.execute_child_workflow( await workflow.execute_child_workflow(
'process_notifications', 'process_notifications',

View File

@@ -5,4 +5,5 @@ redis
couchbase couchbase
pymongo pymongo
jinja2 jinja2
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.5 git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.8
prometheus-client

View File

@@ -816,10 +816,5 @@ async def test_filter_notification_reports(formatters):
'notification_id': 'test_notification_id_3' 'notification_id': 'test_notification_id_3'
} }
] ]
},
'test_group_2': {
'group_name': 'test_group_2',
'contents': ['core_alerts'],
'notifications': []
} }
} }

View File

@@ -1,4 +1,4 @@
from unittest.mock import AsyncMock, patch, ANY, call from unittest.mock import AsyncMock, MagicMock, patch, ANY, call
from pytest import fixture, mark from pytest import fixture, mark
from orchestrator.workflows.alerts import Alerts from orchestrator.workflows.alerts import Alerts
from orchestrator.activities.activities import Activities from orchestrator.activities.activities import Activities
@@ -115,3 +115,22 @@ async def test_run_no_data(workflow_mock, alerts):
]) ])
workflow_mock.execute_local_activity_method.assert_not_called() workflow_mock.execute_local_activity_method.assert_not_called()
@mark.asyncio
@patch("orchestrator.workflows.alerts.workflow", new_callable=AsyncMock)
async def test_run_no_log_report(workflow_mock, alerts):
workflow_mock.execute_child_workflow.side_effect = [
MagicMock(),
[]
]
input_data = {
'schedule_name': 'test-schedule-name',
'notification_ttl': 300,
'sent_ttl': 600
}
await alerts.run(input_data)
workflow_mock.execute_activity_method.assert_not_called()

View File

@@ -97,3 +97,19 @@ async def test_run_no_data(workflow_mock, reports):
]) ])
workflow_mock.execute_local_activity_method.assert_not_called() workflow_mock.execute_local_activity_method.assert_not_called()
@mark.asyncio
@patch("orchestrator.workflows.reports.workflow", new_callable=AsyncMock)
async def test_run_no_groups(workflow_mock, reports):
workflow_mock.execute_local_activity_method.return_value = []
input_data = {
'schedule_name': 'test-schedule-name',
'notification_ttl': 300,
'sent_ttl': 600
}
await reports.run(input_data)
workflow_mock.execute_child_workflow.assert_called_once()

View File

@@ -11,7 +11,7 @@ image:
# This sets the pull policy for images. # This sets the pull policy for images.
pullPolicy: Always pullPolicy: Always
# Overrides the image tag whose default is the chart appVersion. # Overrides the image tag whose default is the chart appVersion.
tag: "0.2.7" tag: "0.3.2"
# This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/ # This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/
imagePullSecrets: imagePullSecrets:
@@ -112,19 +112,31 @@ tolerations: []
affinity: {} affinity: {}
services: services:
api: metrics:
enabled: false enabled: true
type: ClusterIP type: ClusterIP
port: 4841 port: 9090
targetPort: 4841 targetPort: 9090
name: api name: metrics
opc: # Configuração do ServiceMonitor para o Prometheus Operator
enabled: false # ref: https://github.com/prometheus-operator/prometheus-operator
type: ClusterIP serviceMonitor:
port: 4840 # Se true, um recurso ServiceMonitor será criado.
targetPort: 4840 enabled: true
name: server # O intervalo no qual as métricas devem ser coletadas (ex: 30s, 1m).
interval: 30s
# O path do endpoint de métricas na sua aplicação.
path: /metrics
# Labels adicionais para o recurso ServiceMonitor.
# Essencial para que o Prometheus Operator o descubra. Se você usa o helm chart kube-prometheus-stack,
# ele procura por ServiceMonitors com o label "release: kube-prometheus-stack".
additionalLabels:
release: kube-prometheus-stack
# Configurações de relabeling adicionais, se necessário.
# ref: https://prometheus.io/docs/prometheus/latest/configuration/configuration/#relabel_config
relabelings: []
port: metrics
env: env:
@@ -132,7 +144,7 @@ env:
- name: GITHUB_REPO_URL - name: GITHUB_REPO_URL
value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git" value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git"
- name: GITHUB_BRANCH - name: GITHUB_BRANCH
value: "SIENTIAPDE-1172-criar-pipeline-de-alertas-orquestrador" value: "SIENTIAPDE-1174-mapear-e-implementar-metricas-a-serem-criadas"
- name: PYTHON_APP - name: PYTHON_APP
value: "orchestrator.worker.worker" value: "orchestrator.worker.worker"
@@ -203,6 +215,8 @@ env:
- name: LOG_LEVEL - name: LOG_LEVEL
value: "DEBUG" value: "DEBUG"
- name: HTTP_METRICS_PORT
value: "9090"
- name: PROJECT_NAME - name: PROJECT_NAME
value: "sientia-orchestrator" value: "sientia-orchestrator"
@@ -223,7 +237,7 @@ ssh:
# kubectl create secret docker-registry docker-hub-secret --namespace sientia --docker-server=http://aignosi.azurecr.io --docker-username=aignosi --docker-password=5I5zpQ6sRaHqX1hD3dr+2mo647yO3FRc359/wu6gsP+ACRDRz5mp # kubectl create secret docker-registry docker-hub-secret --namespace sientia --docker-server=http://aignosi.azurecr.io --docker-username=aignosi --docker-password=5I5zpQ6sRaHqX1hD3dr+2mo647yO3FRc359/wu6gsP+ACRDRz5mp
# helm upgrade --install sientia-orchestrator-worker sientia/sientia-module -n sientia --create-namespace -f ./values.yaml --version 0.4.0-uat # helm upgrade --install sientia-orchestrator-worker sientia/sientia-module -n sientia --create-namespace -f ./values.yaml --version 0.4.0
# kubectl create secret generic git-ssh-key-sientia-orchestrator-worker \ # kubectl create secret generic git-ssh-key-sientia-orchestrator-worker \
# --namespace sientia \ # --namespace sientia \