exercises

Unnamed repository; edit this file 'description' to name the repository.
Log | Files | Refs | README

commit 35ea0ca585c9f712c8ba5f44b30e8adaccbccd1d
parent 6a4c130576362f984076231453ccf75325ed2bc6
Author: ling0x <ling0x@users.noreply.github.com>
Date:   Sun, 30 Aug 2026 15:02:10 +0100

observability

Diffstat:
Mdistributed_observability/docker-compose.yml | 113++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------------
Adistributed_observability/grafana/provisioning/alerting/.gitkeep | 0
Adistributed_observability/grafana/provisioning/dashboards/dashboards.yml | 17+++++++++++++++++
Adistributed_observability/grafana/provisioning/dashboards/otelmart.json | 873+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Adistributed_observability/grafana/provisioning/datasources/datasources.yml | 29+++++++++++++++++++++++++++++
Adistributed_observability/grafana/provisioning/plugins/.gitkeep | 0
Mdistributed_observability/inventory/src/handlers/inventory.rs | 46+++++++++++++++++++++++++++++++++++++++++++++-
Mdistributed_observability/inventory/src/main.rs | 1+
Adistributed_observability/inventory/src/metrics.rs | 55+++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mdistributed_observability/inventory/src/telemetry.rs | 4++--
Mdistributed_observability/orders/src/handlers/orders.rs | 84+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mdistributed_observability/orders/src/main.rs | 1+
Adistributed_observability/orders/src/metrics.rs | 88+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mdistributed_observability/orders/src/telemetry.rs | 4++--
Mdistributed_observability/otel-collector-config.yaml | 6++++++
Mdistributed_observability/otelmart/src/main.rs | 1+
Adistributed_observability/otelmart/src/metrics.rs | 32++++++++++++++++++++++++++++++++
Mdistributed_observability/otelmart/src/proxy/mod.rs | 108++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Mdistributed_observability/otelmart/src/telemetry.rs | 4++--
Mdistributed_observability/products/src/telemetry.rs | 4++--
Adistributed_observability/prometheus.yml | 33+++++++++++++++++++++++++++++++++
Adistributed_observability/scripts/export-dashboards.sh | 85+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Adistributed_observability/scripts/generate_traffic.sh | 626+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Adistributed_observability/tempo.yaml | 50++++++++++++++++++++++++++++++++++++++++++++++++++
24 files changed, 2192 insertions(+), 72 deletions(-)

diff --git a/distributed_observability/docker-compose.yml b/distributed_observability/docker-compose.yml @@ -1,42 +1,66 @@ services: # ============================================================ - # Otel Collector + # Tempo - Distributed Tracing Backend # ============================================================ - otel-collector: - image: otel/opentelemetry-collector-contrib:latest - container_name: otel-collector - command: ["--config=/etc/otel-collector-config.yaml"] + # Receives traces via OTLP and is queried through Grafana (no separate UI). + # Its metrics-generator derives RED metrics and service graphs from spans + # and remote-writes them to Prometheus, powering Traces Drilldown. + tempo: + image: grafana/tempo:latest + container_name: tempo + command: ["-config.file=/etc/tempo.yaml"] volumes: - - ./otel-collector-config.yaml:/etc/otel-collector-config.yaml + - ./tempo.yaml:/etc/tempo.yaml:ro + - tempo_data:/var/tempo ports: - - "4317:4317" - - "4318:4318" - depends_on: - jaeger: - condition: service_healthy + - "3200:3200" # Tempo query API + - "4317:4317" # OTLP gRPC receiver + - "4318:4318" # OTLP HTTP receiver networks: - app-network + restart: unless-stopped # ============================================================ - # Jaeger + # Prometheus - Metrics Backend # ============================================================ - jaeger: - image: jaegertracing/all-in-one:latest - container_name: jaeger - environment: - - COLLECTOR_OTLP_ENABLED=true + # Prometheus 3.0 receives metrics via OTLP HTTP push (no scrape targets needed) + prometheus: + image: prom/prometheus:v3.0.0 + container_name: prometheus + command: + - '--config.file=/etc/prometheus/prometheus.yml' + - '--web.enable-otlp-receiver' + - '--web.enable-remote-write-receiver' # accepts Tempo's span metrics + - '--enable-feature=exemplar-storage' # trace exemplars on metrics + volumes: + - ./prometheus.yml:/etc/prometheus/prometheus.yml:ro + - prometheus_data:/prometheus ports: - # OTLP ports are NOT published: the collector owns 4317/4318 on the host. - # Jaeger still receives OTLP on 4317 over app-network. - - "16686:16686" - healthcheck: - test: ["CMD", "wget", "--no-verbose", "--tries=1", "--spider", "http://localhost:14269/"] - interval: 5s - timeout: 3s - retries: 12 - start_period: 10s + - "9090:9090" networks: - app-network + restart: unless-stopped + + # ============================================================ + # Grafana - Dashboards & Visualization + # ============================================================ + # Queries both Prometheus (metrics) and Tempo (traces). + # Datasources are provisioned from ./grafana/provisioning so they survive + # container recreates instead of living in the container's ephemeral DB. + grafana: + image: grafana/grafana:latest + container_name: grafana + volumes: + - ./grafana/provisioning:/etc/grafana/provisioning:ro + - grafana_data:/var/lib/grafana + ports: + - "3000:3000" + depends_on: + - prometheus + - tempo + networks: + - app-network + restart: unless-stopped # ============================================================ # Database Service @@ -78,13 +102,16 @@ services: - "3001:3001" environment: - APP_DATABASE_URL=postgres://opentel_user:opentel_pass@postgres:5432/opentel_db - - OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-collector:4317 + - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 + - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics - OTEL_SERVICE_NAME=products - RUST_LOG=info,otel::tracing=trace depends_on: postgres: condition: service_healthy - otel-collector: + tempo: + condition: service_started + prometheus: condition: service_started networks: - app-network @@ -103,13 +130,16 @@ services: - "3002:3002" environment: - APP_DATABASE_URL=postgres://opentel_user:opentel_pass@postgres:5432/opentel_db - - OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-collector:4317 + - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 + - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics - OTEL_SERVICE_NAME=inventory - RUST_LOG=info,otel::tracing=trace depends_on: postgres: condition: service_healthy - otel-collector: + tempo: + condition: service_started + prometheus: condition: service_started networks: - app-network @@ -128,13 +158,16 @@ services: - "3003:3003" environment: - APP_DATABASE_URL=postgres://opentel_user:opentel_pass@postgres:5432/opentel_db - - OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-collector:4317 + - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 + - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics - OTEL_SERVICE_NAME=orders - RUST_LOG=info,otel::tracing=trace depends_on: postgres: condition: service_healthy - otel-collector: + tempo: + condition: service_started + prometheus: condition: service_started networks: - app-network @@ -187,7 +220,8 @@ services: - APP_SERVICES__ORDERS_PORT=3003 - APP_SERVICES__INVENTORY_SERVICE_URL=http://inventory:3002 - APP_SERVICES__ORDERS_SERVICE_URL=http://orders:3003 - - OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-collector:4317 + - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 + - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics - OTEL_SERVICE_NAME=otelmart - RUST_LOG=info,otel::tracing=trace extra_hosts: @@ -195,6 +229,10 @@ services: depends_on: postgres: condition: service_healthy + tempo: + condition: service_started + prometheus: + condition: service_started data-ingestion: condition: service_completed_successfully products: @@ -203,8 +241,6 @@ services: condition: service_started orders: condition: service_started - otel-collector: - condition: service_started networks: - app-network restart: unless-stopped @@ -222,3 +258,10 @@ networks: volumes: postgres_data: name: opentel_postgres_data + tempo_data: + name: opentel_tempo_data + prometheus_data: + name: opentel_prometheus_data + grafana_data: + name: opentel_grafana_data + diff --git a/distributed_observability/grafana/provisioning/alerting/.gitkeep b/distributed_observability/grafana/provisioning/alerting/.gitkeep diff --git a/distributed_observability/grafana/provisioning/dashboards/dashboards.yml b/distributed_observability/grafana/provisioning/dashboards/dashboards.yml @@ -0,0 +1,17 @@ +# Loads dashboard JSON from this directory on Grafana startup, so dashboards +# live in git rather than only inside the grafana_data volume. +apiVersion: 1 + +providers: + - name: otelmart + orgId: 1 + folder: '' # General + type: file + disableDeletion: false + updateIntervalSeconds: 30 + # Keep the dashboard editable in the UI. Edits save to Grafana's DB; + # re-export to this file to commit them (see repo notes). + allowUiUpdates: true + options: + path: /etc/grafana/provisioning/dashboards + foldersFromFilesStructure: false diff --git a/distributed_observability/grafana/provisioning/dashboards/otelmart.json b/distributed_observability/grafana/provisioning/dashboards/otelmart.json @@ -0,0 +1,873 @@ +{ + "annotations": { + "list": [ + { + "builtIn": 1, + "datasource": { + "type": "grafana", + "uid": "-- Grafana --" + }, + "enable": true, + "hide": true, + "iconColor": "rgba(0, 211, 255, 1)", + "name": "Annotations & Alerts", + "type": "dashboard" + } + ] + }, + "editable": true, + "fiscalYearStartMonth": 0, + "graphTooltip": 0, + "liveNow": false, + "panels": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 0, + "gradientMode": "none", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "none" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 0 + }, + "id": 2, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "sum(rate(http_server_request_duration_seconds_count[5m])) by (job)", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Request Rate: The Store Traffic", + "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 25, + "gradientMode": "none", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "none" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 12, + "y": 0 + }, + "id": 1, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "sum(rate(inventory_reservation_attempts_total[5m])) by (job)", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Inventory Reservation Rate", + "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 0, + "gradientMode": "none", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "none" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 8 + }, + "id": 3, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "histogram_quantile(0.5, sum(rate(http_server_request_duration_seconds_bucket[5m])) by (le, job))", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Latency: Checkout Speed (Wait Times - 0.5 percentile)", + "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic", + "seriesBy": "last" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 27, + "gradientMode": "opacity", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineStyle": { + "fill": "solid" + }, + "lineWidth": 1, + "pointSize": 4, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "none" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "percentage", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "#EAB839", + "value": 60 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 12, + "y": 8 + }, + "id": 5, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "(sum(rate(orders_checkout_attempts_total[5m])) - sum(rate(orders_checkout_failures_total[5m]))) / sum(rate(orders_checkout_attempts_total[5m])) * 100", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Business Metrics: Checkout Success Rate", + "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 0, + "gradientMode": "none", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "none" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 16 + }, + "id": 7, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "histogram_quantile(0.9, sum(rate(http_server_request_duration_seconds_bucket[5m])) by (le, job))", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Latency: Checkout Speed (Wait Times - 0.9 percentile)", + "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 0, + "gradientMode": "none", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "none" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 12, + "y": 16 + }, + "id": 4, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "db_pool_connections_idle{job=~\".+\"}", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Idle Connections (Pool Headroom)", + "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 0, + "gradientMode": "none", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "none" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 24 + }, + "id": 8, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "histogram_quantile(0.99, sum(rate(http_server_request_duration_seconds_bucket[5m])) by (le, job))", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Latency: Checkout Speed (Wait Times - 0.99 percentile)", + "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "palette-classic" + }, + "custom": { + "axisBorderShow": false, + "axisCenteredZero": false, + "axisColorMode": "text", + "axisLabel": "", + "axisPlacement": "auto", + "barAlignment": 0, + "barWidthFactor": 0.6, + "drawStyle": "line", + "fillOpacity": 25, + "gradientMode": "none", + "hideFrom": { + "legend": false, + "tooltip": false, + "viz": false + }, + "insertNulls": false, + "lineInterpolation": "linear", + "lineWidth": 1, + "pointSize": 5, + "scaleDistribution": { + "type": "linear" + }, + "showPoints": "auto", + "showValues": false, + "spanNulls": false, + "stacking": { + "group": "A", + "mode": "normal" + }, + "thresholdsStyle": { + "mode": "off" + } + }, + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": 0 + }, + { + "color": "red", + "value": 80 + } + ] + } + } + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 32 + }, + "id": 6, + "options": { + "annotations": { + "clustering": -1, + "multiLane": false + }, + "legend": { + "calcs": [], + "displayMode": "list", + "enableFacetedFilter": false, + "overflow": "ellipsis", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "hideZeros": false, + "mode": "single", + "sort": "none" + } + }, + "pluginVersion": "13.1.3", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "prometheus" + }, + "editorMode": "code", + "expr": "(sum(rate(http_server_request_duration_seconds_count{http_response_status_code=~\"4..|5..\"}[5m])) by (job) / sum(rate(http_server_request_duration_seconds_count[5m])) by (job) * 100) or (sum(rate(http_server_request_duration_seconds_count[5m])) by (job) * 0)", + "legendFormat": "__auto", + "range": true, + "refId": "A" + } + ], + "title": "Error Rate: Broken Registers", + "type": "timeseries" + } + ], + "preload": false, + "refresh": "", + "schemaVersion": 42, + "time": { + "from": "now-6h", + "to": "now" + }, + "timepicker": { + "refresh_intervals": [ + "5s", + "10s", + "30s", + "1m", + "5m", + "15m", + "30m", + "1h", + "2h", + "1d" + ] + }, + "timezone": "browser", + "title": "Otelmart Dashboard", + "uid": "ad79bw7", + "version": 15 +} diff --git a/distributed_observability/grafana/provisioning/datasources/datasources.yml b/distributed_observability/grafana/provisioning/datasources/datasources.yml @@ -0,0 +1,29 @@ +# Grafana datasources, provisioned as code so they survive container recreates. +apiVersion: 1 + +datasources: + - name: prometheus + type: prometheus + uid: prometheus + access: proxy + url: http://prometheus:9090 + isDefault: true + jsonData: + timeInterval: 60s # matches the OTel PeriodicReader export interval + exemplarTraceIdDestinations: + - name: trace_id + datasourceUid: tempo + + - name: tempo + type: tempo + uid: tempo + access: proxy + url: http://tempo:3200 + jsonData: + # Link a trace's spans to the matching Prometheus metrics. + tracesToMetrics: + datasourceUid: prometheus + serviceMap: + datasourceUid: prometheus + nodeGraph: + enabled: true diff --git a/distributed_observability/grafana/provisioning/plugins/.gitkeep b/distributed_observability/grafana/provisioning/plugins/.gitkeep diff --git a/distributed_observability/inventory/src/handlers/inventory.rs b/distributed_observability/inventory/src/handlers/inventory.rs @@ -8,10 +8,13 @@ use axum::{ http::StatusCode, response::{IntoResponse, Json}, }; +use opentelemetry::KeyValue; use sqlx::PgPool; +use std::time::Instant; use uuid::Uuid; use crate::db; +use crate::metrics::metrics; use crate::models::{ ConfirmSaleRequest, InventoryQueryParams, InventoryResponse, ReleaseStockRequest, ReserveStockRequest, StockOperationResponse, UpdateStockRequest, @@ -122,10 +125,24 @@ pub async fn reserve_stock( State(pool): State<PgPool>, Json(request): Json<ReserveStockRequest>, ) -> impl IntoResponse { + // Start timing and record a reservation attempt + let start = Instant::now(); + metrics().reservation_attempts.add(1, &[]); + // Delegate to repository layer for stock reservation match db::reserve_stock(&pool, request.product_uuid, request.quantity).await { Ok(success) => { + let duration = start.elapsed().as_secs_f64(); + if success { + // Record successful reservation metrics + metrics() + .reservation_duration + .record(duration, &[KeyValue::new("outcome", "success")]); + metrics() + .reserved_quantity + .add(request.quantity as u64, &[]); + ( StatusCode::OK, Json(StockOperationResponse { @@ -137,6 +154,18 @@ pub async fn reserve_stock( ) .into_response() } else { + // Insufficient stock — record failure metrics + metrics() + .reservation_failures + .add(1, &[KeyValue::new("failure.reason", "insufficient_stock")]); + metrics().reservation_duration.record( + duration, + &[ + KeyValue::new("outcome", "failure"), + KeyValue::new("failure.reason", "insufficient_stock"), + ], + ); + ( StatusCode::CONFLICT, Json(StockOperationResponse { @@ -149,7 +178,22 @@ pub async fn reserve_stock( .into_response() } } - Err(e) => internal_error("Failed to reserve stock", e.to_string()), + Err(e) => { + // Database error — record failure metrics + let duration = start.elapsed().as_secs_f64(); + metrics() + .reservation_failures + .add(1, &[KeyValue::new("failure.reason", "database_error")]); + metrics().reservation_duration.record( + duration, + &[ + KeyValue::new("outcome", "failure"), + KeyValue::new("failure.reason", "database_error"), + ], + ); + + internal_error("Failed to reserve stock", e.to_string()) + } } } diff --git a/distributed_observability/inventory/src/main.rs b/distributed_observability/inventory/src/main.rs @@ -36,6 +36,7 @@ mod config; mod db; mod handlers; +mod metrics; mod models; mod telemetry; mod utils; diff --git a/distributed_observability/inventory/src/metrics.rs b/distributed_observability/inventory/src/metrics.rs @@ -0,0 +1,55 @@ +//! Business-specific metrics for the Inventory service. +//! +//! Defines counters and histograms that track stock reservation KPIs: +//! attempts, failures, duration, and total reserved quantity. +//! Uses the `OnceLock` singleton pattern so metrics are initialized +//! once and reusable from any handler. + +use opentelemetry::{ + global, + metrics::{Counter, Histogram}, +}; +use std::sync::OnceLock; + +/// Collection of inventory-related metric instruments. +pub struct InventoryMetrics { + /// Total number of stock reservation attempts + pub reservation_attempts: Counter<u64>, + /// Total number of failed reservations, broken down by failure.reason + pub reservation_failures: Counter<u64>, + /// Wall-clock time to complete a stock reservation (seconds) + pub reservation_duration: Histogram<f64>, + /// Cumulative units reserved across all successful reservations + pub reserved_quantity: Counter<u64>, +} + +/// Singleton storage for the metrics instruments. +static METRICS: OnceLock<InventoryMetrics> = OnceLock::new(); + +/// Returns a reference to the lazily-initialized inventory metrics. +pub fn metrics() -> &'static InventoryMetrics { + METRICS.get_or_init(|| { + let meter = global::meter("inventory-service"); + + InventoryMetrics { + reservation_attempts: meter + .u64_counter("inventory.reservation.attempts") + .with_description("Total stock reservation attempts") + .build(), + reservation_failures: meter + .u64_counter("inventory.reservation.failures") + .with_description("Failed stock reservations") + .build(), + reservation_duration: meter + .f64_histogram("inventory.reservation.duration") + .with_description("Time to complete stock reservation") + .with_unit("s") + .build(), + reserved_quantity: meter + .u64_counter("inventory.reserved.quantity") + .with_description("Total units reserved") + .with_unit("units") + .build(), + } + }) +} diff --git a/distributed_observability/inventory/src/telemetry.rs b/distributed_observability/inventory/src/telemetry.rs @@ -67,11 +67,11 @@ pub fn init_telemetry(service_name: &str) -> TelemetryGuard { // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:4317".into()); + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + .unwrap_or_else(|_| "http://localhost:4317".into()); // Initialize both providers let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); diff --git a/distributed_observability/orders/src/handlers/orders.rs b/distributed_observability/orders/src/handlers/orders.rs @@ -1,18 +1,24 @@ //! Order management API handlers +use std::time::Instant; + use axum::{ extract::{Path, Query, State}, http::StatusCode, response::{IntoResponse, Json}, }; use bigdecimal::BigDecimal; +use opentelemetry::KeyValue; use sqlx::PgPool; use tracing::{info, warn}; use uuid::Uuid; -use crate::db::{self, with_transaction, OrderTotals}; use crate::models::{CreateOrderRequest, OrderQueryParams, OrdersResponse}; use crate::AppState; +use crate::{ + db::{self, with_transaction, OrderTotals}, + metrics::metrics, +}; /// Product detail response from products service #[derive(Debug, serde::Deserialize)] @@ -389,18 +395,28 @@ pub async fn create_order( State(state): State<AppState>, Json(request): Json<CreateOrderRequest>, ) -> impl IntoResponse { + // Start timing and record a checkout attempt + let start = Instant::now(); + metrics().checkout_attempts.add(1, &[]); + + // Business KPI: Conversion funnel started + metrics().funnel_started.add(1, &[]); + // Validate all products exist, are available, and have correct prices for item in &request.items { let product = match validate_product(&state, item.product_uuid).await { Ok(product) => product, Err(response) => { - // Product validation failed - return error immediately + // Product validation failed — record checkout failure + record_checkout_failure("product_validation", &start); + // and then return error immediately return response; } }; // Validate price matches actual product price (prevent price manipulation) if item.unit_price != product.final_price { + record_checkout_failure("price_mismatch", &start); return ( StatusCode::BAD_REQUEST, Json(serde_json::json!({ @@ -416,6 +432,9 @@ pub async fn create_order( } } + // Business KPI: Payment info validated (all products and prices confirmed) + metrics().funnel_payment_info.add(1, &[]); + // Reserve stock for all items // Track reserved items so we can release them if something fails let mut reserved_items: Vec<(Uuid, i32)> = Vec::new(); @@ -508,6 +527,25 @@ pub async fn create_order( match result { Ok(created_order) => { // Transaction committed successfully! + + // Record successful checkout metrics + let duration = start.elapsed().as_secs_f64(); + metrics() + .checkout_duration + .record(duration, &[KeyValue::new("outcome", "success")]); + // Business KPI: Conversion funnel completed + metrics().funnel_completed.add(1, &[]); + // Business KPI: Time from checkout start to completion + metrics().time_to_checkout.record(duration, &[]); + // Record the order total as a dollar-value histogram + metrics() + .order_total_amount + .record(bigdecimal_to_f64(&total), &[]); + // Record how many line items were in this order + metrics() + .order_items_count + .record(request.items.len() as u64, &[]); + // Now confirm the sale with inventory service to convert reserved stock to sold confirm_all_sales(&state, created_order.uuid, &reserved_items).await; @@ -525,6 +563,8 @@ pub async fn create_order( Err(e) => { // Transaction failed (rolled back automatically) + // Record failure + record_checkout_failure("order_creation", &start); // Release reserved stock release_all_reserved_stock(&state, &reserved_items).await; @@ -706,3 +746,43 @@ pub async fn get_order_by_id( } } } + +// --------------------------------------------------------------------------- +// Metric helper functions +// --------------------------------------------------------------------------- + +/// Records a checkout failure with the given reason. +/// +/// Increments the failure counter and records the duration histogram +/// with `outcome=failure` and a `failure.reason` attribute. +fn record_checkout_failure(reason: &str, start: &Instant) { + let duration = start.elapsed().as_secs_f64(); + metrics() + .checkout_failures + .add(1, &[KeyValue::new("failure.reason", reason.to_string())]); + metrics().checkout_duration.record( + duration, + &[ + KeyValue::new("outcome", "failure"), + KeyValue::new("failure.reason", reason.to_string()), + ], + ); +} + +/// Produces a short hex hash of the input string (e.g. an email address) +/// for use as a PII-safe span attribute. +#[allow(dead_code)] +fn hash_short(input: &str) -> String { + use std::collections::hash_map::DefaultHasher; + use std::hash::{Hash, Hasher}; + let mut hasher = DefaultHasher::new(); + input.hash(&mut hasher); + format!("{:x}", hasher.finish()) +} + +/// Converts a `BigDecimal` to `f64` for metric recording. +/// Falls back to 0.0 if the conversion is lossy. +fn bigdecimal_to_f64(value: &BigDecimal) -> f64 { + use std::str::FromStr; + f64::from_str(&value.to_string()).unwrap_or(0.0) +} diff --git a/distributed_observability/orders/src/main.rs b/distributed_observability/orders/src/main.rs @@ -33,6 +33,7 @@ mod config; mod db; mod handlers; +mod metrics; mod models; mod telemetry; diff --git a/distributed_observability/orders/src/metrics.rs b/distributed_observability/orders/src/metrics.rs @@ -0,0 +1,88 @@ +use std::sync::OnceLock; + +use opentelemetry::{ + global, + metrics::{Counter, Histogram}, +}; + +/// Collection of order-related metric instruments. +pub struct OrdersMetrics { + /// Total number of checkout attempts (success + failure) + pub checkout_attempts: Counter<u64>, + /// Total number of failed checkouts, broken down by failure.reason + pub checkout_failures: Counter<u64>, + /// Wall-clock time to complete a checkout (seconds) + pub checkout_duration: Histogram<f64>, + /// Dollar amount of each completed order + pub order_total_amount: Histogram<f64>, + /// Number of line items per order + pub order_items_count: Histogram<u64>, + + // Business KPI metrics: Conversion funnel + /// Checkout funnel started (order creation initiated) + pub funnel_started: Counter<u64>, + /// Checkout funnel reached payment validation + pub funnel_payment_info: Counter<u64>, + /// Checkout funnel completed successfully + pub funnel_completed: Counter<u64>, + + // Business KPI metrics: Time tracking + /// Time from checkout start to completion + pub time_to_checkout: Histogram<f64>, +} + +/// Singleton storage for the metrics instruments. +static METRICS: OnceLock<OrdersMetrics> = OnceLock::new(); + +/// Returns a reference to the lazily-initialized orders metrics. +pub fn metrics() -> &'static OrdersMetrics { + METRICS.get_or_init(|| { + let meter = global::meter("orders-service"); + + OrdersMetrics { + checkout_attempts: meter + .u64_counter("orders.checkout.attempts") + .with_description("Total checkout attempts") + .build(), + checkout_failures: meter + .u64_counter("orders.checkout.failures") + .with_description("Total failed checkouts") + .build(), + checkout_duration: meter + .f64_histogram("orders.checkout.duration") + .with_description("Time to complete checkout") + .with_unit("s") + .build(), + order_total_amount: meter + .f64_histogram("orders.total.amount") + .with_description("Order total in dollars") + .with_unit("USD") + .build(), + order_items_count: meter + .u64_histogram("orders.items.count") + .with_description("Number of items per order") + .build(), + + // Business KPI metrics: Conversion funnel + funnel_started: meter + .u64_counter("checkout.funnel.started") + .with_description("Checkout attempts initiated") + .build(), + funnel_payment_info: meter + .u64_counter("checkout.funnel.payment_info") + .with_description("Checkout reached payment validation") + .build(), + funnel_completed: meter + .u64_counter("checkout.funnel.completed") + .with_description("Checkout successfully completed") + .build(), + + // Business KPI metrics: Time tracking + time_to_checkout: meter + .f64_histogram("checkout.time_to_completion") + .with_description("Time from checkout start to completion") + .with_unit("s") + .build(), + } + }) +} diff --git a/distributed_observability/orders/src/telemetry.rs b/distributed_observability/orders/src/telemetry.rs @@ -67,11 +67,11 @@ pub fn init_telemetry(service_name: &str) -> TelemetryGuard { // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:4317".into()); + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + .unwrap_or_else(|_| "http://localhost:4317".into()); // Initialize both providers let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); diff --git a/distributed_observability/otel-collector-config.yaml b/distributed_observability/otel-collector-config.yaml @@ -16,6 +16,8 @@ exporters: endpoint: jaeger:4317 # Forward to Jaeger tls: insecure: true # No TLS in dev, use TLS in prod + otlp_http/prometheus: + endpoint: http://prometheus:9090/api/v1/otlp # Forward metrics to Prometheus OTLP receiver debug: verbosity: detailed @@ -25,3 +27,7 @@ service: receivers: [otlp] processors: [batch] exporters: [otlp_grpc/jaeger, debug] + metrics: + receivers: [otlp] + processors: [batch] + exporters: [otlp_http/prometheus] diff --git a/distributed_observability/otelmart/src/main.rs b/distributed_observability/otelmart/src/main.rs @@ -2,6 +2,7 @@ mod auth; mod config; mod db; mod handlers; +mod metrics; mod models; mod proxy; mod telemetry; diff --git a/distributed_observability/otelmart/src/metrics.rs b/distributed_observability/otelmart/src/metrics.rs @@ -0,0 +1,32 @@ +//! Business-specific metrics for the OtelMart gateway service. +//! +//! Tracks the duration of proxied requests to upstream microservices +//! so operators can monitor backend latency from the gateway's perspective. +//! Uses the `OnceLock` singleton pattern for lazy initialization. + +use opentelemetry::{global, metrics::Histogram}; +use std::sync::OnceLock; + +/// Collection of gateway-related metric instruments. +pub struct GatewayMetrics { + /// Duration of requests forwarded to upstream services (seconds) + pub upstream_request_duration: Histogram<f64>, +} + +/// Singleton storage for the metrics instruments. +static METRICS: OnceLock<GatewayMetrics> = OnceLock::new(); + +/// Returns a reference to the lazily-initialized gateway metrics. +pub fn metrics() -> &'static GatewayMetrics { + METRICS.get_or_init(|| { + let meter = opentelemetry::global::meter("otelmart-gateway"); + + GatewayMetrics { + upstream_request_duration: meter + .f64_histogram("http.client.request.duration") + .with_description("Duration of requests to upstream services") + .with_unit("s") + .build(), + } + }) +} diff --git a/distributed_observability/otelmart/src/proxy/mod.rs b/distributed_observability/otelmart/src/proxy/mod.rs @@ -2,12 +2,17 @@ pub mod inventory; pub mod orders; pub mod products; +use std::time::Instant; + use axum::{ body::Body, extract::Request, http::{HeaderValue, StatusCode}, response::{IntoResponse, Response}, }; +use opentelemetry::KeyValue; + +use crate::metrics::metrics; /// Generic proxy handler that forwards requests to a backend service /// @@ -40,6 +45,7 @@ pub async fn proxy_request( // Extract method and headers let method = req.method().clone(); + let method_str = method.to_string(); let headers = req.headers().clone(); // Extract body @@ -70,33 +76,83 @@ pub async fn proxy_request( client_req = client_req.body(body_bytes.to_vec()); } - // Send request to backend service - match client_req.send().await { - Ok(response) => { - let status = response.status(); - let mut builder = Response::builder().status(status); - - // Forward response headers - let resp_headers = response.headers().clone(); - for (name, value) in resp_headers.iter() { - if let Ok(val) = HeaderValue::from_bytes(value.as_bytes()) { - builder = builder.header(name.as_str(), val); - } - } + // Send request to backend service and measure duration + // The TracingMiddleware automatically: + // 1. Creates a client span with HTTP semantic convention attributes + // 2. Injects trace context (traceparent/tracestate) headers + // 3. Records response status on span completion + let start = Instant::now(); + let response = client_req.send().await; + let duration = start.elapsed().as_secs_f64(); - // Get response body - match response.bytes().await { - Ok(body) => builder.body(Body::from(body)).unwrap(), - Err(_) => Response::builder() - .status(StatusCode::INTERNAL_SERVER_ERROR) - .body(Body::from("Failed to read response body")) - .unwrap(), - } + // Derive the upstream service name from the target URL + let upstream_service = derive_upstream_service(service_url); + + // Handle request failure + let response = match response { + Ok(resp) => resp, + Err(_e) => { + // Record upstream request failure metric + metrics().upstream_request_duration.record( + duration, + &[ + KeyValue::new("upstream.service", upstream_service.clone()), + KeyValue::new("http.request.method", method_str), + KeyValue::new("error.type", "connection_error"), + ], + ); + + return Response::builder() + .status(StatusCode::SERVICE_UNAVAILABLE) + .header("content-type", "application/json") + .body(Body::from(r#"{"error":"Service unavailable"}"#)) + .unwrap(); + } + }; + + // Extract status and headers before consuming the response + let status = response.status(); + let resp_headers = response.headers().clone(); + + // Record successful upstream request duration metric + metrics().upstream_request_duration.record( + duration, + &[ + KeyValue::new("upstream.service", upstream_service), + KeyValue::new("http.request.method", method_str), + KeyValue::new("http.response.status_code", i64::from(status.as_u16())), + ], + ); + + // Get response body (this consumes the response) + let body = match response.bytes().await { + Ok(bytes) => bytes, + Err(_) => { + return Response::builder() + .status(StatusCode::INTERNAL_SERVER_ERROR) + .body(Body::from("Failed to read response body")) + .unwrap(); + } + }; + + // Build response with forwarded headers + let mut builder = Response::builder().status(status); + for (name, value) in resp_headers.iter() { + if let Ok(val) = HeaderValue::from_bytes(value.as_bytes()) { + builder = builder.header(name.as_str(), val); } - Err(_) => Response::builder() - .status(StatusCode::SERVICE_UNAVAILABLE) - .header("content-type", "application/json") - .body(Body::from(r#"{"error":"Service unavailable"}"#)) - .unwrap(), } + + builder.body(Body::from(body)).unwrap() +} + +/// Derives a short service name from a backend URL +/// (e.g. "http://products:3001" → "products"). +fn derive_upstream_service(url: &str) -> String { + url.trim_start_matches("http://") + .trim_start_matches("https://") + .split(":") + .next() + .unwrap_or("unknown") + .to_string() } diff --git a/distributed_observability/otelmart/src/telemetry.rs b/distributed_observability/otelmart/src/telemetry.rs @@ -67,11 +67,11 @@ pub fn init_telemetry(service_name: &str) -> TelemetryGuard { // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:4317".into()); + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + .unwrap_or_else(|_| "http://localhost:4317".into()); // Initialize both providers let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); diff --git a/distributed_observability/products/src/telemetry.rs b/distributed_observability/products/src/telemetry.rs @@ -67,11 +67,11 @@ pub fn init_telemetry(service_name: &str) -> TelemetryGuard { // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) let traces_endpoint = std::env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:4317".into()); + .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); let metrics_endpoint = std::env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT") .or_else(|_| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")) - .unwrap_or_else(|_| "http://localhost:9090/api/v1/otlp/v1/metrics".into()); + .unwrap_or_else(|_| "http://localhost:4317".into()); // Initialize both providers let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); diff --git a/distributed_observability/prometheus.yml b/distributed_observability/prometheus.yml @@ -0,0 +1,33 @@ +# Prometheus configuration for the distributed observability stack +# +# Metrics arrive by OTLP push (services export to the /api/v1/otlp endpoint, +# enabled via the --web.enable-otlp-receiver flag in docker-compose.yml), +# so the only scrape job here is Prometheus scraping itself. + +global: + scrape_interval: 15s + evaluation_interval: 15s + external_labels: + cluster: otelmart + environment: development + +# OTLP receiver settings for pushed metrics +otlp: + # Keep OTel resource attributes as labels on the ingested series, + # so metrics can be filtered by service in Grafana. + promote_resource_attributes: + - service.name + - service.namespace + - service.version + - service.instance.id + - deployment.environment + +storage: + tsdb: + # OTLP pushes can arrive slightly out of order across services. + out_of_order_time_window: 30m + +scrape_configs: + - job_name: prometheus + static_configs: + - targets: ["localhost:9090"] diff --git a/distributed_observability/scripts/export-dashboards.sh b/distributed_observability/scripts/export-dashboards.sh @@ -0,0 +1,85 @@ +#!/usr/bin/env bash +# Export Grafana dashboards to grafana/provisioning/dashboards/ so UI edits +# can be committed to git. +# +# Grafana loads dashboards from those JSON files on startup, but edits made in +# the UI are saved to Grafana's internal DB — not back to the files. Run this +# after editing a dashboard, then `git add` the result. +# +# Usage: +# ./scripts/export-dashboards.sh # export every dashboard +# ./scripts/export-dashboards.sh ad79bw7 # export one, by uid + +set -euo pipefail + +GRAFANA_URL="${GRAFANA_URL:-http://localhost:3000}" +GRAFANA_AUTH="${GRAFANA_AUTH:-admin:admin}" + +# Resolve the output dir relative to this script, so it works from any cwd. +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +OUT_DIR="$REPO_ROOT/grafana/provisioning/dashboards" +mkdir -p "$OUT_DIR" + +# Collect the uids to export: either the one passed in, or all dashboards. +if [ $# -gt 0 ]; then + uids=("$@") +else + mapfile -t uids < <( + curl -sf -u "$GRAFANA_AUTH" "$GRAFANA_URL/api/search?type=dash-db" \ + | python3 -c 'import sys,json; [print(d["uid"]) for d in json.load(sys.stdin)]' + ) +fi + +if [ ${#uids[@]} -eq 0 ]; then + echo "No dashboards found in $GRAFANA_URL" >&2 + exit 1 +fi + +for uid in "${uids[@]}"; do + curl -sf -u "$GRAFANA_AUTH" "$GRAFANA_URL/api/dashboards/uid/$uid" \ + | OUT_DIR="$OUT_DIR" python3 -c ' +import json, os, re, sys + +payload = json.load(sys.stdin) +dash = payload["dashboard"] + +# `id` is local to a Grafana instance and must not be provisioned; `uid` and +# `title` are kept so this file updates the existing dashboard in place +# instead of creating a duplicate. +dash.pop("id", None) + +out_dir = os.environ["OUT_DIR"] + +# Reuse the file that already holds this uid, so re-exporting (or renaming a +# dashboard) updates it in place. Two files sharing a uid would make Grafana +# provisioning fight over which one wins. +path = None +for name in sorted(os.listdir(out_dir)): + if not name.endswith(".json"): + continue + candidate = os.path.join(out_dir, name) + try: + with open(candidate) as fh: + if json.load(fh).get("uid") == dash["uid"]: + path = candidate + break + except (json.JSONDecodeError, OSError): + continue + +# No existing file: name it from the title, lowercased, non-alphanumerics dashed. +if path is None: + slug = re.sub(r"[^a-z0-9]+", "-", dash["title"].lower()).strip("-") or dash["uid"] + path = os.path.join(out_dir, slug + ".json") + +with open(path, "w") as fh: + json.dump(dash, fh, indent=2) + fh.write("\n") + +title = dash["title"] +print(f" {title} -> grafana/provisioning/dashboards/{os.path.basename(path)}") +' +done + +echo +echo "Done. Review and commit:" +echo " git add grafana/provisioning/dashboards/ && git commit -m 'update dashboards'" diff --git a/distributed_observability/scripts/generate_traffic.sh b/distributed_observability/scripts/generate_traffic.sh @@ -0,0 +1,626 @@ +#!/bin/bash +# Realistic E-Commerce Traffic Generator for Chapter 6 +# +# Simulates realistic user behavior patterns: +# - User sessions (30s-5min) with multiple concurrent users +# - Product popularity distribution (price-based) +# - User personas: Window Shoppers, Buyers, Power Buyers, Bots +# - Realistic timing patterns and cart abandonment +# - Clustered error scenarios (payment outages, stock depletion) +# +# Prerequisites: +# - The Chapter 6 docker stack is up: docker compose up -d +# - bash, curl, jq are available +# - The gateway responds at $GATEWAY_URL/health (default: http://localhost:4200) +# +# Usage: +# ./generate_traffic.sh # Run for 2 minutes (default) +# ./generate_traffic.sh --pv # Also print Prometheus verification queries at the end +# ./generate_traffic.sh -h | --help # Show this help and exit +# DURATION=300 ./generate_traffic.sh # Run for 5 minutes +# GATEWAY_URL=http://otelmart:4200 ./generate_traffic.sh +# # Hit the gateway by container name (run from inside docker) +# +# Environment variables: +# GATEWAY_URL Base URL of the otelmart gateway. Default: http://localhost:4200 +# DURATION How long to run, in seconds. Default: 120 +# +# Running on Windows (no native bash): +# Use the alpine/curl image on the same docker network as the stack: +# +# docker run --rm --network chapter_6_app-network \ +# -v ${PWD}/scripts:/scripts \ +# -e GATEWAY_URL=http://otelmart:4200 \ +# -e DURATION=120 \ +# alpine/curl:latest \ +# sh -c "apk add --no-cache jq bash >/dev/null && bash /scripts/generate_traffic.sh" +# +# After the run, point Prometheus / Grafana at http://localhost:9090 and try the +# example PromQL queries from Chapter 6 (sections 6.2, 6.4, 6.6). + +set -e + +# Configuration +GATEWAY_URL="${GATEWAY_URL:-http://localhost:4200}" +DURATION="${DURATION:-120}" # Default: 2 minutes +VERIFY_PROMETHEUS=false + +print_help() { + sed -n '2,40p' "$0" | sed 's/^# \{0,1\}//' +} + +# Parse arguments +while [[ $# -gt 0 ]]; do + case $1 in + --pv|--prometheus-verify) + VERIFY_PROMETHEUS=true + shift + ;; + -h|--help) + print_help + exit 0 + ;; + *) + echo "Unknown option: $1" + echo "Usage: $0 [--pv] [-h|--help]" + echo "Run '$0 --help' for full instructions." + exit 1 + ;; + esac +done + +# Colors for output +RED='\033[0;31m' +GREEN='\033[0;32m' +YELLOW='\033[1;33m' +BLUE='\033[0;34m' +CYAN='\033[0;36m' +NC='\033[0m' # No Color + +echo "=========================================" +echo "OpenTel E-Commerce Traffic Generator" +echo "=========================================" +echo "Gateway URL: $GATEWAY_URL" +echo "Duration: ${DURATION}s" +echo "Prometheus Verification: $VERIFY_PROMETHEUS" +echo "" + +# Check if gateway is reachable +echo "Checking service health..." +if ! curl -s -f "$GATEWAY_URL/health" > /dev/null 2>&1; then + echo -e "${RED}ERROR: Gateway service not reachable at $GATEWAY_URL${NC}" + echo "Make sure the services are running: docker-compose up" + exit 1 +fi + +echo -e "${GREEN}✓ Gateway is healthy${NC}" +echo "" + +# Counters +TOTAL_REQUESTS=0 +SUCCESS_COUNT=0 +ERROR_404_COUNT=0 +ERROR_409_COUNT=0 +ERROR_422_COUNT=0 +ERROR_500_COUNT=0 +ORDER_SUCCESS_COUNT=0 +ORDER_FAILURE_COUNT=0 +SESSION_COUNT=0 + +# Fetch products once +echo "Fetching products..." +PRODUCTS_JSON=$(curl -s "$GATEWAY_URL/api/products?page_size=50") + +if [ -z "$PRODUCTS_JSON" ] || [ "$PRODUCTS_JSON" = "null" ]; then + echo -e "${RED}ERROR: Failed to fetch products from gateway${NC}" + exit 1 +fi + +# Parse products into arrays +PRODUCT_UUIDS=($(echo "$PRODUCTS_JSON" | jq -r '.products[].eid')) +PRODUCT_NAMES=($(echo "$PRODUCTS_JSON" | jq -r '.products[].product_name')) +PRODUCT_PRICES=($(echo "$PRODUCTS_JSON" | jq -r '.products[].final_price')) + +if [ ${#PRODUCT_UUIDS[@]} -eq 0 ]; then + echo -e "${RED}ERROR: No products found. Make sure the database is populated.${NC}" + exit 1 +fi + +echo -e "${GREEN}✓ Found ${#PRODUCT_UUIDS[@]} products${NC}" + +# Create product popularity index (price-based: expensive = more popular) +# Sort products by price descending, top 20% are "popular" +POPULAR_COUNT=$((${#PRODUCT_UUIDS[@]} / 5)) +if [ $POPULAR_COUNT -lt 3 ]; then + POPULAR_COUNT=3 +fi + +# Create popularity tiers +POPULAR_PRODUCTS=() +REGULAR_PRODUCTS=() + +for i in "${!PRODUCT_PRICES[@]}"; do + if [ $i -lt $POPULAR_COUNT ]; then + POPULAR_PRODUCTS+=("$i") + else + REGULAR_PRODUCTS+=("$i") + fi +done + +echo -e "${CYAN}✓ Product popularity: ${#POPULAR_PRODUCTS[@]} popular, ${#REGULAR_PRODUCTS[@]} regular${NC}" +echo "" + +# Helper: Get product with popularity bias +get_product_by_popularity() { + # 80% chance to select popular product, 20% regular + if [ $((RANDOM % 100)) -lt 80 ] && [ ${#POPULAR_PRODUCTS[@]} -gt 0 ]; then + local idx=${POPULAR_PRODUCTS[$((RANDOM % ${#POPULAR_PRODUCTS[@]}))]} + else + local idx=${REGULAR_PRODUCTS[$((RANDOM % ${#REGULAR_PRODUCTS[@]}))]} + fi + echo "${PRODUCT_UUIDS[$idx]}|${PRODUCT_NAMES[$idx]}|${PRODUCT_PRICES[$idx]}" +} + +# Helper: Random delay with activity context +think_time() { + local activity=$1 + case $activity in + "quick_scan") + sleep $(awk 'BEGIN{srand(); print 1+rand()*2}') # 1-3s + ;; + "product_view") + sleep $(awk 'BEGIN{srand(); print 5+rand()*10}') # 5-15s + ;; + "decision") + sleep $(awk 'BEGIN{srand(); print 2+rand()*2}') # 2-4s + ;; + "checkout") + sleep $(awk 'BEGIN{srand(); print 10+rand()*20}') # 10-30s + ;; + "rapid") + sleep $(awk 'BEGIN{srand(); print 0.1+rand()*0.4}') # 0.1-0.5s (bot) + ;; + *) + sleep $(awk 'BEGIN{srand(); print 0.5+rand()*1.5}') # 0.5-2s default + ;; + esac +} + +# User Persona: Window Shopper (50% of users) +# Browses 5-10 products, checks stock, abandons cart +simulate_window_shopper() { + local session_id=$1 + local browse_count=$((5 + RANDOM % 6)) # 5-10 products + + echo -e "${BLUE}[Session $session_id] Window Shopper: browsing $browse_count products${NC}" + + # Browse product list + curl -s "$GATEWAY_URL/api/products?page_size=20" > /dev/null + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + think_time "quick_scan" + + # View individual products + for i in $(seq 1 $browse_count); do + local product_info=$(get_product_by_popularity) + local uuid=$(echo "$product_info" | cut -d'|' -f1) + + curl -s "$GATEWAY_URL/api/products/$uuid" > /dev/null + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + think_time "product_view" + + # 30% chance to check stock (showing intent) + if [ $((RANDOM % 100)) -lt 30 ]; then + curl -s "$GATEWAY_URL/api/inventory/$uuid/stock" > /dev/null + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + think_time "decision" + fi + done + + echo -e "${YELLOW}[Session $session_id] Window Shopper: abandoned (no purchase)${NC}" +} + +# User Persona: Casual Buyer (30% of users) +# Browses 2-5 products, completes purchase +simulate_casual_buyer() { + local session_id=$1 + local browse_count=$((2 + RANDOM % 4)) # 2-5 products + + echo -e "${BLUE}[Session $session_id] Casual Buyer: browsing $browse_count products${NC}" + + # Browse products + for i in $(seq 1 $browse_count); do + local product_info=$(get_product_by_popularity) + local uuid=$(echo "$product_info" | cut -d'|' -f1) + + curl -s "$GATEWAY_URL/api/products/$uuid" > /dev/null + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + think_time "product_view" + done + + # Check stock before purchase + local product_info=$(get_product_by_popularity) + local uuid=$(echo "$product_info" | cut -d'|' -f1) + local name=$(echo "$product_info" | cut -d'|' -f2) + local price=$(echo "$product_info" | cut -d'|' -f3) + + curl -s "$GATEWAY_URL/api/inventory/$uuid/stock" > /dev/null + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + think_time "decision" + + # Create order + think_time "checkout" + create_order "$uuid" "$name" "$price" 1 "$session_id" +} + +# User Persona: Power Buyer (15% of users) +# Quick browse, multi-item orders +simulate_power_buyer() { + local session_id=$1 + local item_count=$((2 + RANDOM % 4)) # 2-5 items + + echo -e "${BLUE}[Session $session_id] Power Buyer: purchasing $item_count items${NC}" + + # Quick product check + curl -s "$GATEWAY_URL/api/products?page_size=10" > /dev/null + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + think_time "quick_scan" + + # Create multi-item order + local product_info=$(get_product_by_popularity) + local uuid=$(echo "$product_info" | cut -d'|' -f1) + local name=$(echo "$product_info" | cut -d'|' -f2) + local price=$(echo "$product_info" | cut -d'|' -f3) + + think_time "checkout" + create_order "$uuid" "$name" "$price" "$item_count" "$session_id" +} + +# User Persona: Bot/Scraper (5% of users) +# Rapid requests, high error rate +simulate_bot() { + local session_id=$1 + local request_count=$((10 + RANDOM % 20)) # 10-30 rapid requests + + echo -e "${CYAN}[Session $session_id] Bot: scraping $request_count pages${NC}" + + for i in $(seq 1 $request_count); do + # Mix of valid and invalid requests + if [ $((RANDOM % 100)) -lt 30 ]; then + # 30% invalid requests (404s) + curl -s "$GATEWAY_URL/api/products/00000000-0000-0000-0000-000000000000" > /dev/null 2>&1 + ERROR_404_COUNT=$((ERROR_404_COUNT + 1)) + else + # Valid product requests + curl -s "$GATEWAY_URL/api/products?page_size=50" > /dev/null + fi + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + think_time "rapid" + done +} + +# Create order with realistic error handling +create_order() { + local uuid=$1 + local name=$2 + local price=$3 + local quantity=$4 + local session_id=$5 + + local response=$(curl -s -w "\n%{http_code}" -X POST "$GATEWAY_URL/api/orders" \ + -H "Content-Type: application/json" \ + -d "{ + \"customer_email\": \"user${session_id}@example.com\", + \"items\": [{ + \"product_uuid\": \"$uuid\", + \"product_name\": \"$name\", + \"quantity\": $quantity, + \"unit_price\": \"$price\" + }], + \"shipping_address\": { + \"first_name\": \"User\", + \"last_name\": \"${session_id}\", + \"address_line1\": \"123 Main St\", + \"city\": \"San Francisco\", + \"state\": \"CA\", + \"postal_code\": \"94105\", + \"country\": \"US\" + }, + \"payment\": { + \"payment_method\": \"credit_card\", + \"card_last4\": \"4242\", + \"card_brand\": \"visa\" + } + }") + + local http_code=$(echo "$response" | tail -n1) + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + + if [ "$http_code" = "201" ] || [ "$http_code" = "200" ]; then + SUCCESS_COUNT=$((SUCCESS_COUNT + 1)) + ORDER_SUCCESS_COUNT=$((ORDER_SUCCESS_COUNT + 1)) + echo -e "${GREEN}[Session $session_id] ✓ Order created ($quantity items)${NC}" + elif [ "$http_code" = "409" ]; then + ERROR_409_COUNT=$((ERROR_409_COUNT + 1)) + ORDER_FAILURE_COUNT=$((ORDER_FAILURE_COUNT + 1)) + echo -e "${YELLOW}[Session $session_id] ⚠ Order failed: Insufficient stock (409)${NC}" + else + ORDER_FAILURE_COUNT=$((ORDER_FAILURE_COUNT + 1)) + echo -e "${RED}[Session $session_id] ✗ Order failed: HTTP $http_code${NC}" + fi +} + +# Simulate payment gateway outage (500 errors) +simulate_payment_outage() { + echo -e "${RED}[INCIDENT] Payment gateway outage (30s)${NC}" + local outage_end=$(($(date +%s) + 30)) + + while [ $(date +%s) -lt $outage_end ]; do + # Generate validation errors (simulating gateway timeout) + local product_info=$(get_product_by_popularity) + local uuid=$(echo "$product_info" | cut -d'|' -f1) + local name=$(echo "$product_info" | cut -d'|' -f2) + local price=$(echo "$product_info" | cut -d'|' -f3) + + curl -s -X POST "$GATEWAY_URL/api/orders" \ + -H "Content-Type: application/json" \ + -d "{ + \"customer_email\": \"outage@example.com\", + \"items\": [{ + \"product_uuid\": \"$uuid\", + \"product_name\": \"$name\", + \"quantity\": 1, + \"unit_price\": \"$price\" + }], + \"shipping_address\": { + \"first_name\": \"Test\", + \"last_name\": \"User\", + \"address_line1\": \"123 Main St\", + \"city\": \"San Francisco\", + \"state\": \"CA\", + \"postal_code\": \"94105\", + \"country\": \"INVALID\" + }, + \"payment\": { + \"payment_method\": \"credit_card\", + \"card_last4\": \"4242\", + \"card_brand\": \"visa\" + } + }" > /dev/null 2>&1 + + TOTAL_REQUESTS=$((TOTAL_REQUESTS + 1)) + ERROR_422_COUNT=$((ERROR_422_COUNT + 1)) + sleep 0.5 + done + + echo -e "${GREEN}[INCIDENT] Payment gateway recovered${NC}" +} + +# Exhaust stock on popular products +exhaust_popular_stock() { + echo -e "${YELLOW}[INCIDENT] Stock depletion wave on popular products${NC}" + + for i in {1..3}; do + if [ ${#POPULAR_PRODUCTS[@]} -gt $i ]; then + local idx=${POPULAR_PRODUCTS[$i]} + local uuid="${PRODUCT_UUIDS[$idx]}" + local name="${PRODUCT_NAMES[$idx]}" + local price="${PRODUCT_PRICES[$idx]}" + + # Attempt large orders to exhaust stock + for j in {1..5}; do + create_order "$uuid" "$name" "$price" 10 "stock_depletion_$i" + sleep 0.2 + done + fi + done +} + +# Main traffic generation loop +echo "=========================================" +echo "Generating realistic traffic for ${DURATION}s..." +echo "=========================================" +echo "" + +START_TIME=$(date +%s) +END_TIME=$((START_TIME + DURATION)) + +# Background session spawner +spawn_user_session() { + local persona=$1 + SESSION_COUNT=$((SESSION_COUNT + 1)) + local session_id=$SESSION_COUNT + + ( + case $persona in + "window_shopper") + simulate_window_shopper $session_id + ;; + "casual_buyer") + simulate_casual_buyer $session_id + ;; + "power_buyer") + simulate_power_buyer $session_id + ;; + "bot") + simulate_bot $session_id + ;; + esac + ) & +} + +PHASE=1 +INCIDENT_TRIGGERED=false + +while [ $(date +%s) -lt $END_TIME ]; do + CURRENT_TIME=$(date +%s) + ELAPSED=$((CURRENT_TIME - START_TIME)) + REMAINING=$((END_TIME - CURRENT_TIME)) + + # Progress indicator + if [ $((ELAPSED % 15)) -eq 0 ] && [ $ELAPSED -gt 0 ]; then + echo -e "${BLUE}[${ELAPSED}s/${DURATION}s] Phase $PHASE - Sessions: $SESSION_COUNT, Requests: $TOTAL_REQUESTS, Orders: $ORDER_SUCCESS_COUNT${NC}" + fi + + # Phase 1 (0-30s): Morning Traffic - Low Volume + if [ $ELAPSED -lt 30 ]; then + PHASE=1 + # Spawn 1-2 concurrent sessions + persona_roll=$((RANDOM % 100)) + if [ $persona_roll -lt 70 ]; then + spawn_user_session "window_shopper" + elif [ $persona_roll -lt 90 ]; then + spawn_user_session "casual_buyer" + else + spawn_user_session "power_buyer" + fi + sleep 2 + + # Phase 2 (30-60s): Lunch Rush - Burst Traffic + elif [ $ELAPSED -lt 60 ]; then + PHASE=2 + # Spawn 3-5 concurrent sessions + for i in {1..3}; do + persona_roll=$((RANDOM % 100)) + if [ $persona_roll -lt 50 ]; then + spawn_user_session "window_shopper" + elif [ $persona_roll -lt 90 ]; then + spawn_user_session "casual_buyer" + else + spawn_user_session "power_buyer" + fi + done + sleep 1 + + # Phase 3 (60-90s): Afternoon - Incidents + elif [ $ELAPSED -lt 90 ]; then + PHASE=3 + + # Trigger payment outage once + if [ $ELAPSED -eq 65 ] && [ "$INCIDENT_TRIGGERED" = false ]; then + simulate_payment_outage & + INCIDENT_TRIGGERED=true + fi + + # Trigger stock depletion + if [ $ELAPSED -eq 75 ]; then + exhaust_popular_stock & + fi + + # Normal traffic with bots + persona_roll=$((RANDOM % 100)) + if [ $persona_roll -lt 50 ]; then + spawn_user_session "window_shopper" + elif [ $persona_roll -lt 75 ]; then + spawn_user_session "casual_buyer" + elif [ $persona_roll -lt 90 ]; then + spawn_user_session "power_buyer" + else + spawn_user_session "bot" + fi + sleep 1.5 + + # Phase 4 (90-120s): Evening Peak - Sustained High Traffic + else + PHASE=4 + # Spawn 2-4 concurrent sessions with higher buyer ratio + for i in {1..2}; do + persona_roll=$((RANDOM % 100)) + if [ $persona_roll -lt 40 ]; then + spawn_user_session "window_shopper" + elif [ $persona_roll -lt 85 ]; then + spawn_user_session "casual_buyer" + else + spawn_user_session "power_buyer" + fi + done + sleep 1 + fi +done + +# Wait for all background sessions to complete +echo "" +echo "Waiting for active sessions to complete..." +wait + +echo "" +echo "=========================================" +echo "Traffic Generation Complete!" +echo "=========================================" +echo "Duration: ${DURATION}s" +echo "Total Sessions: $SESSION_COUNT" +echo "Total Requests: $TOTAL_REQUESTS" +echo "Successful Responses: $SUCCESS_COUNT" +echo "" +echo "Error Breakdown:" +echo " 404 (Not Found): $ERROR_404_COUNT" +echo " 409 (Conflict/Stock): $ERROR_409_COUNT" +echo " 422 (Validation): $ERROR_422_COUNT" +echo "" +echo "Order Statistics:" +echo " Successful Orders: $ORDER_SUCCESS_COUNT" +echo " Failed Orders: $ORDER_FAILURE_COUNT" +if [ $((ORDER_SUCCESS_COUNT + ORDER_FAILURE_COUNT)) -gt 0 ]; then + echo " Success Rate: $(awk "BEGIN {printf \"%.1f%%\", ($ORDER_SUCCESS_COUNT/($ORDER_SUCCESS_COUNT+$ORDER_FAILURE_COUNT))*100}")" +fi +echo "" + +# Optional: Verify metrics in Prometheus +if [ "$VERIFY_PROMETHEUS" = true ]; then + echo "=========================================" + echo "Verifying Metrics in Prometheus" + echo "=========================================" + + PROM_URL="http://localhost:9090" + + echo "Checking Prometheus availability..." + if ! curl -s -f "$PROM_URL/-/healthy" > /dev/null 2>&1; then + echo -e "${YELLOW}⚠ Prometheus not reachable at $PROM_URL${NC}" + echo "Skipping verification." + else + echo -e "${GREEN}✓ Prometheus is healthy${NC}" + echo "" + + echo "Querying metrics..." + + # HTTP Server Metrics (RED) + HTTP_METRICS=$(curl -s "$PROM_URL/api/v1/query?query=http_server_request_duration_seconds_count" | jq -r '.data.result | length') + echo " HTTP Server Metrics: $HTTP_METRICS series" + + # Order Metrics + ORDER_ATTEMPTS=$(curl -s "$PROM_URL/api/v1/query?query=orders_checkout_attempts_total" | jq -r '.data.result[0].value[1] // "0"') + echo " Order Attempts: $ORDER_ATTEMPTS" + + # Funnel Metrics + FUNNEL_STARTED=$(curl -s "$PROM_URL/api/v1/query?query=checkout_funnel_started_total" | jq -r '.data.result[0].value[1] // "0"') + FUNNEL_COMPLETED=$(curl -s "$PROM_URL/api/v1/query?query=checkout_funnel_completed_total" | jq -r '.data.result[0].value[1] // "0"') + echo " Funnel Started: $FUNNEL_STARTED" + echo " Funnel Completed: $FUNNEL_COMPLETED" + + if [ "$FUNNEL_STARTED" != "0" ]; then + CONVERSION=$(awk "BEGIN {printf \"%.1f%%\", ($FUNNEL_COMPLETED/$FUNNEL_STARTED)*100}") + echo " Conversion Rate: $CONVERSION" + fi + + # Inventory Metrics + INV_ATTEMPTS=$(curl -s "$PROM_URL/api/v1/query?query=inventory_reservation_attempts_total" | jq -r '.data.result[0].value[1] // "0"') + echo " Inventory Reservations: $INV_ATTEMPTS" + + # Gateway Metrics + GATEWAY_METRICS=$(curl -s "$PROM_URL/api/v1/query?query=http_client_request_duration_seconds_count" | jq -r '.data.result | length') + echo " Gateway Upstream Metrics: $GATEWAY_METRICS series" + + echo "" + echo -e "${GREEN}✓ Metrics verification complete${NC}" + fi + echo "" +fi + +echo "View metrics in Prometheus: http://localhost:9090" +echo "View dashboards in Grafana: http://localhost:3000" +echo "" +echo "Example PromQL queries:" +echo " Request rate: rate(http_server_request_duration_seconds_count[5m])" +echo " Error rate: (sum(rate(http_server_request_duration_seconds_count{http_response_status_code=~\"4..|5..\"}[5m])) / sum(rate(http_server_request_duration_seconds_count[5m])) * 100) or vector(0)" +echo " P95 latency: histogram_quantile(0.95, rate(http_server_request_duration_seconds_bucket[5m]))" +echo " Conversion rate: ((rate(checkout_funnel_completed_total[5m]) / rate(checkout_funnel_started_total[5m])) * 100) or vector(0)" +echo "" diff --git a/distributed_observability/tempo.yaml b/distributed_observability/tempo.yaml @@ -0,0 +1,50 @@ +# Tempo configuration for the distributed observability stack +# +# Receives traces over OTLP (gRPC 4317 / HTTP 4318) directly from the services +# and is queried through Grafana at :3200 — Tempo has no UI of its own. +# +# Note: this targets Tempo v3.x, which replaced the v2 `ingester` and +# `compactor` config blocks with a block-builder architecture. Those keys +# are rejected outright by v3, so they are intentionally absent here. + +server: + http_listen_port: 3200 + grpc_listen_port: 9096 + +distributor: + receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:4317 + http: + endpoint: 0.0.0.0:4318 + +# Derives RED metrics and service-graph data from spans and remote-writes +# them to Prometheus. This is what powers Traces Drilldown's rate/error/ +# duration views and the service map. +metrics_generator: + registry: + external_labels: + source: tempo + storage: + path: /var/tempo/generator/wal + remote_write: + - url: http://prometheus:9090/api/v1/write + send_exemplars: true + +storage: + trace: + backend: local # local filesystem is fine for dev; use S3/GCS in prod + wal: + path: /var/tempo/wal + local: + path: /var/tempo/blocks + +overrides: + defaults: + metrics_generator: + processors: + - service-graphs + - span-metrics + - local-blocks # required for TraceQL metrics in Traces Drilldown