commit b9a894ee4ddce1f74a65473618b24c59eb7bf5d8 parent 35ea0ca585c9f712c8ba5f44b30e8adaccbccd1d Author: ling0x <ling0x@users.noreply.github.com> Date: Sun, 30 Aug 2026 19:04:33 +0100 add logging and loki Diffstat:
17 files changed, 518 insertions(+), 131 deletions(-)
diff --git a/distributed_observability/Cargo.lock b/distributed_observability/Cargo.lock @@ -266,6 +266,15 @@ dependencies = [ ] [[package]] +name = "block-buffer" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d2f6c7dbe95a6ed67ad9f18e57daf93a2f034c524b99fd2b76d18fdfeb6660aa" +dependencies = [ + "hybrid-array", +] + +[[package]] name = "blowfish" version = "0.9.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -348,7 +357,7 @@ version = "0.4.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773f3b9af64447d2ce9850330c473515014aa235e6a783b02db81ff39e4a3dad" dependencies = [ - "crypto-common", + "crypto-common 0.1.7", "inout", ] @@ -406,6 +415,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" [[package]] +name = "const-oid" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6ef517f0926dd24a1582492c791b6a4818a4d94e789a334894aa15b0d12f55c" + +[[package]] name = "const-random" version = "0.1.18" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -525,12 +540,21 @@ dependencies = [ ] [[package]] +name = "crypto-common" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ce6e4c961d6cd6c9a86db418387425e8bdeaf05b3c8bc1411e6dca4c252f1453" +dependencies = [ + "hybrid-array", +] + +[[package]] name = "der" version = "0.7.10" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e7c1832837b905bbfb5101e07cc24c8deddf52f93225eee6ead5f4d63d53ddcb" dependencies = [ - "const-oid", + "const-oid 0.9.6", "pem-rfc7468", "zeroize", ] @@ -550,13 +574,24 @@ version = "0.10.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" dependencies = [ - "block-buffer", - "const-oid", - "crypto-common", + "block-buffer 0.10.4", + "const-oid 0.9.6", + "crypto-common 0.1.7", "subtle", ] [[package]] +name = "digest" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" +dependencies = [ + "block-buffer 0.12.1", + "const-oid 0.10.2", + "crypto-common 0.2.2", +] + +[[package]] name = "displaydoc" version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -902,7 +937,7 @@ version = "0.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6c49c37c09c17a53d937dfbb742eb3a961d65a994e6bcdcf37e7399d0cc8ab5e" dependencies = [ - "digest", + "digest 0.10.7", ] [[package]] @@ -967,6 +1002,15 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" [[package]] +name = "hybrid-array" +version = "0.4.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "707114b52a152fa7bdb290cd7cd5912d9467273b6d74e21b8d81aca1f8533f6b" +dependencies = [ + "typenum", +] + +[[package]] name = "hyper" version = "1.8.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1201,13 +1245,16 @@ dependencies = [ "chrono", "config", "dotenvy", + "hex", "opentelemetry", + "opentelemetry-appender-tracing", "opentelemetry-otlp", "opentelemetry_sdk", "reqwest-middleware", "reqwest-tracing", "serde", "serde_json", + "sha2 0.11.0", "sqlx", "thiserror", "tokio", @@ -1426,7 +1473,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d89e7ee0cfbedfc4da3340218492196241d89eefb6dab27de5df917a6d2e78cf" dependencies = [ "cfg-if", - "digest", + "digest 0.10.7", ] [[package]] @@ -1576,6 +1623,18 @@ dependencies = [ ] [[package]] +name = "opentelemetry-appender-tracing" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2c0080f0dc1d7c786f467cd85a4e395fcab11ee852004f39a29a18ab7c25d837" +dependencies = [ + "opentelemetry", + "tracing", + "tracing-core", + "tracing-subscriber", +] + +[[package]] name = "opentelemetry-http" version = "0.32.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1640,6 +1699,8 @@ dependencies = [ "portable-atomic", "rand 0.9.5", "thiserror", + "tokio", + "tokio-stream", ] [[package]] @@ -1664,7 +1725,9 @@ dependencies = [ "chrono", "config", "dotenvy", + "hex", "opentelemetry", + "opentelemetry-appender-tracing", "opentelemetry-otlp", "opentelemetry_sdk", "reqwest", @@ -1672,6 +1735,7 @@ dependencies = [ "reqwest-tracing", "serde", "serde_json", + "sha2 0.11.0", "sqlx", "thiserror", "tokio", @@ -1695,8 +1759,10 @@ dependencies = [ "chrono", "config", "dotenvy", + "hex", "jsonwebtoken", "opentelemetry", + "opentelemetry-appender-tracing", "opentelemetry-otlp", "opentelemetry_sdk", "reqwest", @@ -1704,6 +1770,7 @@ dependencies = [ "reqwest-tracing", "serde", "serde_json", + "sha2 0.11.0", "sqlx", "thiserror", "tokio", @@ -1816,7 +1883,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72f27a2cfee9f9039c4d86faa5af122a0ac3851441a34865b8a043b46be0065a" dependencies = [ "pest", - "sha2", + "sha2 0.10.9", ] [[package]] @@ -1929,13 +1996,16 @@ dependencies = [ "chrono", "config", "dotenvy", + "hex", "opentelemetry", + "opentelemetry-appender-tracing", "opentelemetry-otlp", "opentelemetry_sdk", "reqwest-middleware", "reqwest-tracing", "serde", "serde_json", + "sha2 0.11.0", "sqlx", "thiserror", "tokio", @@ -2275,8 +2345,8 @@ version = "0.9.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40a0376c50d0358279d9d643e4bf7b7be212f1f4ff1da9070a7b54d22ef75c88" dependencies = [ - "const-oid", - "digest", + "const-oid 0.9.6", + "digest 0.10.7", "num-bigint-dig", "num-integer", "num-traits", @@ -2538,7 +2608,7 @@ checksum = "e3bf829a2d51ab4a5ddf1352d8470c140cadc8301b2ae1789db023f01cedd6ba" dependencies = [ "cfg-if", "cpufeatures 0.2.17", - "digest", + "digest 0.10.7", ] [[package]] @@ -2549,7 +2619,18 @@ checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" dependencies = [ "cfg-if", "cpufeatures 0.2.17", - "digest", + "digest 0.10.7", +] + +[[package]] +name = "sha2" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "446ba717509524cb3f22f17ecc096f10f4822d76ab5c0b9822c5f9c284e825f4" +dependencies = [ + "cfg-if", + "cpufeatures 0.3.0", + "digest 0.11.3", ] [[package]] @@ -2582,7 +2663,7 @@ version = "2.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de" dependencies = [ - "digest", + "digest 0.10.7", "rand_core 0.6.4", ] @@ -2699,7 +2780,7 @@ dependencies = [ "rustls", "serde", "serde_json", - "sha2", + "sha2 0.10.9", "smallvec", "thiserror", "tokio", @@ -2738,7 +2819,7 @@ dependencies = [ "quote", "serde", "serde_json", - "sha2", + "sha2 0.10.9", "sqlx-core", "sqlx-mysql", "sqlx-postgres", @@ -2762,7 +2843,7 @@ dependencies = [ "bytes", "chrono", "crc", - "digest", + "digest 0.10.7", "dotenvy", "either", "futures-channel", @@ -2783,7 +2864,7 @@ dependencies = [ "rsa", "serde", "sha1", - "sha2", + "sha2 0.10.9", "smallvec", "sqlx-core", "stringprep", @@ -2824,7 +2905,7 @@ dependencies = [ "rand 0.8.5", "serde", "serde_json", - "sha2", + "sha2 0.10.9", "smallvec", "sqlx-core", "stringprep", @@ -3350,9 +3431,9 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "typenum" -version = "1.19.0" +version = "1.20.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "562d481066bde0658276a35467c4af00bdc6ee726305698a55b86e61d7ad82bb" +checksum = "b6f5e870be6c3b371b77fe0ee0bafb859fa4964b4404c27de1d380043c4dda20" [[package]] name = "ucd-trie" diff --git a/distributed_observability/Cargo.toml b/distributed_observability/Cargo.toml @@ -86,10 +86,23 @@ tracing-subscriber = { version = "0.3.22", features = ["env-filter", "json"] } # OpenTelemetry core opentelemetry = "0.32.0" -opentelemetry_sdk = "0.32.1" +opentelemetry_sdk = { version = "0.32.1", features = [ + "rt-tokio", # Tokio runtime support + "metrics", # Metrics support + "logs", # Logs support +] } +opentelemetry-appender-tracing = "0.32.0" # OTLP exporter (sends traces to collectors/backends) -opentelemetry-otlp = { version = "0.32.0", features = ["grpc-tonic"] } +opentelemetry-otlp = { version = "0.32.0", default-features = false, features = [ + "grpc-tonic", # gRPC for traces + "http-proto", # HTTP transport for metrics and Logs exporters + "reqwest-blocking-client", # HTTP client implementation + "trace", # Traces export + "metrics", # Metrics export + "logs", # Logs export + "internal-logs", # Internal logs export +] } # Bridge between tracing and OpenTelemetry tracing-opentelemetry = "0.33.0" @@ -100,6 +113,10 @@ reqwest-middleware = { version = "0.5.2", features = ["json"] } reqwest-tracing = { version = "0.7.1", features = ["opentelemetry_0_32"] } axum-otel-metrics = "0.14.1" +# Hash utils +hex = "0.4.3" +sha2 = "0.11.0" + [profile.release] opt-level = 3 lto = true diff --git a/distributed_observability/docker-compose.yml b/distributed_observability/docker-compose.yml @@ -42,6 +42,23 @@ services: restart: unless-stopped # ============================================================ + # Loki - Log Aggregation Backend + # ============================================================ + # Receives logs via OTLP HTTP from OpenTelemetry SDK + loki: + image: grafana/loki:3.0.0 + container_name: loki + ports: + - "3100:3100" + volumes: + - ./loki-config.yaml:/etc/loki/local-config.yaml:ro + - loki_data:/loki + command: -config.file=/etc/loki/local-config.yaml + networks: + - app-network + restart: unless-stopped + + # ============================================================ # Grafana - Dashboards & Visualization # ============================================================ # Queries both Prometheus (metrics) and Tempo (traces). @@ -58,6 +75,7 @@ services: depends_on: - prometheus - tempo + - loki networks: - app-network restart: unless-stopped @@ -104,8 +122,9 @@ services: - APP_DATABASE_URL=postgres://opentel_user:opentel_pass@postgres:5432/opentel_db - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics + - OTEL_EXPORTER_OTLP_LOGS_ENDPOINT=http://loki:3100/otlp/v1/logs - OTEL_SERVICE_NAME=products - - RUST_LOG=info,otel::tracing=trace + - RUST_LOG=info,open_telemetry=off,otel::tracing=trace depends_on: postgres: condition: service_healthy @@ -132,8 +151,9 @@ services: - APP_DATABASE_URL=postgres://opentel_user:opentel_pass@postgres:5432/opentel_db - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics + - OTEL_EXPORTER_OTLP_LOGS_ENDPOINT=http://loki:3100/otlp/v1/logs - OTEL_SERVICE_NAME=inventory - - RUST_LOG=info,otel::tracing=trace + - RUST_LOG=info,open_telemetry=off,otel::tracing=trace depends_on: postgres: condition: service_healthy @@ -160,8 +180,9 @@ services: - APP_DATABASE_URL=postgres://opentel_user:opentel_pass@postgres:5432/opentel_db - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics + - OTEL_EXPORTER_OTLP_LOGS_ENDPOINT=http://loki:3100/otlp/v1/logs - OTEL_SERVICE_NAME=orders - - RUST_LOG=info,otel::tracing=trace + - RUST_LOG=info,open_telemetry=off,otel::tracing=trace depends_on: postgres: condition: service_healthy @@ -222,8 +243,9 @@ services: - APP_SERVICES__ORDERS_SERVICE_URL=http://orders:3003 - OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://tempo:4317 - OTEL_EXPORTER_OTLP_METRICS_ENDPOINT=http://prometheus:9090/api/v1/otlp/v1/metrics + - OTEL_EXPORTER_OTLP_LOGS_ENDPOINT=http://loki:3100/otlp/v1/logs - OTEL_SERVICE_NAME=otelmart - - RUST_LOG=info,otel::tracing=trace + - RUST_LOG=info,open_telemetry=off,otel::tracing=trace extra_hosts: - "host.docker.internal:host-gateway" depends_on: @@ -262,6 +284,7 @@ volumes: name: opentel_tempo_data prometheus_data: name: opentel_prometheus_data + loki_data: + name: opentel_loki_data grafana_data: name: opentel_grafana_data - diff --git a/distributed_observability/grafana/provisioning/dashboards/otelmart.json b/distributed_observability/grafana/provisioning/dashboards/otelmart.json @@ -160,7 +160,7 @@ "spanNulls": false, "stacking": { "group": "A", - "mode": "none" + "mode": "normal" }, "thresholdsStyle": { "mode": "off" @@ -187,7 +187,7 @@ "x": 12, "y": 0 }, - "id": 1, + "id": 6, "options": { "annotations": { "clustering": -1, @@ -215,13 +215,13 @@ "uid": "prometheus" }, "editorMode": "code", - "expr": "sum(rate(inventory_reservation_attempts_total[5m])) by (job)", + "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": "Inventory Reservation Rate", + "title": "Error Rate: Broken Registers", "type": "timeseries" }, { @@ -334,8 +334,7 @@ "fieldConfig": { "defaults": { "color": { - "mode": "palette-classic", - "seriesBy": "last" + "mode": "palette-classic" }, "custom": { "axisBorderShow": false, @@ -346,8 +345,8 @@ "barAlignment": 0, "barWidthFactor": 0.6, "drawStyle": "line", - "fillOpacity": 27, - "gradientMode": "opacity", + "fillOpacity": 25, + "gradientMode": "none", "hideFrom": { "legend": false, "tooltip": false, @@ -355,11 +354,8 @@ }, "insertNulls": false, "lineInterpolation": "linear", - "lineStyle": { - "fill": "solid" - }, "lineWidth": 1, - "pointSize": 4, + "pointSize": 5, "scaleDistribution": { "type": "linear" }, @@ -375,17 +371,13 @@ } }, "thresholds": { - "mode": "percentage", + "mode": "absolute", "steps": [ { "color": "green", "value": 0 }, { - "color": "#EAB839", - "value": 60 - }, - { "color": "red", "value": 80 } @@ -399,7 +391,7 @@ "x": 12, "y": 8 }, - "id": 5, + "id": 1, "options": { "annotations": { "clustering": -1, @@ -427,13 +419,13 @@ "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", + "expr": "sum(rate(inventory_reservation_attempts_total[5m])) by (job)", "legendFormat": "__auto", "range": true, "refId": "A" } ], - "title": "Business Metrics: Checkout Success Rate", + "title": "Inventory Reservation Rate", "type": "timeseries" }, { @@ -546,7 +538,8 @@ "fieldConfig": { "defaults": { "color": { - "mode": "palette-classic" + "mode": "palette-classic", + "seriesBy": "last" }, "custom": { "axisBorderShow": false, @@ -557,8 +550,8 @@ "barAlignment": 0, "barWidthFactor": 0.6, "drawStyle": "line", - "fillOpacity": 0, - "gradientMode": "none", + "fillOpacity": 27, + "gradientMode": "opacity", "hideFrom": { "legend": false, "tooltip": false, @@ -566,8 +559,11 @@ }, "insertNulls": false, "lineInterpolation": "linear", + "lineStyle": { + "fill": "solid" + }, "lineWidth": 1, - "pointSize": 5, + "pointSize": 4, "scaleDistribution": { "type": "linear" }, @@ -583,13 +579,17 @@ } }, "thresholds": { - "mode": "absolute", + "mode": "percentage", "steps": [ { "color": "green", "value": 0 }, { + "color": "#EAB839", + "value": 60 + }, + { "color": "red", "value": 80 } @@ -603,7 +603,7 @@ "x": 12, "y": 16 }, - "id": 4, + "id": 5, "options": { "annotations": { "clustering": -1, @@ -631,13 +631,13 @@ "uid": "prometheus" }, "editorMode": "code", - "expr": "db_pool_connections_idle{job=~\".+\"}", + "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": "Idle Connections (Pool Headroom)", + "title": "Business Metrics: Checkout Success Rate", "type": "timeseries" }, { @@ -761,7 +761,7 @@ "barAlignment": 0, "barWidthFactor": 0.6, "drawStyle": "line", - "fillOpacity": 25, + "fillOpacity": 0, "gradientMode": "none", "hideFrom": { "legend": false, @@ -780,7 +780,7 @@ "spanNulls": false, "stacking": { "group": "A", - "mode": "normal" + "mode": "none" }, "thresholdsStyle": { "mode": "off" @@ -804,10 +804,10 @@ "gridPos": { "h": 8, "w": 12, - "x": 0, - "y": 32 + "x": 12, + "y": 24 }, - "id": 6, + "id": 4, "options": { "annotations": { "clustering": -1, @@ -835,13 +835,13 @@ "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)", + "expr": "db_pool_connections_idle{job=~\".+\"}", "legendFormat": "__auto", "range": true, "refId": "A" } ], - "title": "Error Rate: Broken Registers", + "title": "Idle Connections (Pool Headroom)", "type": "timeseries" } ], @@ -869,5 +869,5 @@ "timezone": "browser", "title": "Otelmart Dashboard", "uid": "ad79bw7", - "version": 15 + "version": 17 } diff --git a/distributed_observability/grafana/provisioning/datasources/datasources.yml b/distributed_observability/grafana/provisioning/datasources/datasources.yml @@ -1,7 +1,10 @@ -# Grafana datasources, provisioned as code so they survive container recreates. +# Grafana datasource provisioning +# Auto-configures Prometheus, Tempo, and Loki datasources on startup + apiVersion: 1 datasources: + # Prometheus — metrics backend - name: prometheus type: prometheus uid: prometheus @@ -14,16 +17,41 @@ datasources: - name: trace_id datasourceUid: tempo + # Loki — log aggregation backend (provisioned before Jaeger so tracesToLogsV2 can reference its UID) + - name: Loki + type: loki + access: proxy + url: http://loki:3100 + uid: loki + editable: true + jsonData: + derivedFields: + # Logs → Traces: turn the `trace_id` structured-metadata field on every + # log record into a clickable link that opens the matching trace in + # Jaeger. Without this, the trace_id column in Loki Explore is plain + # text and the chapter's "click trace_id to jump to Jaeger" workflow + # silently doesn't work. + - name: TraceID + matcherType: label + matcherRegex: trace_id + datasourceUid: jaeger + url: '$${__value.raw}' + urlDisplayLabel: 'View Trace' + + # Tempo — distributed tracing backend (after Loki so its UID is available) - name: tempo type: tempo uid: tempo access: proxy url: http://tempo:3200 + editable: true jsonData: - # Link a trace's spans to the matching Prometheus metrics. - tracesToMetrics: - datasourceUid: prometheus - serviceMap: - datasourceUid: prometheus - nodeGraph: - enabled: true + tracesToLogsV2: + # Link from trace spans to Loki logs using trace_id structured metadata + datasourceUid: loki + spanStartTimeShift: "-1m" + spanEndTimeShift: "1m" + filterByTraceID: false + filterBySpanID: false + customQuery: true + query: '{service_name=~".+"} | trace_id = `$${__span.traceId}`' diff --git a/distributed_observability/inventory/Cargo.toml b/distributed_observability/inventory/Cargo.toml @@ -55,3 +55,8 @@ axum-tracing-opentelemetry = { workspace = true } reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } axum-otel-metrics = { workspace = true } +opentelemetry-appender-tracing = { workspace = true } + +# Hash utils +hex = { workspace = true } +sha2 = { workspace = true } diff --git a/distributed_observability/inventory/src/handlers/inventory.rs b/distributed_observability/inventory/src/handlers/inventory.rs @@ -11,6 +11,7 @@ use axum::{ use opentelemetry::KeyValue; use sqlx::PgPool; use std::time::Instant; +use tracing::{error, info, instrument, warn}; use uuid::Uuid; use crate::db; @@ -121,6 +122,17 @@ pub async fn update_stock( /// /// # Endpoint /// `POST /inventory/reserve` +#[instrument( + name = "reserve_stock", + skip(pool), + fields( + otel.kind = "client", + db.system.name = "postgresql", + db.namespace = "inventory", + product.uuid = %request.product_uuid, + operation.result = tracing::field::Empty, + ) +)] pub async fn reserve_stock( State(pool): State<PgPool>, Json(request): Json<ReserveStockRequest>, @@ -135,7 +147,12 @@ pub async fn reserve_stock( let duration = start.elapsed().as_secs_f64(); if success { - // Record successful reservation metrics + // Record successful reservation logs and metrics + + tracing::Span::current().record("operation.result", "success"); + // also on the log record + info!(operation.result = "success", "Stock reserved successfully"); + metrics() .reservation_duration .record(duration, &[KeyValue::new("outcome", "success")]); @@ -154,7 +171,14 @@ pub async fn reserve_stock( ) .into_response() } else { - // Insufficient stock — record failure metrics + // Insufficient stock — record failure logs and metrics + tracing::Span::current().record("operation.result", "insufficient_stock"); + warn!( + operation.result = "insufficient_stock", + quantity.requested = request.quantity, + "Insufficient stock" + ); + metrics() .reservation_failures .add(1, &[KeyValue::new("failure.reason", "insufficient_stock")]); @@ -179,7 +203,10 @@ pub async fn reserve_stock( } } Err(e) => { - // Database error — record failure metrics + // Database error — record failure logs and metrics + tracing::Span::current().record("operation.result", "error"); + error!(operation.result = "error", error.r#type = "database", error.message = %e, "Database error during stock reservation"); + let duration = start.elapsed().as_secs_f64(); metrics() .reservation_failures diff --git a/distributed_observability/inventory/src/telemetry.rs b/distributed_observability/inventory/src/telemetry.rs @@ -8,14 +8,15 @@ use opentelemetry::global; use opentelemetry::trace::TracerProvider as _; use opentelemetry::KeyValue; +use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge; use opentelemetry_otlp::WithExportConfig; use opentelemetry_sdk::{ - metrics::SdkMeterProvider, propagation::TraceContextPropagator, trace::SdkTracerProvider, - Resource, + logs::SdkLoggerProvider, metrics::SdkMeterProvider, propagation::TraceContextPropagator, + trace::SdkTracerProvider, Resource, }; -use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; +use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter, Layer}; -/// Holds both the tracer and meter providers, ensuring they are +/// Holds the tracer, meter and logger providers, ensuring they are /// shut down gracefully when the application exits. Dropping the guard /// flushes pending telemetry and shuts down each provider in turn, so /// callers only need to keep the value alive for the lifetime of the @@ -23,17 +24,25 @@ use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilte pub struct TelemetryGuard { tracer_provider: SdkTracerProvider, meter_provider: SdkMeterProvider, + logger_provider: SdkLoggerProvider, } impl Drop for TelemetryGuard { fn drop(&mut self) { tracing::info!("Shutting down telemetry..."); + + // Logger first — prevents new log records from being generated + // during the shutdown of other providers. + if let Err(e) = self.logger_provider.shutdown() { + eprintln!("Error shutting down logger provider: {:?}", e); + } if let Err(e) = self.meter_provider.shutdown() { eprintln!("Error shutting down meter provider: {:?}", e); } if let Err(e) = self.tracer_provider.shutdown() { eprintln!("Error shutting down tracer provider: {:?}", e); } + tracing::info!("Telemetry shutdown complete"); } } @@ -65,31 +74,76 @@ pub fn init_telemetry(service_name: &str) -> TelemetryGuard { .build(); // Resolve separate endpoints for traces (gRPC) and metrics (HTTP) + // Traces → Jaeger (gRPC on port 4317) 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:9090/api/v1/otlp/v1/metrics".into()); - + // Metrics → Prometheus 3.0 (HTTP on port 9090) 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:4317".into()); + // Logs → Loki (HTTP on port 3100) + let logs_endpoint = std::env::var("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT") + .unwrap_or_else(|_| "http://localhost:3100/otlp/v1/logs".into()); - // Initialize both providers - let tracer_provider = init_tracer_provider(service_name, &traces_endpoint, &resource); + // Initialize all three providers + let tracer_provider = init_tracer_provider(&traces_endpoint, &resource); let meter_provider = init_meter_provider(&metrics_endpoint, &resource); + let logger_provider = init_logger_provider(&logs_endpoint, &resource); + + // Build the subscriber layers + // OpenTelemetry traces layer — converts tracing spans to OTel spans for Jaeger + let tracer = tracer_provider.tracer(service_name.to_string()); + let otel_trace_layer = tracing_opentelemetry::layer().with_tracer(tracer); + // OpenTelemetry logs layer — bridges tracing events to OTel log records for Loki + let otel_logs_layer = OpenTelemetryTracingBridge::new(&logger_provider); + + // Environment filter for log levels + // Include opentelemetry=off to prevent the SDK's internal diagnostics + // from feeding back into the OpenTelemetryTracingBridge (infinite recursion) + let env_filter = EnvFilter::try_from_default_env() + .unwrap_or_else(|_| EnvFilter::new("info,opentelemetry=off")); + + // Console output layer: JSON in production, compact in development + let is_production = std::env::var("ENVIRONMENT") + .map(|e| e == "production") + .unwrap_or(false); + + let fmt_layer = if is_production { + tracing_subscriber::fmt::layer() + .json() + .with_target(true) + .with_thread_ids(false) + .with_file(false) + .with_line_number(false) + .boxed() + } else { + tracing_subscriber::fmt::layer() + .with_target(true) + .with_thread_ids(false) + .compact() + .boxed() + }; + + // Combine all layers into a subscriber and set it as the global default + // Layer ordering: env_filter at the top filters all downstream layers + tracing_subscriber::registry() + .with(env_filter) + .with(fmt_layer) + .with(otel_trace_layer) + .with(otel_logs_layer) + .init(); TelemetryGuard { tracer_provider, meter_provider, + logger_provider, } } /// Creates the OTLP gRPC trace exporter and tracer provider, /// then wires it into the global tracing subscriber. -fn init_tracer_provider( - service_name: &str, - endpoint: &str, - resource: &Resource, -) -> SdkTracerProvider { +fn init_tracer_provider(endpoint: &str, resource: &Resource) -> SdkTracerProvider { // Create the OTLP exporter configured to send traces via gRPC let exporter = opentelemetry_otlp::SpanExporter::builder() .with_tonic() @@ -103,30 +157,9 @@ fn init_tracer_provider( .with_resource(resource.clone()) .build(); - // Create a tracer from the provider for the OpenTelemetry layer - let tracer = provider.tracer(service_name.to_string()); - - // Create the OpenTelemetry tracing layer - let otel_layer = tracing_opentelemetry::layer().with_tracer(tracer); - - // Create an environment filter for log levels (defaults to "info") - // Include otel::tracing=trace to allow axum-tracing-opentelemetry spans through - let env_filter = EnvFilter::try_from_default_env() - .unwrap_or_else(|_| EnvFilter::new("info,otel::tracing=trace")); - - // Create a formatting layer for console output - let fmt_layer = tracing_subscriber::fmt::layer() - .with_target(true) - .with_thread_ids(false) - .compact(); - - // Combine all layers into a subscriber and set it as the global default - tracing_subscriber::registry() - .with(env_filter) - .with(fmt_layer) - .with(otel_layer) - .init(); - + // NOTE: the global subscriber is assembled once in `init_telemetry`, which + // combines the trace layer with the logs bridge. Setting it here too would + // panic with "a global default trace dispatcher has already been set". provider } @@ -151,6 +184,21 @@ fn init_meter_provider(endpoint: &str, resource: &Resource) -> SdkMeterProvider provider } +/// Initialize the logger provider with OTLP HTTP exporter. +/// Exports log records to Loki via its native OTLP endpoint. +fn init_logger_provider(endpoint: &str, resource: &Resource) -> SdkLoggerProvider { + let exporter = opentelemetry_otlp::LogExporter::builder() + .with_http() + .with_endpoint(endpoint) + .build() + .expect("Failed to create OTLP log exporter"); + + SdkLoggerProvider::builder() + .with_resource(resource.clone()) + .with_batch_exporter(exporter) + .build() +} + /// Registers observable gauges that track database connection pool health. /// /// Creates three instruments: diff --git a/distributed_observability/loki-config.yaml b/distributed_observability/loki-config.yaml @@ -0,0 +1,40 @@ +# Loki configuration for Chapter 7 — Log aggregation via OTLP +# Receives logs from OpenTelemetry SDK via HTTP (OTLP/HTTP) on port 3100 + +auth_enabled: false + +server: + http_listen_port: 3100 + grpc_listen_port: 9096 + +common: + instance_addr: 127.0.0.1 + path_prefix: /loki + storage: + filesystem: + chunks_directory: /loki/chunks + rules_directory: /loki/rules + replication_factor: 1 + ring: + kvstore: + store: inmemory + +schema_config: + configs: + - from: "2024-01-01" + store: tsdb + object_store: filesystem + schema: v13 + index: + prefix: index_ + period: 24h + +limits_config: + allow_structured_metadata: true + volume_enabled: true + otlp_config: + resource_attributes: + attributes_config: + - action: index_label + attributes: + - service.name diff --git a/distributed_observability/orders/Cargo.toml b/distributed_observability/orders/Cargo.toml @@ -58,3 +58,8 @@ axum-tracing-opentelemetry = { workspace = true } reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } axum-otel-metrics = { workspace = true } +opentelemetry-appender-tracing = { workspace = true } + +# Hash utils +hex = { workspace = true } +sha2 = { workspace = true } diff --git a/distributed_observability/orders/src/handlers/orders.rs b/distributed_observability/orders/src/handlers/orders.rs @@ -10,9 +10,10 @@ use axum::{ use bigdecimal::BigDecimal; use opentelemetry::KeyValue; use sqlx::PgPool; -use tracing::{info, warn}; +use tracing::{debug, error, info, instrument, warn}; use uuid::Uuid; +use crate::logging::hash_email; use crate::models::{CreateOrderRequest, OrderQueryParams, OrdersResponse}; use crate::AppState; use crate::{ @@ -82,24 +83,29 @@ async fn release_stock(state: &AppState, product_uuid: Uuid, quantity: i32) { match response { Ok(resp) => { if resp.status().is_success() { - info!( - product_uuid = %product_uuid, + debug!( + product.uuid = %product_uuid, quantity = quantity, - "Successfully released stock" + "Stock released successfully" ); } else { - warn!( - product_uuid = %product_uuid, - status = %resp.status(), - "Failed to release stock" + // This is serious — we have leaked inventory + error!( + product.uuid = %product_uuid, + quantity = quantity, + error.r#type = "compensation_failure", + downstream.status = %resp.status(), + "CRITICAL: Failed to release reserved stock" ); } } Err(e) => { - warn!( - product_uuid = %product_uuid, - error = %e, - "Error calling inventory service to release stock" + error!( + product.uuid = %product_uuid, + quantity = quantity, + error.r#type = "compensation_failure", + error.message = %e, + "CRITICAL: Failed to release reserved stock" ); } } @@ -114,14 +120,16 @@ async fn release_all_reserved_stock(state: &AppState, reserved_items: &[(Uuid, i return; } - warn!( + info!( item_count = reserved_items.len(), - "Rolling back stock reservations" + "Releasing reserved stock" ); for (product_uuid, quantity) in reserved_items { release_stock(state, *product_uuid, *quantity).await; } + + info!(item_count = reserved_items.len(), "Stock release complete"); } /// Confirm a sale with the inventory service @@ -391,6 +399,15 @@ impl From<sqlx::Error> for OrderError { /// /// Uses the `with_transaction` wrapper to ensure all database operations /// are atomic and instrumented with OpenTelemetry semantic conventions. +#[instrument( + name = "create_order", + skip(state, request), + fields( + customer.email_hash = %hash_email(&request.customer_email), + order.item_count = request.items.len(), + order.number = tracing::field::Empty, + ) +)] pub async fn create_order( State(state): State<AppState>, Json(request): Json<CreateOrderRequest>, @@ -495,6 +512,10 @@ pub async fn create_order( Box::pin(async move { // Generate order number let order_number = db::generate_order_number(tx).await?; + // This fills in the order.number which was empty, tracing::field::Empty, + // at instrumentation macro + tracing::Span::current().record("order.number", &order_number); + // From this point, every Log in the span includes order.number // Create order record let created_order = db::create_order( @@ -549,6 +570,15 @@ pub async fn create_order( // Now confirm the sale with inventory service to convert reserved stock to sold confirm_all_sales(&state, created_order.uuid, &reserved_items).await; + tracing::info!( + order.number = %created_order.order_number, + order.uuid = %created_order.uuid, + order.total = %total, + order.status = "processing", + duration_ms = (duration * 1000.0) as u64, + "Order created successfully" + ); + ( StatusCode::CREATED, Json(serde_json::json!({ @@ -563,14 +593,25 @@ pub async fn create_order( Err(e) => { // Transaction failed (rolled back automatically) + tracing::warn!( + reserved_count = reserved_items.len(), + "Order creation failed, releasing reserved stock" + ); // Record failure record_checkout_failure("order_creation", &start); // Release reserved stock release_all_reserved_stock(&state, &reserved_items).await; + let duration = start.elapsed().as_secs_f64(); + let error_msg = match e { OrderError::Database(db_err) => { - tracing::error!(error = %db_err, "Database error during order creation"); + tracing::error!( + error.r#type = "database", + error.message = %db_err, + duration_ms = (duration * 1000.0) as u64, + "Database error during order creation" + ); format!("Database error: {}", db_err) } }; diff --git a/distributed_observability/orders/src/logging.rs b/distributed_observability/orders/src/logging.rs @@ -0,0 +1,26 @@ +//! PII-safe logging utilities. +//! +//! Provides hashing functions for sensitive data (emails, IDs) so that +//! log entries preserve correlation without exposing raw PII. + +use sha2::{Digest, Sha256}; + +/// Hash an email for logging (preserves privacy while enabling correlation). +/// The same email always produces the same hash, so you can correlate +/// log entries for a customer without storing their raw email. +pub fn hash_email(email: &str) -> String { + let mut hasher = Sha256::new(); + hasher.update(email.as_bytes()); + let result = hasher.finalize(); + // Truncate to 8 bytes (64 bits) for readability + format!("sha256:{}", hex::encode(&result[..8])) +} + +/// Hash a numeric ID for logging. +#[allow(dead_code)] +pub fn hash_id(id: i32) -> String { + let mut hasher = Sha256::new(); + hasher.update(id.to_le_bytes()); + let result = hasher.finalize(); + format!("sha256:{}", hex::encode(&result[..8])) +} 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 logging; mod metrics; mod models; mod telemetry; diff --git a/distributed_observability/otelmart/Cargo.toml b/distributed_observability/otelmart/Cargo.toml @@ -58,6 +58,7 @@ tracing-subscriber = { workspace = true } # OpenTelemetry core opentelemetry = { workspace = true } opentelemetry_sdk = { workspace = true } +opentelemetry-appender-tracing = { workspace = true } # OTLP exporter (sends traces to collectors/backends) opentelemetry-otlp = { workspace = true } @@ -71,4 +72,8 @@ reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } axum-otel-metrics = { workspace = true } +# Hash utils +hex = { workspace = true } +sha2 = { workspace = true } + [dev-dependencies] diff --git a/distributed_observability/otelmart/src/proxy/mod.rs b/distributed_observability/otelmart/src/proxy/mod.rs @@ -91,7 +91,7 @@ pub async fn proxy_request( // Handle request failure let response = match response { Ok(resp) => resp, - Err(_e) => { + Err(e) => { // Record upstream request failure metric metrics().upstream_request_duration.record( duration, @@ -102,6 +102,16 @@ pub async fn proxy_request( ], ); + // Log the downstream service failure with structured fields + tracing::error!( + error.r#type = "http_client", + error.message = %e, + downstream.service = %upstream_service, + downstream.url = %target_url, + downstream.duration_ms = format!("{:.1}", duration * 1000.0), + "Downstream service unavailable" + ); + return Response::builder() .status(StatusCode::SERVICE_UNAVAILABLE) .header("content-type", "application/json") @@ -118,12 +128,20 @@ pub async fn proxy_request( metrics().upstream_request_duration.record( duration, &[ - KeyValue::new("upstream.service", upstream_service), + KeyValue::new("upstream.service", upstream_service.clone()), KeyValue::new("http.request.method", method_str), KeyValue::new("http.response.status_code", i64::from(status.as_u16())), ], ); + // Log completed downstream call with response status + tracing::info!( + downstream.service = %upstream_service, + downstream.status = status.as_u16(), + downstream.duration_ms = format!("{:.1}", duration * 1000.0), + "Downstream call completed" + ); + // Get response body (this consumes the response) let body = match response.bytes().await { Ok(bytes) => bytes, diff --git a/distributed_observability/products/Cargo.toml b/distributed_observability/products/Cargo.toml @@ -55,3 +55,8 @@ axum-tracing-opentelemetry = { workspace = true } reqwest-middleware = { workspace = true } reqwest-tracing = { workspace = true } axum-otel-metrics = { workspace = true } +opentelemetry-appender-tracing = { workspace = true } + +# Hash utils +hex = { workspace = true } +sha2 = { workspace = true } diff --git a/distributed_observability/products/src/handlers/products.rs b/distributed_observability/products/src/handlers/products.rs @@ -9,6 +9,7 @@ use axum::{ response::{IntoResponse, Json}, }; use sqlx::PgPool; +use tracing::instrument; use uuid::Uuid; use crate::db; @@ -69,6 +70,14 @@ pub async fn list_products( /// # Errors /// - `404 NOT FOUND` - Product not found or has been deleted /// - `500 INTERNAL SERVER ERROR` - Database error +#[instrument( + name = "get_product_by_id", + skip(pool), + fields( + product.uuid = %uuid, + product.found = tracing::field::Empty + ) +)] pub async fn get_product_by_id( State(pool): State<PgPool>, Path(uuid): Path<Uuid>, @@ -76,14 +85,22 @@ pub async fn get_product_by_id( // Delegate to repository layer for database query match db::get_product_by_uuid(&pool, uuid).await { Ok(Some(mut product)) => { + tracing::Span::current().record("product.found", true); + tracing::debug!(product.uuid = true, product.name = %product.product_name, "Product found"); + // Set the string representation of the product ID product.set_product_id(); (StatusCode::OK, Json(product)).into_response() } - Ok(None) => not_found_error( - "Product not found", - serde_json::json!({"uuid": uuid.to_string()}), - ), + Ok(None) => { + tracing::Span::current().record("product.found", false); + tracing::info!(product.found = false, "Product not found"); // INFO not ERROR + + not_found_error( + "Product not found", + serde_json::json!({"uuid": uuid.to_string()}), + ) + } Err(e) => internal_error("Failed to fetch product", e.to_string()), } }