from .base import SnowflakeJob from src.utils.monitoring import datadog_client from src.utils.snowflake.constants import QueryType, Role, Warehouse import pandas as pd from typing import Optional import time from dagster import OpExecutionContext class StripeDisputesMonitor(SnowflakeJob): def __init__(self): super().__init__( name="stripe_disputes_monitor", description="Counts the number of disputes in Stripe and reports to Datadog", schedule="0 * * * *", # Runs at minute 0 of every hour query_file="queries/stripe_disputes_monitor.sql", query_type=QueryType.SELECT, warehouse=Warehouse.SMALL, role=Role.ACCOUNTADMIN, group_name="finops", monitored=True, owners=["team:core-pod"], metadata={ "slack": "#tech-anti-bots", }, tags={"team": "core-pod", "category": "finops", "tech-alerts": "true"}, ) def post_execute(self, result: Optional[pd.DataFrame], context: OpExecutionContext) -> None: """Optional hook for child classes to process query results. Args: result: DataFrame for SELECT queries, None for DDL queries """ if result is not None: now = int(time.time()) # Current Unix timestamp metrics = [] for _, row in result.iterrows(): metrics.append( { "metric": "stripe.disputes.count", "points": [(now, float(row["TOTAL_DISPUTES"]))], # Must be float "type": "gauge", "tags": [f"status:{row['STATUS']}"], } ) datadog_client.Metric.send(metrics=metrics) stripe_disputes_monitor = StripeDisputesMonitor()