diff --git a/.env.example b/.env.example index 1f4afa7..414ab1d 100644 --- a/.env.example +++ b/.env.example @@ -53,3 +53,9 @@ METADATA_PORT=8002 UPLOAD_PORT=8003 NGINX_VOD_PORT=8081 FRONTEND_PORT=3000 +PROMETHEUS_PORT=9090 +GRAFANA_PORT=3001 +GRAFANA_USER=admin +GRAFANA_PASSWORD=admin +TRANSCODER_METRICS_PORT=8004 +LOKI_PORT=3100 diff --git a/deploy/docker-compose.yml b/deploy/docker-compose.yml index fa733fd..ef842f9 100644 --- a/deploy/docker-compose.yml +++ b/deploy/docker-compose.yml @@ -198,6 +198,9 @@ services: ENCODE_MODE: ${ENCODE_MODE:-cbr} FFMPEG_HWACCEL: ${FFMPEG_HWACCEL:-auto} TRANSCODER_PIPE_INPUT: ${TRANSCODER_PIPE_INPUT:-true} + METRICS_PORT: ${TRANSCODER_METRICS_PORT:-8004} + ports: + - "${TRANSCODER_METRICS_PORT:-8004}:8004" depends_on: rabbitmq: condition: service_healthy @@ -257,7 +260,81 @@ services: - gateway restart: unless-stopped + prometheus: + image: prom/prometheus:v2.55.1 + command: + - --config.file=/etc/prometheus/prometheus.yml + - --storage.tsdb.path=/prometheus + - --web.console.libraries=/usr/share/prometheus/console_libraries + - --web.console.templates=/usr/share/prometheus/consoles + volumes: + - ./prometheus/prometheus.yml:/etc/prometheus/prometheus.yml:ro + - ./prometheus/alert.rules.yml:/etc/prometheus/alert.rules.yml:ro + - promdata:/prometheus + ports: + - "${PROMETHEUS_PORT:-9090}:9090" + restart: unless-stopped + healthcheck: + test: ["CMD", "wget", "-qO-", "http://localhost:9090/-/healthy"] + interval: 15s + timeout: 5s + retries: 5 + + grafana: + image: grafana/grafana:11.4.0 + environment: + GF_SECURITY_ADMIN_USER: ${GRAFANA_USER:-admin} + GF_SECURITY_ADMIN_PASSWORD: ${GRAFANA_PASSWORD:-admin} + GF_SERVER_HTTP_PORT: 3000 + volumes: + - grafanadata:/var/lib/grafana + - ./grafana/provisioning:/etc/grafana/provisioning:ro + - ./grafana/dashboards:/etc/grafana/dashboards:ro + ports: + - "${GRAFANA_PORT:-3001}:3000" + depends_on: + prometheus: + condition: service_healthy + loki: + condition: service_healthy + restart: unless-stopped + healthcheck: + test: ["CMD-SHELL", "wget -qO- http://localhost:3000/api/health | grep -q ok || curl -f http://localhost:3000/api/health"] + interval: 15s + timeout: 5s + retries: 10 + + loki: + image: grafana/loki:3.0.0 + command: -config.file=/etc/loki/loki.yml + volumes: + - ./loki/loki.yml:/etc/loki/loki.yml:ro + - lokidata:/loki + ports: + - "${LOKI_PORT:-3100}:3100" + restart: unless-stopped + healthcheck: + test: ["CMD-SHELL", "wget -qO- http://localhost:3100/ready | grep -q ready || curl -f http://localhost:3100/ready"] + interval: 15s + timeout: 5s + retries: 10 + + promtail: + image: grafana/promtail:3.0.0 + command: -config.file=/etc/promtail/promtail.yml + volumes: + - ./promtail/promtail.yml:/etc/promtail/promtail.yml:ro + - /var/lib/docker/containers:/var/lib/docker/containers:ro + - /var/run/docker.sock:/var/run/docker.sock:ro + depends_on: + loki: + condition: service_healthy + restart: unless-stopped + volumes: pgdata: miniodata: rabbitdata: + promdata: + grafanadata: + lokidata: diff --git a/deploy/grafana/dashboards/flowix.json b/deploy/grafana/dashboards/flowix.json new file mode 100644 index 0000000..92c87bc --- /dev/null +++ b/deploy/grafana/dashboards/flowix.json @@ -0,0 +1,45 @@ +{ + "title": "Flowix — Overview", + "uid": "flowix-overview", + "tags": ["flowix"], + "timezone": "browser", + "schemaVersion": 30, + "version": 1, + "panels": [ + { + "type": "stat", + "title": "RabbitMQ queue depth", + "datasource": "Prometheus", + "targets": [{ "expr": "rabbitmq_queue_depth", "instant": true }], + "gridPos": { "h": 4, "w": 6, "x": 0, "y": 0 } + }, + { + "type": "graph", + "title": "FFmpeg duration (histogram count)", + "datasource": "Prometheus", + "targets": [{ "expr": "sum by (quality) (increase(ffmpeg_duration_seconds_count[5m]))", "legendFormat": "{{quality}}" }], + "gridPos": { "h": 8, "w": 12, "x": 6, "y": 0 } + }, + { + "type": "stat", + "title": "Upload bytes (total)", + "datasource": "Prometheus", + "targets": [{ "expr": "sum(upload_bytes)", "instant": true }], + "gridPos": { "h": 4, "w": 6, "x": 0, "y": 4 } + }, + { + "type": "stat", + "title": "VOD cache hits (total)", + "datasource": "Prometheus", + "targets": [{ "expr": "sum(vod_cache_hit)", "instant": true }], + "gridPos": { "h": 4, "w": 6, "x": 6, "y": 4 } + }, + { + "type": "graph", + "title": "HTTP requests by service (from metrics if available)", + "datasource": "Prometheus", + "targets": [{ "expr": "sum by (job) (rate(http_requests_total[5m]))", "legendFormat": "{{job}}" }], + "gridPos": { "h": 8, "w": 12, "x": 0, "y": 8 } + } + ] +} diff --git a/deploy/grafana/provisioning/dashboards/dashboard.yml b/deploy/grafana/provisioning/dashboards/dashboard.yml new file mode 100644 index 0000000..50ca277 --- /dev/null +++ b/deploy/grafana/provisioning/dashboards/dashboard.yml @@ -0,0 +1,10 @@ +apiVersion: 1 +providers: + - name: 'Flowix' + orgId: 1 + folder: '' + type: file + disableDeletion: false + editable: true + options: + path: /etc/grafana/dashboards diff --git a/deploy/grafana/provisioning/datasources/datasource.yml b/deploy/grafana/provisioning/datasources/datasource.yml new file mode 100644 index 0000000..0387fce --- /dev/null +++ b/deploy/grafana/provisioning/datasources/datasource.yml @@ -0,0 +1,13 @@ +apiVersion: 1 +datasources: + - name: Prometheus + type: prometheus + access: proxy + url: http://prometheus:9090 + isDefault: true + editable: true + - name: Loki + type: loki + access: proxy + url: http://loki:3100 + editable: true diff --git a/deploy/loki/loki.yml b/deploy/loki/loki.yml new file mode 100644 index 0000000..e47af0a --- /dev/null +++ b/deploy/loki/loki.yml @@ -0,0 +1,55 @@ +auth_enabled: false +server: + http_listen_port: 3100 + grpc_listen_port: 9096 + +common: + path_prefix: /loki + storage: + filesystem: + chunks_directory: /loki/chunks + rules_directory: /loki/rules + replication_factor: 1 + ring: + kvstore: + store: inmemory + +ingester: + lifecycler: + ring: + kvstore: + store: inmemory + replication_factor: 1 + final_sleep: 0s + chunk_idle_period: 5m + chunk_retain_period: 30s + wal: + enabled: false + +schema_config: + configs: + - from: 2020-05-15 + store: boltdb-shipper + object_store: filesystem + schema: v11 + index: + prefix: index_ + period: 24h + +storage_config: + boltdb_shipper: + active_index_directory: /loki/index + cache_location: /loki/cache + cache_ttl: 24h + filesystem: + directory: /loki/chunks + +limits_config: + allow_structured_metadata: false + volume_enabled: true + retention_period: 744h + reject_old_samples: true + reject_old_samples_max_age: 168h + +analytics: + reporting_enabled: false diff --git a/deploy/prometheus/alert.rules.yml b/deploy/prometheus/alert.rules.yml new file mode 100644 index 0000000..4ad4f3d --- /dev/null +++ b/deploy/prometheus/alert.rules.yml @@ -0,0 +1,25 @@ +groups: + - name: flowix + interval: 30s + rules: + - alert: QueueDepthHigh + expr: rabbitmq_queue_messages > 100 + for: 2m + labels: + severity: warning + annotations: + summary: "RabbitMQ queue depth >100" + - alert: TranscoderOOM + expr: increase(container_memory_failures_total{name=~"transcoder.*"}[5m]) > 0 + for: 1m + labels: + severity: critical + annotations: + summary: "Transcoder OOM detected" + - alert: MinIODiskHigh + expr: (minio_cluster_capacity_usable_free_bytes / minio_cluster_capacity_usable_total_bytes) < 0.2 + for: 5m + labels: + severity: warning + annotations: + summary: "MinIO disk >80% used" diff --git a/deploy/prometheus/prometheus.yml b/deploy/prometheus/prometheus.yml new file mode 100644 index 0000000..2c39f6c --- /dev/null +++ b/deploy/prometheus/prometheus.yml @@ -0,0 +1,45 @@ +global: + scrape_interval: 15s + evaluation_interval: 15s + +scrape_configs: + - job_name: 'prometheus' + static_configs: + - targets: ['localhost:9090'] + + - job_name: 'gateway' + static_configs: + - targets: ['gateway:8080'] + metrics_path: /metrics + + - job_name: 'metadata' + static_configs: + - targets: ['metadata:8002'] + metrics_path: /metrics + + - job_name: 'upload' + static_configs: + - targets: ['upload:8003'] + metrics_path: /metrics + + - job_name: 'auth' + static_configs: + - targets: ['auth:8001'] + metrics_path: /metrics + + - job_name: 'transcoder' + static_configs: + - targets: ['transcoder:8004'] + metrics_path: /metrics + + - job_name: 'rabbitmq' + static_configs: + - targets: ['rabbitmq:15672'] + + - job_name: 'minio' + static_configs: + - targets: ['minio:9000'] + metrics_path: /minio/v2/metrics/cluster + +rule_files: + - /etc/prometheus/alert.rules.yml diff --git a/deploy/promtail/promtail.yml b/deploy/promtail/promtail.yml new file mode 100644 index 0000000..4042a15 --- /dev/null +++ b/deploy/promtail/promtail.yml @@ -0,0 +1,37 @@ +server: + http_listen_port: 9080 + grpc_listen_port: 0 + +positions: + filename: /tmp/positions.yaml + +clients: + - url: http://loki:3100/loki/api/v1/push + +scrape_configs: + - job_name: docker + docker_sd_configs: + - host: unix:///var/run/docker.sock + refresh_interval: 5s + relabel_configs: + - source_labels: ['__meta_docker_container_name'] + regex: '/(.*)' + target_label: 'container' + - source_labels: ['__meta_docker_container_label_com_docker_compose_service'] + target_label: 'service' + - source_labels: ['__meta_docker_container_label_com_docker_compose_project'] + target_label: 'project' + - source_labels: ['__meta_docker_container_log_stream'] + target_label: 'stream' + pipeline_stages: + - json: + expressions: + level: level + msg: msg + service: service + trace_id: trace_id + request_id: request_id + - labels: + level: + service: + trace_id: diff --git a/services/auth/pyproject.toml b/services/auth/pyproject.toml index fa67cc6..5792af8 100644 --- a/services/auth/pyproject.toml +++ b/services/auth/pyproject.toml @@ -17,6 +17,7 @@ dependencies = [ "pydantic-settings>=2.6,<3", "python-multipart>=0.0.9", "httpx>=0.27,<0.28", + "prometheus_client>=0.21", ] [dependency-groups] diff --git a/services/auth/src/main.py b/services/auth/src/main.py index 3583fe7..93d3578 100644 --- a/services/auth/src/main.py +++ b/services/auth/src/main.py @@ -1,7 +1,34 @@ -from fastapi import FastAPI +import json +import logging +import sys +import time +import uuid + +from fastapi import FastAPI, Request +from fastapi.responses import Response from .routers.auth import router +# Unified JSON logging with trace_id (Phase 14) +_auth_logger = logging.getLogger("auth") +if not _auth_logger.handlers: + _h = logging.StreamHandler(sys.stdout) + _h.setFormatter(logging.Formatter("%(message)s")) + _auth_logger.addHandler(_h) +_auth_logger.setLevel(logging.INFO) + +try: + from prometheus_client import CONTENT_TYPE_LATEST, Counter, Gauge, Histogram, generate_latest + + auth_requests = Counter("auth_requests_total", "Auth requests", ["method", "endpoint"]) + rabbitmq_queue_depth = Gauge("rabbitmq_queue_depth", "RabbitMQ queue depth") + upload_bytes = Counter("upload_bytes", "Upload bytes") + vod_cache_hit = Counter("vod_cache_hit", "VOD cache hits") + ffmpeg_duration = Histogram("ffmpeg_duration_seconds", "FFmpeg duration", ["quality"]) + METRICS_ENABLED = True +except ImportError: + METRICS_ENABLED = False + app = FastAPI(title="flowix-auth", version="0.1.0") # CORS handled by gateway (single entry point) to avoid duplicate # Access-Control-Allow-Origin headers (gateway sets origin, upstream must not). @@ -9,11 +36,60 @@ app.include_router(router) +@app.middleware("http") +async def trace_middleware(request: Request, call_next): + trace_id = ( + request.headers.get("x-request-id") + or request.headers.get("X-Request-ID") + or str(uuid.uuid4()) + ) + request.state.trace_id = trace_id + start = time.time() + response = await call_next(request) + duration = time.time() - start + level = ( + "info" if response.status_code < 400 else "warn" if response.status_code < 500 else "error" + ) + log_data = { + "level": level, + "msg": "request", + "service": "auth", + "method": request.method, + "path": request.url.path, + "status": response.status_code, + "duration": duration, + "trace_id": trace_id, + "request_id": trace_id, + } + try: + _auth_logger.info(json.dumps(log_data)) + except Exception: + pass + response.headers["X-Request-ID"] = trace_id + return response + + @app.get("/health") def health(): return {"status": "ok", "service": "auth"} +if METRICS_ENABLED: + + @app.get("/metrics") + def metrics(): + return Response(generate_latest(), media_type=CONTENT_TYPE_LATEST) + + @app.middleware("http") + async def metrics_middleware(request: Request, call_next): + response = await call_next(request) + try: + auth_requests.labels(method=request.method, endpoint=request.url.path).inc() + except Exception: + pass + return response + + @app.get("/") def root(): return {"service": "auth", "docs": "/docs"} diff --git a/services/auth/uv.lock b/services/auth/uv.lock index 68f93b5..d348cd9 100644 --- a/services/auth/uv.lock +++ b/services/auth/uv.lock @@ -399,6 +399,7 @@ dependencies = [ { name = "fastapi" }, { name = "httpx" }, { name = "passlib", extra = ["argon2"] }, + { name = "prometheus-client" }, { name = "psycopg", extra = ["binary"] }, { name = "pydantic", extra = ["email"] }, { name = "pydantic-settings" }, @@ -426,6 +427,7 @@ requires-dist = [ { name = "fastapi", specifier = ">=0.115,<0.116" }, { name = "httpx", specifier = ">=0.27,<0.28" }, { name = "passlib", extras = ["argon2"], specifier = ">=1.7,<2" }, + { name = "prometheus-client", specifier = ">=0.21" }, { name = "psycopg", extras = ["binary"], specifier = ">=3.1,<4" }, { name = "pydantic", extras = ["email"], specifier = ">=2.9,<3" }, { name = "pydantic-settings", specifier = ">=2.6,<3" }, @@ -792,6 +794,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/54/20/4d324d65cc6d9205fabedc306948156824eb9f0ee1633355a8f7ec5c66bf/pluggy-1.6.0-py3-none-any.whl", hash = "sha256:e920276dd6813095e9377c0bc5566d94c932c33b27a3e3945d8389c374dd4746", size = 20538, upload-time = "2025-05-15T12:30:06.134Z" }, ] +[[package]] +name = "prometheus-client" +version = "0.26.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/52/73/f1334c29c2af4cd9dba6c7817e61b611bd0215e2eb5565c6064a4de18802/prometheus_client-0.26.0.tar.gz", hash = "sha256:04a91bcf94e2cf74a44a1a874d651a2e853ed354b6e822f3b7487751465d5c2b", size = 92910, upload-time = "2026-07-24T19:36:41.893Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/eb/a3/b69efbf4143b5b9859b977770bbbabcc2796b702fa69dc40271e45cd5a56/prometheus_client-0.26.0-py3-none-any.whl", hash = "sha256:fa93d06737aa02bacd05794768508bb97d2fbee28cb3bca04eaae92f0ca953d6", size = 64494, upload-time = "2026-07-24T19:36:40.854Z" }, +] + [[package]] name = "psycopg" version = "3.3.4" diff --git a/services/gateway/cmd/server/main.go b/services/gateway/cmd/server/main.go index 14ae106..e51c3c2 100644 --- a/services/gateway/cmd/server/main.go +++ b/services/gateway/cmd/server/main.go @@ -15,6 +15,7 @@ import ( "flowix/gateway/internal/handler" gwmw "flowix/gateway/internal/middleware" + "flowix/gateway/internal/metrics" "flowix/gateway/internal/proxy" ) @@ -81,6 +82,7 @@ func main() { // health — без прокси, без rate-limit (rate-limit уже пропускает /health) r.Get("/health", healthHandler) r.Get("/healthz", healthHandler) + r.Get("/metrics", metrics.Handler().ServeHTTP) r.Get("/", func(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, map[string]string{"service": "gateway", "status": "ok"}) }) @@ -133,7 +135,7 @@ func main() { // --- HLS / VOD: защищён HLSAuth (private 403 без токена, public пропуск) --- hlsAuth := gwmw.HLSAuth(jwtSecret, internalToken, metadataURL) - r.With(hlsAuth).Handle("/hls/*", vodProxy) + r.With(hlsAuth, metrics.Middleware).Handle("/hls/*", vodProxy) // --- Thumbnails / public MinIO objects via gateway (avoid direct :9000 CORS) --- // frontend uses /thumbnails/{id}/thumb.jpg ; gateway proxies to MinIO bucket `videos` diff --git a/services/gateway/go.mod b/services/gateway/go.mod index 4b6517a..44b61be 100644 --- a/services/gateway/go.mod +++ b/services/gateway/go.mod @@ -5,12 +5,20 @@ go 1.25.0 require ( github.com/go-chi/chi/v5 v5.2.1 github.com/golang-jwt/jwt/v5 v5.3.1 + github.com/prometheus/client_golang v1.22.0 github.com/rs/zerolog v1.35.1 golang.org/x/time v0.15.0 ) require ( + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/mattn/go-colorable v0.1.14 // indirect github.com/mattn/go-isatty v0.0.20 // indirect - golang.org/x/sys v0.29.0 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/prometheus/client_model v0.6.1 // indirect + github.com/prometheus/common v0.62.0 // indirect + github.com/prometheus/procfs v0.15.1 // indirect + golang.org/x/sys v0.30.0 // indirect + google.golang.org/protobuf v1.36.5 // indirect ) diff --git a/services/gateway/go.sum b/services/gateway/go.sum index 0e14d3d..fb306a9 100644 --- a/services/gateway/go.sum +++ b/services/gateway/go.sum @@ -1,15 +1,45 @@ +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/go-chi/chi/v5 v5.2.1 h1:KOIHODQj58PmL80G2Eak4WdvUzjSJSm0vG72crDCqb8= github.com/go-chi/chi/v5 v5.2.1/go.mod h1:L2yAIGWB3H+phAw1NxKwWM+7eUH/lU8pOMm5hHcoops= github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY= github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= +github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= +github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= +github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_golang v1.22.0 h1:rb93p9lokFEsctTys46VnV1kLCDpVZ0a/Y92Vm0Zc6Q= +github.com/prometheus/client_golang v1.22.0/go.mod h1:R7ljNsLXhuQXYZYtw6GAE9AZg8Y7vEW5scdCXrWRXC0= +github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= +github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= +github.com/prometheus/common v0.62.0 h1:xasJaQlnWAeyHdUBeGjXmutelfJHWMRr+Fg4QszZ2Io= +github.com/prometheus/common v0.62.0/go.mod h1:vyBcEuLSvWos9B1+CyL7JZ2up+uFzXhkqml0W5zIY1I= +github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= +github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= github.com/rs/zerolog v1.35.1 h1:m7xQeoiLIiV0BCEY4Hs+j2NG4Gp2o2KPKmhnnLiazKI= github.com/rs/zerolog v1.35.1/go.mod h1:EjML9kdfa/RMA7h/6z6pYmq1ykOuA8/mjWaEvGI+jcw= +github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU= -golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc= +golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= +google.golang.org/protobuf v1.36.5 h1:tPhr+woSbjfYvY6/GPufUoYizxw1cF/yFoxJ2fmpwlM= +google.golang.org/protobuf v1.36.5/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/services/gateway/internal/metrics/metrics.go b/services/gateway/internal/metrics/metrics.go new file mode 100644 index 0000000..458429d --- /dev/null +++ b/services/gateway/internal/metrics/metrics.go @@ -0,0 +1,43 @@ +package metrics + +import ( + "net/http" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +var ( + RabbitMQQueueDepth = prometheus.NewGauge(prometheus.GaugeOpts{ + Name: "rabbitmq_queue_depth", + Help: "Current depth of video.uploaded queue", + }) + VodCacheHit = prometheus.NewCounter(prometheus.CounterOpts{ + Name: "vod_cache_hit", + Help: "Number of VOD cache hits", + }) + UploadBytes = prometheus.NewCounter(prometheus.CounterOpts{ + Name: "upload_bytes", + Help: "Total uploaded bytes via gateway", + }) + FfmpegDuration = prometheus.NewHistogram(prometheus.HistogramOpts{ + Name: "ffmpeg_duration_seconds", + Help: "FFmpeg transcoding duration", + Buckets: prometheus.DefBuckets, + }) +) + +func init() { + prometheus.MustRegister(RabbitMQQueueDepth, VodCacheHit, UploadBytes, FfmpegDuration) +} + +func Handler() http.Handler { + return promhttp.Handler() +} + +func Middleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + VodCacheHit.Inc() + next.ServeHTTP(w, r) + }) +} diff --git a/services/gateway/internal/middleware/logger.go b/services/gateway/internal/middleware/logger.go index 7bc833d..4c5c641 100644 --- a/services/gateway/internal/middleware/logger.go +++ b/services/gateway/internal/middleware/logger.go @@ -23,12 +23,19 @@ func RequestLogger(next http.Handler) http.Handler { } else if ww.status >= 400 { ev = log.Warn() } + reqID := r.Header.Get("X-Request-ID") + if reqID == "" { + reqID = r.Header.Get("X-Request-Id") + } ev.Str("method", r.Method). Str("path", r.URL.Path). Int("status", ww.status). Dur("duration", dur). Str("ip", clientIP(r)). - Str("req_id", r.Header.Get("X-Request-ID")). + Str("req_id", reqID). + Str("trace_id", reqID). + Str("request_id", reqID). + Str("service", "gateway"). Msg("request") // также zerolog global logger доступен как zerolog.Ctx _ = zerolog.Ctx(r.Context()) diff --git a/services/gateway/internal/proxy/proxy.go b/services/gateway/internal/proxy/proxy.go index f90e9e1..c122838 100644 --- a/services/gateway/internal/proxy/proxy.go +++ b/services/gateway/internal/proxy/proxy.go @@ -2,6 +2,7 @@ package proxy import ( + "net" "net/http" "net/http/httputil" "net/url" @@ -21,9 +22,30 @@ func New(target *url.URL) *httputil.ReverseProxy { // Сохраняем оригинальный Host заголовок клиента в X-Forwarded-Host, // а Host выставляем на upstream (важно для nginx-vod Host header). r.Header.Set("X-Forwarded-Host", r.Host) - // Пробрасываем X-Forwarded-For корректно - // (net/http уже делает, но на всякий — добавляем RemoteAddr) - // Не трогаем Authorization / X-User-ID — их ставит Auth middleware + // Прокидываем X-Request-ID (chi middleware.RequestID) во все апстримы для трассировки + if v := r.Header.Get("X-Request-Id"); v != "" { + r.Header.Set("X-Request-ID", v) + } + if v := r.Header.Get("X-Request-ID"); v != "" { + // keep as is, ensure downstream sees it + r.Header.Set("X-Request-ID", v) + } + // X-User-ID уже ставит AuthMiddleware — не трогаем, просто прокидываем если есть + // X-Forwarded-For: добавляем client IP (gateway — единственная точка входа, поэтому формируем цепочку) + host, _, err := net.SplitHostPort(r.RemoteAddr) + if err != nil { + host = r.RemoteAddr + } + if host != "" { + if xff := r.Header.Get("X-Forwarded-For"); xff != "" { + // не дублируем если уже содержит host + if !strings.Contains(xff, host) { + r.Header.Set("X-Forwarded-For", xff+", "+host) + } + } else { + r.Header.Set("X-Forwarded-For", host) + } + } // Заголовок X-Forwarded-Proto if r.TLS != nil { r.Header.Set("X-Forwarded-Proto", "https") diff --git a/services/metadata/cmd/server/main.go b/services/metadata/cmd/server/main.go index fb68f8b..94ba4dd 100644 --- a/services/metadata/cmd/server/main.go +++ b/services/metadata/cmd/server/main.go @@ -21,6 +21,7 @@ import ( _ "flowix/metadata/docs" // generated by swag init "flowix/metadata/internal/handler" + "flowix/metadata/internal/metrics" mw "flowix/metadata/internal/middleware" "flowix/metadata/internal/repository" "flowix/metadata/internal/storage" @@ -28,6 +29,7 @@ import ( "github.com/go-chi/chi/v5/middleware" "github.com/jackc/pgx/v5/pgxpool" "github.com/rs/zerolog" + zlog "github.com/rs/zerolog/log" httpSwagger "github.com/swaggo/http-swagger" ) @@ -58,6 +60,17 @@ func main() { minioAccess := os.Getenv("MINIO_ACCESS_KEY") minioSecret := os.Getenv("MINIO_SECRET_KEY") secure, _ := strconv.ParseBool(os.Getenv("MINIO_SECURE")) + if strings.ToLower(os.Getenv("LOG_FORMAT")) == "console" || os.Getenv("ENV") == "dev" { + zlog.Logger = zlog.Output(zerolog.ConsoleWriter{Out: os.Stdout}) + } else { + zerolog.TimeFieldFormat = zerolog.TimeFormatUnix + } + zerolog.SetGlobalLevel(zerolog.InfoLevel) + if lvl := os.Getenv("LOG_LEVEL"); lvl != "" { + if l, err := zerolog.ParseLevel(lvl); err == nil { + zerolog.SetGlobalLevel(l) + } + } logger := zerolog.New(os.Stdout).With().Timestamp().Logger() pool, err := pgxpool.New(context.Background(), dbURL) @@ -85,11 +98,12 @@ func main() { } r := chi.NewRouter() - r.Use(middleware.Logger, middleware.Recoverer, middleware.RequestID) + r.Use(middleware.RequestID, middleware.Recoverer, mw.RequestLogger) r.Get("/health", func(w http.ResponseWriter, r *http.Request) { writeJSON(w, r, http.StatusOK, map[string]string{"status": "ok", "service": "metadata"}) }) + r.Get("/metrics", metrics.Handler().ServeHTTP) r.Get("/", func(w http.ResponseWriter, r *http.Request) { writeJSON(w, r, http.StatusOK, map[string]string{"service": "metadata", "version": "0.1"}) }) diff --git a/services/metadata/go.mod b/services/metadata/go.mod index 659150b..d31e469 100644 --- a/services/metadata/go.mod +++ b/services/metadata/go.mod @@ -8,6 +8,7 @@ require ( github.com/google/uuid v1.6.0 github.com/jackc/pgx/v5 v5.7.5 github.com/minio/minio-go/v7 v7.3.0 + github.com/prometheus/client_golang v1.22.0 github.com/rs/zerolog v1.33.0 github.com/swaggo/http-swagger v1.3.4 github.com/swaggo/swag v1.16.4 @@ -15,6 +16,7 @@ require ( require ( github.com/KyleBanks/depth v1.2.1 // indirect + github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/dustin/go-humanize v1.0.1 // indirect github.com/go-openapi/jsonpointer v0.19.5 // indirect @@ -33,8 +35,11 @@ require ( github.com/mattn/go-isatty v0.0.20 // indirect github.com/minio/crc64nvme v1.1.1 // indirect github.com/minio/md5-simd v1.1.2 // indirect - github.com/minio/minio-go/v7 v7.3.0 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/philhofer/fwd v1.2.0 // indirect + github.com/prometheus/client_model v0.6.1 // indirect + github.com/prometheus/common v0.62.0 // indirect + github.com/prometheus/procfs v0.15.1 // indirect github.com/rogpeppe/go-internal v1.16.0 // indirect github.com/rs/xid v1.6.0 // indirect github.com/swaggo/files v0.0.0-20220610200504-28940afbdbfe // indirect @@ -47,6 +52,7 @@ require ( golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.41.0 // indirect golang.org/x/tools v0.48.0 // indirect + google.golang.org/protobuf v1.36.10 // indirect gopkg.in/ini.v1 v1.67.3 // indirect gopkg.in/yaml.v2 v2.4.0 // indirect ) diff --git a/services/metadata/go.sum b/services/metadata/go.sum index afaba92..35f89d1 100644 --- a/services/metadata/go.sum +++ b/services/metadata/go.sum @@ -1,5 +1,7 @@ github.com/KyleBanks/depth v1.2.1 h1:5h8fQADFrWtarTdtDudMmGsC7GPbOAu6RVB3ffsVFHc= github.com/KyleBanks/depth v1.2.1/go.mod h1:jzSb9d0L43HxTQfT+oSA1EEp2q+ne2uh6XgeJcm8brE= +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc= @@ -24,6 +26,8 @@ github.com/go-openapi/swag v0.19.15/go.mod h1:QYRuS/SOXUCsnplDa677K7+DxSOj6IPNl/ github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8= github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= @@ -44,17 +48,18 @@ github.com/klauspost/cpuid/v2 v2.4.0/go.mod h1:19jmZ9mjzoF//ddRSUsv0zfBTJWh3QJh9 github.com/klauspost/crc32 v1.3.0 h1:sSmTt3gUt81RP655XGZPElI0PelVTZ6YwCRnPSupoFM= github.com/klauspost/crc32 v1.3.0/go.mod h1:D7kQaZhnkX/Y0tstFGf8VUzv2UofNGqCjnC3zdHB0Hw= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= -github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= -github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= +github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/mailru/easyjson v0.0.0-20190614124828-94de47d64c63/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc= github.com/mailru/easyjson v0.0.0-20190626092158-b2ccc519800e/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc= github.com/mailru/easyjson v0.7.6 h1:8yTIVnZgCoiM1TgqoeTl+LfU5Jg6/xL3QhGQnimLYnA= github.com/mailru/easyjson v0.7.6/go.mod h1:xzfreul335JAWq5oZzymOObrkdz5UnU4kGfJJLY9Nlc= -github.com/mattn/go-colorable v0.1.13 h1:fFA4WZxdEF4tXPZVKMLwD8oUnCTTo08duU7wxecdEvA= github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg= github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= @@ -68,12 +73,22 @@ github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34= github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM= github.com/minio/minio-go/v7 v7.3.0 h1:HM4pFCSQq/TK+j0/zmorSh5ddh81iDgRgU0BG0Vz/YU= github.com/minio/minio-go/v7 v7.3.0/go.mod h1:KUPWdecEO1LWyUz+sTGXAuf2jZHrPh5fCsRH86QbPfk= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno= github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM= github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_golang v1.22.0 h1:rb93p9lokFEsctTys46VnV1kLCDpVZ0a/Y92Vm0Zc6Q= +github.com/prometheus/client_golang v1.22.0/go.mod h1:R7ljNsLXhuQXYZYtw6GAE9AZg8Y7vEW5scdCXrWRXC0= +github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= +github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= +github.com/prometheus/common v0.62.0 h1:xasJaQlnWAeyHdUBeGjXmutelfJHWMRr+Fg4QszZ2Io= +github.com/prometheus/common v0.62.0/go.mod h1:vyBcEuLSvWos9B1+CyL7JZ2up+uFzXhkqml0W5zIY1I= +github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= +github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= github.com/rogpeppe/go-internal v1.16.0 h1:O9DK+vNMDVGLr2BeZqmpLeMjiMNkuXfcqntWbZV6S5g= github.com/rogpeppe/go-internal v1.16.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs= github.com/rs/xid v1.5.0/go.mod h1:trrq9SKmegXys3aeAKXMUTdJsYXVwGY3RLcfgqegfbg= @@ -90,8 +105,6 @@ github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/ github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= -github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk= -github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= @@ -103,24 +116,19 @@ github.com/swaggo/swag v1.16.4 h1:clWJtd9LStiG3VeijiCfOVODP6VpHtKdQy9ELFG3s1A= github.com/swaggo/swag v1.16.4/go.mod h1:VBsHJRsDvfYvqoiMKnsdwhNV9LEMHgEDZcyVYX0sxPg= github.com/tinylib/msgp v1.6.4 h1:mOwYbyYDLPj35mkA2BjjYejgJk9BuHxDdvRnb6v2ZcQ= github.com/tinylib/msgp v1.6.4/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77roAqEA= +github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= +github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0= github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw= go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg= -golang.org/x/crypto v0.37.0 h1:kJNSjF/Xp7kU0iB2Z+9viTPMW4EqqsrywMXLJOOsXSE= -golang.org/x/crypto v0.37.0/go.mod h1:vg+k43peMZ0pUMhYmVAWysMK35e6ioLh3wB8ZCAfbVc= golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= -golang.org/x/mod v0.21.0 h1:vvrHzRwRfVKSiLrG+d4FMl/Qi4ukBCE6kZlTUkDYRT0= -golang.org/x/mod v0.21.0/go.mod h1:6SkKJ3Xj0I0BrPOZoBy3bdMptDDU9oJrpohJ3eWZ1fY= golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= golang.org/x/net v0.0.0-20210805182204-aaa1db679c0d/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= -golang.org/x/net v0.30.0 h1:AcW1SDZMkb8IpzCdQUaIq2sP4sZ4zw+55h6ynffypl4= -golang.org/x/net v0.30.0/go.mod h1:2wGyMJ5iFasEhkwi13ChkO/t1ECNC4X4eBKkVFyYFlU= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU= -golang.org/x/sync v0.13.0 h1:AauUjRAJ9OSnvULf/ARrrVywoJDy0YS2AwQ98I37610= -golang.org/x/sync v0.13.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= @@ -128,21 +136,17 @@ golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.32.0 h1:s77OFDvIQeibCmezSnk/q6iAfkdiQaJi4VzroCFrN20= -golang.org/x/sys v0.32.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.24.0 h1:dd5Bzh4yt5KYA8f9CJHCP4FB4D51c2c6JvN37xJJkJ0= -golang.org/x/text v0.24.0/go.mod h1:L8rBsPeo2pSS+xqN0d5u2ikmjtmoJbDBT1b7nHvFCdU= golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= -golang.org/x/tools v0.26.0 h1:v/60pFQmzmT9ExmjDv2gGIfi3OqfKoEP6I5+umXlbnQ= -golang.org/x/tools v0.26.0/go.mod h1:TPVVj70c7JJ3WCazhD8OdXcZg/og+b9+tH/KxylGwH0= golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= +google.golang.org/protobuf v1.36.10 h1:AYd7cD/uASjIL6Q9LiTjz8JLcrh/88q5UObnmY3aOOE= +google.golang.org/protobuf v1.36.10/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/services/metadata/internal/metrics/metrics.go b/services/metadata/internal/metrics/metrics.go new file mode 100644 index 0000000..fd56e8a --- /dev/null +++ b/services/metadata/internal/metrics/metrics.go @@ -0,0 +1,36 @@ +package metrics + +import ( + "net/http" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +var ( + RabbitMQQueueDepth = prometheus.NewGauge(prometheus.GaugeOpts{ + Name: "rabbitmq_queue_depth", + Help: "RabbitMQ queue depth", + }) + VodCacheHit = prometheus.NewCounter(prometheus.CounterOpts{ + Name: "vod_cache_hit", + Help: "VOD cache hits", + }) + UploadBytes = prometheus.NewCounter(prometheus.CounterOpts{ + Name: "upload_bytes", + Help: "Uploaded bytes", + }) + FfmpegDuration = prometheus.NewHistogram(prometheus.HistogramOpts{ + Name: "ffmpeg_duration_seconds", + Help: "FFmpeg duration", + Buckets: prometheus.DefBuckets, + }) +) + +func init() { + prometheus.MustRegister(RabbitMQQueueDepth, VodCacheHit, UploadBytes, FfmpegDuration) +} + +func Handler() http.Handler { + return promhttp.Handler() +} diff --git a/services/metadata/internal/middleware/logger.go b/services/metadata/internal/middleware/logger.go new file mode 100644 index 0000000..dd7a988 --- /dev/null +++ b/services/metadata/internal/middleware/logger.go @@ -0,0 +1,49 @@ +package middleware + +import ( + "net/http" + "time" + + "github.com/rs/zerolog/log" +) + +// RequestLogger — structured logging via zerolog, includes trace_id=request_id +func RequestLogger(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + start := time.Now() + ww := &respWriter{ResponseWriter: w, status: http.StatusOK} + next.ServeHTTP(ww, r) + dur := time.Since(start) + reqID := r.Header.Get("X-Request-ID") + if reqID == "" { + reqID = r.Header.Get("X-Request-Id") + } + if reqID == "" { + reqID = r.Header.Get("X-Correlation-ID") + } + ev := log.Info() + if ww.status >= 500 { + ev = log.Error() + } else if ww.status >= 400 { + ev = log.Warn() + } + ev.Str("method", r.Method). + Str("path", r.URL.Path). + Int("status", ww.status). + Dur("duration", dur). + Str("trace_id", reqID). + Str("request_id", reqID). + Str("service", "metadata"). + Msg("request") + }) +} + +type respWriter struct { + http.ResponseWriter + status int +} + +func (w *respWriter) WriteHeader(code int) { + w.status = code + w.ResponseWriter.WriteHeader(code) +} diff --git a/services/transcoder/app/consumer.py b/services/transcoder/app/consumer.py index f1fb6db..d17a528 100644 --- a/services/transcoder/app/consumer.py +++ b/services/transcoder/app/consumer.py @@ -12,6 +12,23 @@ import requests from minio import Minio +try: + from prometheus_client import Counter, Gauge, Histogram, start_http_server + + ffmpeg_duration_metric = Histogram( + "ffmpeg_duration_seconds", "FFmpeg transcoding duration", ["quality"] + ) + rabbitmq_queue_depth_metric = Gauge("rabbitmq_queue_depth", "RabbitMQ queue depth") + upload_bytes_metric = Counter("upload_bytes", "Total uploaded bytes") + vod_cache_hit_metric = Counter("vod_cache_hit", "VOD cache hits") + METRICS_AVAILABLE = True +except ImportError: + ffmpeg_duration_metric = None # type: ignore + rabbitmq_queue_depth_metric = None # type: ignore + upload_bytes_metric = None # type: ignore + vod_cache_hit_metric = None # type: ignore + METRICS_AVAILABLE = False + log = logging.getLogger(__name__) logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s") @@ -401,7 +418,13 @@ def transcode_one( ENCODE_MODE, " ".join(cmd), ) + start = time.time() subprocess.run(cmd, check=True, capture_output=True, timeout=900) + if METRICS_AVAILABLE and ffmpeg_duration_metric is not None: + try: + ffmpeg_duration_metric.labels(quality=f"{height}p").observe(time.time() - start) + except Exception: + pass def transcode_one_pipe( @@ -486,6 +509,7 @@ def transcode_one_pipe( cmd += ["-c:a", "copy"] cmd += ["-movflags", "+faststart", output_path] log.info("ffmpeg pipe %dx%d %dk enc=%s: %s", width, height, bitrate_k, enc, " ".join(cmd)) + start = time.time() proc = subprocess.Popen( cmd, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE ) @@ -515,6 +539,11 @@ def transcode_one_pipe( proc.stdin.write(chunk) proc.stdin.close() stdout, stderr = proc.communicate(timeout=900) + if METRICS_AVAILABLE and ffmpeg_duration_metric is not None: + try: + ffmpeg_duration_metric.labels(quality=f"{height}p").observe(time.time() - start) + except Exception: + pass if proc.returncode != 0: raise subprocess.CalledProcessError(proc.returncode, cmd, output=stdout, stderr=stderr) finally: @@ -825,6 +854,13 @@ def _get_status(video_id: str) -> str | None: def main(): signal.signal(signal.SIGTERM, _handle_sigterm) signal.signal(signal.SIGINT, _handle_sigterm) + if METRICS_AVAILABLE: + try: + mp = int(os.getenv("METRICS_PORT", "8004")) + start_http_server(mp) + log.info("metrics server started on :%d", mp) + except Exception as e: + log.warning("metrics server failed: %s", e) params = pika.URLParameters(RABBITMQ_URL) # Phase 10 fix: long transcoding (2GB ~5min) blocks heartbeat thread → broker closes connection (104). # Default heartbeat 60s is too short; set to 600s (10min) to cover 5-6GB files. For larger files Phase 11 will use chunked. diff --git a/services/transcoder/pyproject.toml b/services/transcoder/pyproject.toml index d69e58c..978523e 100644 --- a/services/transcoder/pyproject.toml +++ b/services/transcoder/pyproject.toml @@ -13,6 +13,7 @@ dependencies = [ "redis>=5", "pika>=1.3", "requests>=2.32", + "prometheus_client>=0.21", ] [dependency-groups] diff --git a/services/transcoder/uv.lock b/services/transcoder/uv.lock index 8266bcc..18d7e1d 100644 --- a/services/transcoder/uv.lock +++ b/services/transcoder/uv.lock @@ -467,6 +467,7 @@ dependencies = [ { name = "celery" }, { name = "minio" }, { name = "pika" }, + { name = "prometheus-client" }, { name = "pydantic" }, { name = "pydantic-settings" }, { name = "redis" }, @@ -488,6 +489,7 @@ requires-dist = [ { name = "celery", specifier = ">=5.4,<6" }, { name = "minio", specifier = ">=7.2,<8" }, { name = "pika", specifier = ">=1.3" }, + { name = "prometheus-client", specifier = ">=0.21" }, { name = "pydantic", specifier = ">=2.9,<3" }, { name = "pydantic-settings", specifier = ">=2.6,<3" }, { name = "redis", specifier = ">=5" }, @@ -736,6 +738,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/54/20/4d324d65cc6d9205fabedc306948156824eb9f0ee1633355a8f7ec5c66bf/pluggy-1.6.0-py3-none-any.whl", hash = "sha256:e920276dd6813095e9377c0bc5566d94c932c33b27a3e3945d8389c374dd4746", size = 20538, upload-time = "2025-05-15T12:30:06.134Z" }, ] +[[package]] +name = "prometheus-client" +version = "0.26.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/52/73/f1334c29c2af4cd9dba6c7817e61b611bd0215e2eb5565c6064a4de18802/prometheus_client-0.26.0.tar.gz", hash = "sha256:04a91bcf94e2cf74a44a1a874d651a2e853ed354b6e822f3b7487751465d5c2b", size = 92910, upload-time = "2026-07-24T19:36:41.893Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/eb/a3/b69efbf4143b5b9859b977770bbbabcc2796b702fa69dc40271e45cd5a56/prometheus_client-0.26.0-py3-none-any.whl", hash = "sha256:fa93d06737aa02bacd05794768508bb97d2fbee28cb3bca04eaae92f0ca953d6", size = 64494, upload-time = "2026-07-24T19:36:40.854Z" }, +] + [[package]] name = "prompt-toolkit" version = "3.0.53" diff --git a/services/upload/cmd/server/main.go b/services/upload/cmd/server/main.go index 5084403..64fe754 100644 --- a/services/upload/cmd/server/main.go +++ b/services/upload/cmd/server/main.go @@ -13,16 +13,19 @@ import ( "net/http" "os" "strconv" + "strings" _ "flowix/upload/docs" "flowix/upload/internal/client" "flowix/upload/internal/handler" + "flowix/upload/internal/metrics" mw "flowix/upload/internal/middleware" "flowix/upload/internal/queue" "flowix/upload/internal/storage" "github.com/go-chi/chi/v5" "github.com/go-chi/chi/v5/middleware" "github.com/rs/zerolog" + zlog "github.com/rs/zerolog/log" httpSwagger "github.com/swaggo/http-swagger" ) @@ -55,7 +58,17 @@ func main() { minioSecret := os.Getenv("MINIO_SECRET_KEY") secure, _ := strconv.ParseBool(os.Getenv("MINIO_SECURE")) - + if strings.ToLower(os.Getenv("LOG_FORMAT")) == "console" || os.Getenv("ENV") == "dev" { + zlog.Logger = zlog.Output(zerolog.ConsoleWriter{Out: os.Stdout}) + } else { + zerolog.TimeFieldFormat = zerolog.TimeFormatUnix + } + zerolog.SetGlobalLevel(zerolog.InfoLevel) + if lvl := os.Getenv("LOG_LEVEL"); lvl != "" { + if l, err := zerolog.ParseLevel(lvl); err == nil { + zerolog.SetGlobalLevel(l) + } + } logger := zerolog.New(os.Stdout).With().Timestamp().Logger() store, err := storage.NewMinioClient(minioEndpoint, minioAccess, minioSecret, bucket, secure) @@ -73,12 +86,13 @@ func main() { ph := handler.NewPresignHandler(store, pub, metaCl) r := chi.NewRouter() - r.Use(middleware.Logger, middleware.Recoverer, middleware.RequestID) + r.Use(middleware.RequestID, middleware.Recoverer, mw.RequestLogger) r.Get("/health", func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") _, _ = w.Write([]byte(`{"status":"ok","service":"upload"}`)) }) + r.Get("/metrics", metrics.Handler().ServeHTTP) r.Get("/swagger/*", httpSwagger.WrapHandler) r.Get("/openapi.json", func(w http.ResponseWriter, r *http.Request) { http.Redirect(w, r, "/swagger/doc.json", http.StatusMovedPermanently) diff --git a/services/upload/go.mod b/services/upload/go.mod index 35f3482..93e8133 100644 --- a/services/upload/go.mod +++ b/services/upload/go.mod @@ -8,6 +8,7 @@ require ( github.com/go-chi/chi/v5 v5.2.1 github.com/golang-jwt/jwt/v5 v5.2.2 github.com/minio/minio-go/v7 v7.0.70 + github.com/prometheus/client_golang v1.22.0 github.com/rabbitmq/amqp091-go v1.10.0 github.com/rs/zerolog v1.33.0 github.com/swaggo/http-swagger v1.3.4 @@ -16,6 +17,8 @@ require ( require ( github.com/KyleBanks/depth v1.2.1 // indirect + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/dustin/go-humanize v1.0.1 // indirect github.com/go-openapi/jsonpointer v0.19.5 // indirect github.com/go-openapi/jsonreference v0.20.0 // indirect @@ -24,12 +27,16 @@ require ( github.com/goccy/go-json v0.10.2 // indirect github.com/google/uuid v1.6.0 // indirect github.com/josharian/intern v1.0.0 // indirect - github.com/klauspost/compress v1.17.7 // indirect + github.com/klauspost/compress v1.18.0 // indirect github.com/klauspost/cpuid/v2 v2.2.6 // indirect github.com/mailru/easyjson v0.7.6 // indirect github.com/mattn/go-colorable v0.1.13 // indirect github.com/mattn/go-isatty v0.0.20 // indirect github.com/minio/md5-simd v1.1.2 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/prometheus/client_model v0.6.1 // indirect + github.com/prometheus/common v0.62.0 // indirect + github.com/prometheus/procfs v0.15.1 // indirect github.com/rs/xid v1.5.0 // indirect github.com/swaggo/files v0.0.0-20220610200504-28940afbdbfe // indirect golang.org/x/crypto v0.33.0 // indirect @@ -37,6 +44,7 @@ require ( golang.org/x/sys v0.32.0 // indirect golang.org/x/text v0.24.0 // indirect golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d // indirect + google.golang.org/protobuf v1.36.5 // indirect gopkg.in/ini.v1 v1.67.0 // indirect gopkg.in/yaml.v2 v2.4.0 // indirect ) diff --git a/services/upload/go.sum b/services/upload/go.sum index 1476321..0078194 100644 --- a/services/upload/go.sum +++ b/services/upload/go.sum @@ -1,5 +1,9 @@ github.com/KyleBanks/depth v1.2.1 h1:5h8fQADFrWtarTdtDudMmGsC7GPbOAu6RVB3ffsVFHc= github.com/KyleBanks/depth v1.2.1/go.mod h1:jzSb9d0L43HxTQfT+oSA1EEp2q+ne2uh6XgeJcm8brE= +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -24,20 +28,26 @@ github.com/goccy/go-json v0.10.2/go.mod h1:6MelG93GURQebXPDq3khkgXZkazVtN9CRI+MG github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8= github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/josharian/intern v1.0.0 h1:vlS4z54oSdjm0bgjRigI+G1HpF+tI+9rE5LLzOg8HmY= github.com/josharian/intern v1.0.0/go.mod h1:5DoeVV0s6jJacbCEi61lwdGj/aVlrQvzHFFd8Hwg//Y= -github.com/klauspost/compress v1.17.7 h1:ehO88t2UGzQK66LMdE8tibEd1ErmzZjNEqWkjLAKQQg= -github.com/klauspost/compress v1.17.7/go.mod h1:Di0epgTjJY877eYKx5yC51cX2A2Vl2ibi7bDH9ttBbw= +github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= +github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= github.com/klauspost/cpuid/v2 v2.0.1/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= github.com/klauspost/cpuid/v2 v2.2.6 h1:ndNyv040zDGIDh8thGkXYjnFtiN02M1PVVF+JE/48xc= github.com/klauspost/cpuid/v2 v2.2.6/go.mod h1:Lcz8mBdAVJIBVzewtcLocK12l3Y+JytZYpaMropDUws= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= +github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/mailru/easyjson v0.0.0-20190614124828-94de47d64c63/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc= github.com/mailru/easyjson v0.0.0-20190626092158-b2ccc519800e/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc= github.com/mailru/easyjson v0.7.6 h1:8yTIVnZgCoiM1TgqoeTl+LfU5Jg6/xL3QhGQnimLYnA= @@ -52,13 +62,24 @@ github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34= github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM= github.com/minio/minio-go/v7 v7.0.70 h1:1u9NtMgfK1U42kUxcsl5v0yj6TEOPR497OAQxpJnn2g= github.com/minio/minio-go/v7 v7.0.70/go.mod h1:4yBA8v80xGA30cfM3fz0DKYMXunWl/AV/6tWEs9ryzo= -github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e h1:fD57ERR4JtEqsWbfPhv4DMiApHyliiK5xCTNVSPiaAs= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_golang v1.22.0 h1:rb93p9lokFEsctTys46VnV1kLCDpVZ0a/Y92Vm0Zc6Q= +github.com/prometheus/client_golang v1.22.0/go.mod h1:R7ljNsLXhuQXYZYtw6GAE9AZg8Y7vEW5scdCXrWRXC0= +github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= +github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= +github.com/prometheus/common v0.62.0 h1:xasJaQlnWAeyHdUBeGjXmutelfJHWMRr+Fg4QszZ2Io= +github.com/prometheus/common v0.62.0/go.mod h1:vyBcEuLSvWos9B1+CyL7JZ2up+uFzXhkqml0W5zIY1I= +github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= +github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= github.com/rabbitmq/amqp091-go v1.10.0 h1:STpn5XsHlHGcecLmMFCtg7mqq0RnD+zFr4uzukfVhBw= github.com/rabbitmq/amqp091-go v1.10.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= +github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ= +github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog= github.com/rs/xid v1.5.0 h1:mKX4bl4iPYJtEIxp6CYiUuLQ/8DYMoz0PUdtGgMFRVc= github.com/rs/xid v1.5.0/go.mod h1:trrq9SKmegXys3aeAKXMUTdJsYXVwGY3RLcfgqegfbg= github.com/rs/zerolog v1.33.0 h1:1cU2KZkvPxNyfgEmhHAz/1A9Bz+llsdYzklWFzgp0r8= @@ -66,8 +87,8 @@ github.com/rs/zerolog v1.33.0/go.mod h1:/7mN4D5sKwJLZQ2b/znpjC3/GQWY/xaDXUM0kKWR github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= -github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY= -github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/swaggo/files v0.0.0-20220610200504-28940afbdbfe h1:K8pHPVoTgxFJt1lXuIzzOX7zZhZFldJQK/CgKx9BFIc= github.com/swaggo/files v0.0.0-20220610200504-28940afbdbfe/go.mod h1:lKJPbtWzJ9JhsTN1k1gZgleJWY/cqq0psdoMmaThG3w= github.com/swaggo/http-swagger v1.3.4 h1:q7t/XLx0n15H1Q9/tk3Y9L4n210XzJF5WtnDX64a5ww= @@ -100,10 +121,13 @@ golang.org/x/text v0.24.0/go.mod h1:L8rBsPeo2pSS+xqN0d5u2ikmjtmoJbDBT1b7nHvFCdU= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d h1:vU5i/LfpvrRCpgM/VPfJLg5KjxD3E+hfT1SH+d9zLwg= golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk= +google.golang.org/protobuf v1.36.5 h1:tPhr+woSbjfYvY6/GPufUoYizxw1cF/yFoxJ2fmpwlM= +google.golang.org/protobuf v1.36.5/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f h1:BLraFXnmrev5lT+xlilqcH8XK9/i0At2xKjWk4p6zsU= gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/ini.v1 v1.67.0 h1:Dgnx+6+nfE+IfzjUEISNeydPJh9AXNNsWbGP9KzCsOA= gopkg.in/ini.v1 v1.67.0/go.mod h1:pNLf8WUiyNEtQjuu5G5vTm06TEv9tsIgeAvK8hOrP4k= gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= @@ -111,5 +135,5 @@ gopkg.in/yaml.v2 v2.4.0 h1:D8xgwECY7CYvx+Y2n4sBz93Jn9JRvxdiyyo8CTfuKaY= gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.0-20200615113413-eeeca48fe776/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gopkg.in/yaml.v3 v3.0.0 h1:hjy8E9ON/egN1tAYqKb61G10WtihqetD4sz2H+8nIeA= -gopkg.in/yaml.v3 v3.0.0/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/services/upload/internal/handler/presign.go b/services/upload/internal/handler/presign.go index cd8c6c2..2d184e8 100644 --- a/services/upload/internal/handler/presign.go +++ b/services/upload/internal/handler/presign.go @@ -9,6 +9,7 @@ import ( "strings" "time" + "flowix/upload/internal/metrics" mw "flowix/upload/internal/middleware" "github.com/go-chi/chi/v5" ) @@ -143,6 +144,7 @@ func (h *PresignHandler) Complete(w http.ResponseWriter, r *http.Request) { http.Error(w, `{"error":"object not found, upload via presigned URL first"}`, 404) return } + metrics.UploadBytes.Inc() ev := VideoUploadedEvent{VideoID: videoID, S3Key: s3Key, OwnerID: ownerID} if err := h.publisher.Publish(r.Context(), ev); err != nil { http.Error(w, `{"error":"queue: `+err.Error()+`"}`, 500) diff --git a/services/upload/internal/handler/upload.go b/services/upload/internal/handler/upload.go index 1a849ec..cbfccfd 100644 --- a/services/upload/internal/handler/upload.go +++ b/services/upload/internal/handler/upload.go @@ -11,6 +11,7 @@ import ( "strconv" "strings" + "flowix/upload/internal/metrics" mw "flowix/upload/internal/middleware" ) @@ -131,6 +132,9 @@ func (h *UploadHandler) Upload(w http.ResponseWriter, r *http.Request) { return } + if header.Size > 0 { + metrics.UploadBytes.Add(float64(header.Size)) + } // 3. publish event ev := VideoUploadedEvent{VideoID: videoID, S3Key: s3Key, OwnerID: ownerID} if err := h.publisher.Publish(r.Context(), ev); err != nil { diff --git a/services/upload/internal/metrics/metrics.go b/services/upload/internal/metrics/metrics.go new file mode 100644 index 0000000..2fecd00 --- /dev/null +++ b/services/upload/internal/metrics/metrics.go @@ -0,0 +1,36 @@ +package metrics + +import ( + "net/http" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" +) + +var ( + UploadBytes = prometheus.NewCounter(prometheus.CounterOpts{ + Name: "upload_bytes", + Help: "Total uploaded bytes", + }) + RabbitMQQueueDepth = prometheus.NewGauge(prometheus.GaugeOpts{ + Name: "rabbitmq_queue_depth", + Help: "RabbitMQ queue depth", + }) + VodCacheHit = prometheus.NewCounter(prometheus.CounterOpts{ + Name: "vod_cache_hit", + Help: "VOD cache hits", + }) + FfmpegDuration = prometheus.NewHistogram(prometheus.HistogramOpts{ + Name: "ffmpeg_duration_seconds", + Help: "FFmpeg duration", + Buckets: prometheus.DefBuckets, + }) +) + +func init() { + prometheus.MustRegister(UploadBytes, RabbitMQQueueDepth, VodCacheHit, FfmpegDuration) +} + +func Handler() http.Handler { + return promhttp.Handler() +} diff --git a/services/upload/internal/middleware/logger.go b/services/upload/internal/middleware/logger.go new file mode 100644 index 0000000..dd9ab52 --- /dev/null +++ b/services/upload/internal/middleware/logger.go @@ -0,0 +1,49 @@ +package middleware + +import ( + "net/http" + "time" + + "github.com/rs/zerolog/log" +) + +// RequestLogger — structured logging via zerolog, includes trace_id=request_id +func RequestLogger(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + start := time.Now() + ww := &respWriter{ResponseWriter: w, status: http.StatusOK} + next.ServeHTTP(ww, r) + dur := time.Since(start) + reqID := r.Header.Get("X-Request-ID") + if reqID == "" { + reqID = r.Header.Get("X-Request-Id") + } + if reqID == "" { + reqID = r.Header.Get("X-Correlation-ID") + } + ev := log.Info() + if ww.status >= 500 { + ev = log.Error() + } else if ww.status >= 400 { + ev = log.Warn() + } + ev.Str("method", r.Method). + Str("path", r.URL.Path). + Int("status", ww.status). + Dur("duration", dur). + Str("trace_id", reqID). + Str("request_id", reqID). + Str("service", "upload"). + Msg("request") + }) +} + +type respWriter struct { + http.ResponseWriter + status int +} + +func (w *respWriter) WriteHeader(code int) { + w.status = code + w.ResponseWriter.WriteHeader(code) +}