defmodule DaProductApp.Telemetry.ReversalMetrics do @moduledoc """ Production metrics collection for reversal system observability. This module provides structured metrics collection for Prometheus/Grafana dashboards and production monitoring. It acts as a telemetry handler that processes reversal events and maintains aggregated metrics. ## Metrics Collected ### Counters - `reversal_attempts_total` - Total reversal attempts by reason - `reversal_successes_total` - Total successful reversals - `reversal_failures_total` - Total failed reversals by reason - `stuck_transactions_total` - Total stuck transactions detected - `network_transmission_failures_total` - Network-related failures - `stan_batch_operations_total` - STAN/batch service operations ### Gauges - `reversal_queue_size` - Current number of pending reversals - `manual_review_queue_size` - Transactions requiring manual review - `available_stan_numbers` - Remaining STAN numbers per terminal - `available_batch_numbers` - Remaining batch numbers per terminal - `cleanup_worker_health` - Cleanup worker operational status ### Histograms - `reversal_duration_seconds` - Time to complete reversal operations - `packet_generation_duration_seconds` - MTI 0400 packet generation time - `network_transmission_duration_seconds` - Upstream transmission latency - `stan_batch_retrieval_duration_seconds` - Database operation latency ### Summary - `system_health_score` - Overall reversal system health (0-1) ## Usage Automatically attached to telemetry events during application startup. No manual intervention required - metrics are collected in real-time. Access via Prometheus endpoint: `/metrics` """ use GenServer require Logger alias DaProductApp.Acquirer.{ReversalConfig, ReversalTelemetry} # Metric names following Prometheus conventions @counter_metrics [ :reversal_attempts_total, :reversal_successes_total, :reversal_failures_total, :stuck_transactions_total, :network_transmission_failures_total, :stan_batch_operations_total, :cleanup_cycles_total, :db_operation_failures_total ] @gauge_metrics [ :reversal_queue_size, :manual_review_queue_size, :available_stan_numbers, :available_batch_numbers, :cleanup_worker_health, :system_health_score ] @histogram_metrics [ :reversal_duration_seconds, :packet_generation_duration_seconds, :network_transmission_duration_seconds, :stan_batch_retrieval_duration_seconds, :cleanup_cycle_duration_seconds ] @doc """ Start the metrics collector as part of application supervision tree. """ def start_link(opts \\ []) do GenServer.start_link(__MODULE__, opts, name: __MODULE__) end @doc """ Get current metrics snapshot for monitoring dashboards. Returns map with all current metric values. """ def get_metrics_snapshot do GenServer.call(__MODULE__, :get_metrics_snapshot) end @doc """ Get system health score (0.0 to 1.0). Based on recent success rates, queue sizes, and error rates. """ def get_system_health_score do GenServer.call(__MODULE__, :get_system_health_score) end @doc """ Force metrics recalculation (for testing or manual refresh). """ def refresh_metrics do GenServer.cast(__MODULE__, :refresh_metrics) end # GenServer Implementation @impl true def init(_opts) do # Attach to all reversal telemetry events :telemetry.attach_many( "reversal-metrics-collector", ReversalTelemetry.event_names(), &handle_telemetry_event/4, nil ) # Initialize metrics storage initial_state = %{ counters: initialize_counters(), gauges: initialize_gauges(), histograms: initialize_histograms(), last_updated: DateTime.utc_now(), health_history: [] # Last 10 health calculations } # Schedule periodic health calculations schedule_health_calculation() Logger.info("ReversalMetrics collector started - attached to #{length(ReversalTelemetry.event_names())} telemetry events") {:ok, initial_state} end @impl true def handle_call(:get_metrics_snapshot, _from, state) do snapshot = %{ counters: state.counters, gauges: calculate_current_gauges(state), histograms: get_histogram_summaries(state.histograms), last_updated: state.last_updated, system_health: calculate_system_health(state) } {:reply, snapshot, state} end @impl true def handle_call(:get_system_health_score, _from, state) do health_score = calculate_system_health(state) {:reply, health_score, state} end @impl true def handle_cast(:refresh_metrics, state) do updated_state = %{state | gauges: calculate_current_gauges(state), last_updated: DateTime.utc_now() } {:noreply, updated_state} end @impl true def handle_info(:calculate_health, state) do health_score = calculate_system_health(state) # Keep last 10 health scores for trending updated_history = [health_score | state.health_history] |> Enum.take(10) updated_state = %{state | health_history: updated_history, last_updated: DateTime.utc_now() } # Schedule next calculation schedule_health_calculation() {:noreply, updated_state} end # Telemetry Event Handlers @doc false def handle_telemetry_event([:da_product_app, :reversal, :attempt], measurements, metadata, _config) do GenServer.cast(__MODULE__, {:increment_counter, :reversal_attempts_total, [ reason: metadata.reason, attempt_number: measurements.attempt_number ]}) end def handle_telemetry_event([:da_product_app, :reversal, :success], measurements, metadata, _config) do GenServer.cast(__MODULE__, {:increment_counter, :reversal_successes_total, []}) GenServer.cast(__MODULE__, {:record_histogram, :reversal_duration_seconds, measurements.duration_ms / 1000}) end def handle_telemetry_event([:da_product_app, :reversal, :failure], measurements, metadata, _config) do GenServer.cast(__MODULE__, {:increment_counter, :reversal_failures_total, [ reason: extract_failure_category(metadata.reason), attempt_number: measurements.attempt_number ]}) end def handle_telemetry_event([:da_product_app, :reversal, :stuck_transaction], measurements, metadata, _config) do GenServer.cast(__MODULE__, {:increment_counter, :stuck_transactions_total, [ status: metadata.status, age_bucket: categorize_age(measurements.age_seconds) ]}) end def handle_telemetry_event([:da_product_app, :reversal, :cleanup_cycle_complete], measurements, metadata, _config) do GenServer.cast(__MODULE__, {:increment_counter, :cleanup_cycles_total, []}) GenServer.cast(__MODULE__, {:record_histogram, :cleanup_cycle_duration_seconds, measurements.duration_ms / 1000}) end def handle_telemetry_event([:da_product_app, :reversal, :connection_loss], measurements, metadata, _config) do GenServer.cast(__MODULE__, {:increment_counter, :network_transmission_failures_total, [ reason: "connection_loss", affected_count: measurements.affected_count ]}) end def handle_telemetry_event([:da_product_app, :reversal, :db_operation_failure], measurements, metadata, _config) do GenServer.cast(__MODULE__, {:increment_counter, :db_operation_failures_total, [ operation: metadata.operation ]}) end def handle_telemetry_event(_event, _measurements, _metadata, _config) do # Unknown event - ignore :ok end @impl true def handle_cast({:increment_counter, metric, labels}, state) do updated_counters = increment_counter(state.counters, metric, labels) updated_state = %{state | counters: updated_counters, last_updated: DateTime.utc_now()} {:noreply, updated_state} end @impl true def handle_cast({:record_histogram, metric, value}, state) do updated_histograms = record_histogram_value(state.histograms, metric, value) updated_state = %{state | histograms: updated_histograms, last_updated: DateTime.utc_now()} {:noreply, updated_state} end # Private Functions defp initialize_counters do Enum.reduce(@counter_metrics, %{}, fn metric, acc -> Map.put(acc, metric, %{}) end) end defp initialize_gauges do Enum.reduce(@gauge_metrics, %{}, fn metric, acc -> Map.put(acc, metric, 0) end) end defp initialize_histograms do Enum.reduce(@histogram_metrics, %{}, fn metric, acc -> Map.put(acc, metric, []) end) end defp increment_counter(counters, metric, labels) do label_key = labels_to_key(labels) current_value = get_in(counters, [metric, label_key]) || 0 put_in(counters, [metric, label_key], current_value + 1) end defp record_histogram_value(histograms, metric, value) do current_values = Map.get(histograms, metric, []) # Keep last 1000 values for percentile calculations updated_values = [value | current_values] |> Enum.take(1000) Map.put(histograms, metric, updated_values) end defp labels_to_key(labels) do labels |> Enum.sort() |> Enum.map(fn {k, v} -> "#{k}=#{v}" end) |> Enum.join(",") end defp calculate_current_gauges(state) do # These would query the database/system state for current values %{ reversal_queue_size: count_pending_reversals(), manual_review_queue_size: count_manual_review_queue(), available_stan_numbers: estimate_available_stan_numbers(), available_batch_numbers: estimate_available_batch_numbers(), cleanup_worker_health: check_cleanup_worker_health(), system_health_score: calculate_system_health(state) } end defp get_histogram_summaries(histograms) do Enum.reduce(histograms, %{}, fn {metric, values}, acc -> if Enum.empty?(values) do Map.put(acc, metric, %{count: 0, sum: 0, avg: 0, p50: 0, p95: 0, p99: 0}) else sorted_values = Enum.sort(values) count = length(values) sum = Enum.sum(values) avg = sum / count Map.put(acc, metric, %{ count: count, sum: sum, avg: avg, p50: percentile(sorted_values, 50), p95: percentile(sorted_values, 95), p99: percentile(sorted_values, 99) }) end end) end defp percentile(sorted_values, percentile) do count = length(sorted_values) index = round(count * percentile / 100) - 1 index = max(0, min(index, count - 1)) Enum.at(sorted_values, index, 0) end defp calculate_system_health(state) do # Calculate health score based on multiple factors success_rate = calculate_success_rate(state.counters) queue_health = calculate_queue_health() error_rate = calculate_error_rate(state.counters) performance_health = calculate_performance_health(state.histograms) # Weighted average of health factors health_score = (success_rate * 0.4) + (queue_health * 0.2) + ((1 - error_rate) * 0.2) + (performance_health * 0.2) Float.round(health_score, 3) end defp calculate_success_rate(counters) do successes = get_counter_total(counters, :reversal_successes_total) failures = get_counter_total(counters, :reversal_failures_total) total = successes + failures if total > 0 do successes / total else 1.0 # No attempts = perfect success rate end end defp calculate_queue_health do # Health decreases as queue sizes increase reversal_queue = count_pending_reversals() manual_queue = count_manual_review_queue() # Target: <10 pending reversals, <5 manual reviews reversal_health = max(0, (10 - reversal_queue) / 10) manual_health = max(0, (5 - manual_queue) / 5) (reversal_health + manual_health) / 2 end defp calculate_error_rate(counters) do errors = get_counter_total(counters, :reversal_failures_total) + get_counter_total(counters, :network_transmission_failures_total) + get_counter_total(counters, :db_operation_failures_total) attempts = get_counter_total(counters, :reversal_attempts_total) if attempts > 0 do errors / attempts else 0.0 end end defp calculate_performance_health(histograms) do # Health based on response times - lower is better reversal_times = Map.get(histograms, :reversal_duration_seconds, []) if Enum.empty?(reversal_times) do 1.0 else avg_time = Enum.sum(reversal_times) / length(reversal_times) # Target: <5 seconds average, perfect at <1 second max(0, (5 - avg_time) / 5) end end defp get_counter_total(counters, metric) do metric_data = Map.get(counters, metric, %{}) Enum.sum(Map.values(metric_data)) end # Database query helpers (mock implementations - replace with real queries) defp count_pending_reversals do # TODO: Query database for PosReversal where status IN ('PENDING', 'SENT') 0 end defp count_manual_review_queue do # TODO: Query database for PosTempTransaction where status = 'PENDING_MANUAL_REVIEW' 0 end defp estimate_available_stan_numbers do # TODO: Query AcquirerTerminalStan for available STAN count 999999 end defp estimate_available_batch_numbers do # TODO: Query AcquirerBatch for available batch count 999 end defp check_cleanup_worker_health do # TODO: Check if ReversalCleanupWorker is running and healthy case GenServer.whereis(DaProductApp.Acquirer.ReversalCleanupWorker) do nil -> 0 _pid -> 1 end end # Utility functions defp extract_failure_category(reason_string) do cond do String.contains?(reason_string, "timeout") -> "timeout" String.contains?(reason_string, "network") -> "network" String.contains?(reason_string, "db") -> "database" String.contains?(reason_string, "validation") -> "validation" true -> "other" end end defp categorize_age(age_seconds) do cond do age_seconds < 60 -> "under_1m" age_seconds < 300 -> "1m_to_5m" age_seconds < 900 -> "5m_to_15m" age_seconds < 3600 -> "15m_to_1h" true -> "over_1h" end end defp schedule_health_calculation do # Calculate health every 30 seconds Process.send_after(self(), :calculate_health, 30_000) end end