From 17af76c1592150a4e68f93fc4739c6f35c37d7a8 Mon Sep 17 00:00:00 2001 From: Dmitry Gammel Date: Thu, 13 Aug 2026 11:24:25 +0500 Subject: [PATCH] =?UTF-8?q?=D0=92=D1=8B=D0=BD=D0=BE=D1=81=20health-check-s?= =?UTF-8?q?ervice=20=D0=B2=20=D0=BE=D1=82=D0=B4=D0=B5=D0=BB=D1=8C=D0=BD?= =?UTF-8?q?=D1=8B=D0=B9=20=D1=81=D0=B5=D1=80=D0=B2=D0=B8=D1=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Тикер офлайн-детекции (internal/healthcheck) переехал из device-control-service в новый самостоятельный Go-сервис — как и планировалось с самого начала (см. изначальный README-заглушку). Детекция online остаётся в device-control-service: она происходит как побочный эффект уже имеющейся там MQTT-подписки на devices/+/telemetry и devices/+/ack, заводить для этого отдельный сервис с дублирующей MQTT-подпиской избыточно. health-check-service — свой Go-модуль (без зависимости от proto/, сборка из собственного контекста), читает только узкое read+offline-подмножество Device Shadow keyspace в Redis (devices:known, status, last_seen) — полный Shadow API с записью desired/reported state остаётся только в device-control-service. Метрика device_control_health_transitions_total переименована в device_health_transitions_total и теперь публикуется с этим именем из ДВУХ сервисов (device-control-service — direction=online, health-check-service — direction=offline): Prometheus агрегирует одноимённые метрики с разных таргетов прозрачно, поэтому Grafana- дашборд адаптирован простой сменой имени метрики в запросе, без переделки панели. Проверено вживую через docker compose: реальный online→offline переход (публикация тестовой телеметрии + ожидание таймаута) корректно долетает до Redis и RabbitMQ (device.status_changed), Prometheus видит новый scrape-таргет как up. --- .env.example | 7 +- docker-compose.yml | 32 +++++- .../dashboards/automation-devices.json | 2 +- monitoring/prometheus/prometheus.yml | 4 + services/device-control-service/README.md | 12 ++- .../cmd/device-control-service/main.go | 19 +--- .../internal/config/config.go | 22 ----- .../internal/metrics/metrics.go | 8 +- services/health-check-service/Dockerfile | 12 +++ services/health-check-service/README.md | 55 +++++++++-- .../cmd/health-check-service/main.go | 97 +++++++++++++++++++ services/health-check-service/go.mod | 23 +++++ services/health-check-service/go.sum | 54 +++++++++++ .../internal/config/config.go | 84 ++++++++++++++++ .../internal/healthcheck/checker.go | 0 .../internal/healthcheck/checker_test.go | 0 .../internal/metrics/metrics.go | 21 ++++ .../internal/rabbitmq/publisher.go | 78 +++++++++++++++ .../internal/shadow/store.go | 91 +++++++++++++++++ .../internal/shadow/store_test.go | 93 ++++++++++++++++++ 20 files changed, 656 insertions(+), 58 deletions(-) create mode 100644 services/health-check-service/Dockerfile create mode 100644 services/health-check-service/cmd/health-check-service/main.go create mode 100644 services/health-check-service/go.mod create mode 100644 services/health-check-service/go.sum create mode 100644 services/health-check-service/internal/config/config.go rename services/{device-control-service => health-check-service}/internal/healthcheck/checker.go (100%) rename services/{device-control-service => health-check-service}/internal/healthcheck/checker_test.go (100%) create mode 100644 services/health-check-service/internal/metrics/metrics.go create mode 100644 services/health-check-service/internal/rabbitmq/publisher.go create mode 100644 services/health-check-service/internal/shadow/store.go create mode 100644 services/health-check-service/internal/shadow/store_test.go diff --git a/.env.example b/.env.example index 32262de..fd0e274 100644 --- a/.env.example +++ b/.env.example @@ -34,8 +34,11 @@ MQTT_PORT=1883 DEVICE_CONTROL_GRPC_PORT=50051 DEVICE_CONTROL_HTTP_PORT=8090 DEVICE_CONTROL_MQTT_CLIENT_ID=device-control-service -DEVICE_CONTROL_HEALTHCHECK_TIMEOUT=60s -DEVICE_CONTROL_HEALTHCHECK_INTERVAL=15s + +# --- health-check-service --- +HEALTH_CHECK_TIMEOUT=60s +HEALTH_CHECK_INTERVAL=15s +HEALTH_CHECK_METRICS_PORT=9103 # --- ingest-service --- INGEST_MQTT_CLIENT_ID=ingest-service diff --git a/docker-compose.yml b/docker-compose.yml index e9afd92..838703d 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -140,8 +140,6 @@ services: RABBITMQ_PORT: 5672 RABBITMQ_USER: ${RABBITMQ_USER} RABBITMQ_PASSWORD: ${RABBITMQ_PASSWORD} - DEVICE_CONTROL_HEALTHCHECK_TIMEOUT: ${DEVICE_CONTROL_HEALTHCHECK_TIMEOUT} - DEVICE_CONTROL_HEALTHCHECK_INTERVAL: ${DEVICE_CONTROL_HEALTHCHECK_INTERVAL} depends_on: mosquitto: condition: service_started @@ -153,6 +151,33 @@ services: - home-automation restart: unless-stopped + health-check-service: + build: ./services/health-check-service + ports: + - "${HEALTH_CHECK_METRICS_PORT}:9103" + environment: + REDIS_HOST: redis + REDIS_PORT: 6379 + REDIS_PASSWORD: ${REDIS_PASSWORD} + REDIS_DB: ${REDIS_DB} + RABBITMQ_HOST: rabbitmq + RABBITMQ_PORT: 5672 + RABBITMQ_USER: ${RABBITMQ_USER} + RABBITMQ_PASSWORD: ${RABBITMQ_PASSWORD} + HEALTH_CHECK_TIMEOUT: ${HEALTH_CHECK_TIMEOUT} + HEALTH_CHECK_INTERVAL: ${HEALTH_CHECK_INTERVAL} + HEALTH_CHECK_METRICS_PORT: 9103 + depends_on: + redis: + condition: service_started + rabbitmq: + condition: service_healthy + device-control-service: + condition: service_started + networks: + - home-automation + restart: unless-stopped + rule-engine-service: build: context: . @@ -258,5 +283,4 @@ services: - home-automation restart: unless-stopped - # Remaining app services (health-check-service, esp32-emulator) are added - # here as they get implemented — see services/*/README.md for the plan. + # esp32-emulator is added here once implemented — see its README for the plan. diff --git a/monitoring/grafana/provisioning/dashboards/automation-devices.json b/monitoring/grafana/provisioning/dashboards/automation-devices.json index b43f736..80f86c1 100644 --- a/monitoring/grafana/provisioning/dashboards/automation-devices.json +++ b/monitoring/grafana/provisioning/dashboards/automation-devices.json @@ -59,7 +59,7 @@ "gridPos": { "h": 8, "w": 12, "x": 12, "y": 8 }, "targets": [ { - "expr": "sum by (direction) (rate(device_control_health_transitions_total[5m]))", + "expr": "sum by (direction) (rate(device_health_transitions_total[5m]))", "legendFormat": "{{direction}}", "refId": "A" } diff --git a/monitoring/prometheus/prometheus.yml b/monitoring/prometheus/prometheus.yml index 54912a3..8d26505 100644 --- a/monitoring/prometheus/prometheus.yml +++ b/monitoring/prometheus/prometheus.yml @@ -13,3 +13,7 @@ scrape_configs: - job_name: rule-engine-service static_configs: - targets: ["rule-engine-service:9102"] + + - job_name: health-check-service + static_configs: + - targets: ["health-check-service:9103"] diff --git a/services/device-control-service/README.md b/services/device-control-service/README.md index e693472..4227138 100644 --- a/services/device-control-service/README.md +++ b/services/device-control-service/README.md @@ -15,10 +15,10 @@ - Поддерживает паттерн Device Shadow в Redis: `desired_state` (чего хочет пользователь) против `reported_state` (что подтвердило устройство), плюс `status` и `last_seen`. -- Ведёт health-check горутиной-тикером: переводит устройство в `offline`, - когда `last_seen` превышает таймаут, и публикует событие в RabbitMQ - (`device.status_changed`) для будущего notification-service — как при - уходе в офлайн, так и при возврате online. +- Переводит устройство в `online` и публикует `device.status_changed` в + RabbitMQ при получении телеметрии/ack (см. «Архитектурные решения»). + Обратный переход, `offline` по таймауту `last_seen`, — зона + ответственности отдельного health-check-service (см. его README). Не входит в зону ответственности: решение о том, *когда* отправлять команду на основе показаний датчиков — это логика rule-engine-service. Этот сервис @@ -54,7 +54,9 @@ `GET /metrics` на том же порту, что и команды (`DEVICE_CONTROL_HTTP_PORT`) — не открывали отдельный порт ради этого. Счётчики/latency команд (TurnOn/TurnOff/SetLevel по action+outcome), MQTT-событий (telemetry/ack), -переходов online/offline health-check. +переходов в online (`device_health_transitions_total{direction="online"}` +— та же метрика по имени, что публикует health-check-service для +`direction="offline"`, см. его README). ## Запуск diff --git a/services/device-control-service/cmd/device-control-service/main.go b/services/device-control-service/cmd/device-control-service/main.go index 38b7535..513c8d6 100644 --- a/services/device-control-service/cmd/device-control-service/main.go +++ b/services/device-control-service/cmd/device-control-service/main.go @@ -1,6 +1,8 @@ // Command device-control-service exposes a gRPC API for issuing device -// commands, dispatches them over MQTT, maintains the Device Shadow in -// Redis, and runs a health-check loop that flips stale devices offline. +// commands, dispatches them over MQTT, and maintains the Device Shadow in +// Redis. Offline detection lives in the separate health-check-service — +// this service only ever flips a device back to online, as a side effect +// of handling its telemetry/ack MQTT messages. package main import ( @@ -21,7 +23,6 @@ import ( devicecontrol "git.cactoz.su/cacto/home_automatization/proto/device_control" "git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/config" - "git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/healthcheck" "git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/httpapi" "git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/metrics" "git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/mqttclient" @@ -75,18 +76,6 @@ func run(logger *slog.Logger) error { return fmt.Errorf("subscribe ack: %w", err) } - checker := healthcheck.New(shadowStore, cfg.HealthCheckTimeout, cfg.HealthCheckInterval, func(deviceID string) { - metrics.HealthTransitionsTotal.WithLabelValues("offline").Inc() - - pubCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) - defer cancel() - if err := publisher.PublishStatusChanged(pubCtx, deviceID, shadow.StatusOffline); err != nil { - logger.Error("publish offline event failed", "device_id", deviceID, "error", err) - } - }, logger) - checker.Start() - defer checker.Stop() - srv := server.New(shadowStore, mqttClient, logger) lis, err := net.Listen("tcp", fmt.Sprintf(":%d", cfg.GRPCPort)) diff --git a/services/device-control-service/internal/config/config.go b/services/device-control-service/internal/config/config.go index 12263f1..5806599 100644 --- a/services/device-control-service/internal/config/config.go +++ b/services/device-control-service/internal/config/config.go @@ -6,7 +6,6 @@ import ( "fmt" "os" "strconv" - "time" ) type Config struct { @@ -21,9 +20,6 @@ type Config struct { RedisDB int RabbitMQURL string - - HealthCheckTimeout time.Duration - HealthCheckInterval time.Duration } func Load() (Config, error) { @@ -52,12 +48,6 @@ func Load() (Config, error) { if cfg.RedisDB, err = getEnvInt("REDIS_DB", 0); err != nil { return Config{}, err } - if cfg.HealthCheckTimeout, err = getEnvDuration("DEVICE_CONTROL_HEALTHCHECK_TIMEOUT", 60*time.Second); err != nil { - return Config{}, err - } - if cfg.HealthCheckInterval, err = getEnvDuration("DEVICE_CONTROL_HEALTHCHECK_INTERVAL", 15*time.Second); err != nil { - return Config{}, err - } return cfg, nil } @@ -80,15 +70,3 @@ func getEnvInt(key string, fallback int) (int, error) { } return n, nil } - -func getEnvDuration(key string, fallback time.Duration) (time.Duration, error) { - v := os.Getenv(key) - if v == "" { - return fallback, nil - } - d, err := time.ParseDuration(v) - if err != nil { - return 0, fmt.Errorf("%s: %w", key, err) - } - return d, nil -} diff --git a/services/device-control-service/internal/metrics/metrics.go b/services/device-control-service/internal/metrics/metrics.go index 1aae263..e859d46 100644 --- a/services/device-control-service/internal/metrics/metrics.go +++ b/services/device-control-service/internal/metrics/metrics.go @@ -25,8 +25,14 @@ var ( Help: "Total number of telemetry/ack messages received, by event type and outcome.", }, []string{"event_type", "outcome"}) + // Name matches health-check-service's own metric of the same purpose — + // this service only ever emits direction="online" (detected via MQTT + // telemetry/ack), health-check-service only ever emits "offline". + // Prometheus aggregates same-named metrics across scrape targets, so + // Grafana's `sum by (direction) (...)` panel works across both without + // a query change. HealthTransitionsTotal = promauto.NewCounterVec(prometheus.CounterOpts{ - Name: "device_control_health_transitions_total", + Name: "device_health_transitions_total", Help: "Total number of device online/offline transitions detected.", }, []string{"direction"}) ) diff --git a/services/health-check-service/Dockerfile b/services/health-check-service/Dockerfile new file mode 100644 index 0000000..9a59102 --- /dev/null +++ b/services/health-check-service/Dockerfile @@ -0,0 +1,12 @@ +FROM golang:1.25-alpine AS build +WORKDIR /src +COPY go.mod go.sum ./ +RUN go mod download +COPY . . +RUN CGO_ENABLED=0 go build -o /out/health-check-service ./cmd/health-check-service + +FROM alpine:3.20 +RUN adduser -D -u 10001 app +COPY --from=build /out/health-check-service /usr/local/bin/health-check-service +USER app +ENTRYPOINT ["/usr/local/bin/health-check-service"] diff --git a/services/health-check-service/README.md b/services/health-check-service/README.md index d77e36c..c084216 100644 --- a/services/health-check-service/README.md +++ b/services/health-check-service/README.md @@ -1,12 +1,51 @@ # health-check-service (Go) -Статус: пока не реализован. Может начаться как горутина внутри -device-control-service и позже выделиться в отдельный контейнер — поэтому -с первого дня держим отдельную директорию, чтобы это разделение прошло -безболезненно. +Статус: реализован. + +Изначально жил как горутина-тикер внутри device-control-service (см. его +README/git-историю) — выделен в отдельный сервис по мере роста проекта, +как и планировалось с самого начала. Зона ответственности: -- Периодически (по тикеру) проверяет `last_seen` каждого устройства в Redis. -- Когда устройство превышает таймаут офлайна, переводит его `status` в - `offline` и публикует событие для notification-service. -- Хороший повод продемонстрировать горутины + graceful shutdown в Go. +- Периодически (по тикеру, `HEALTH_CHECK_INTERVAL`) обходит все известные + устройства из Redis-множества `devices:known` и проверяет `last_seen` + каждого. +- Когда устройство превышает таймаут офлайна (`HEALTH_CHECK_TIMEOUT`), + переводит его `status` в `offline` в том же Redis (Device Shadow + keyspace, который device-control-service читает/пишет) и публикует + `device.status_changed` в RabbitMQ. + +Не входит в зону ответственности: +- **Детекция перехода в online.** Это происходит как побочный эффект + MQTT-подписки device-control-service на `devices/+/telemetry` и + `devices/+/ack` (там уже есть открытое соединение и подписка ради + reported_state) — заводить для этого отдельный сервис с собственной + MQTT-подпиской избыточно. health-check-service работает только "в одну + сторону": online → offline. +- Хранение/интерпретация телеметрии — это ingest-service. +- Знание о существовании устройств из PostgreSQL — этот сервис, как и + device-control-service, не подключается к Postgres: список устройств для + сканирования берётся из того же Redis-множества `devices:known` + (пополняется device-control-service). + +## Метрики + +`GET /metrics` на `HEALTH_CHECK_METRICS_PORT`. Счётчик +`device_health_transitions_total{direction="offline"}` — намеренно та же +метрика (по имени), что device-control-service публикует для +`direction="online"`: Prometheus агрегирует одноимённые метрики с разных +таргетов прозрачно, так что дашборд Grafana не пришлось переделывать под +разделение на два сервиса. + +## Запуск + +```bash +cd services/health-check-service +go test ./... +go build ./cmd/health-check-service +``` + +Конфигурация — через переменные окружения (секция health-check-service в +корневом `.env.example`). Не зависит от `proto/` — работает только с Redis +и RabbitMQ, поэтому в отличие от device-control-service/rule-engine-service +сборка Docker-образа идёт из собственного контекста (не корня репозитория). diff --git a/services/health-check-service/cmd/health-check-service/main.go b/services/health-check-service/cmd/health-check-service/main.go new file mode 100644 index 0000000..4fe6e27 --- /dev/null +++ b/services/health-check-service/cmd/health-check-service/main.go @@ -0,0 +1,97 @@ +// Command health-check-service periodically scans every known device in +// Redis and flips it to offline once its last_seen exceeds the configured +// timeout, publishing a device.status_changed event to RabbitMQ for each +// transition. It never marks a device online — that happens as a side +// effect of device-control-service's own MQTT telemetry/ack handling — this +// service only ever detects silence. +package main + +import ( + "context" + "fmt" + "log/slog" + "net/http" + "os" + "os/signal" + "syscall" + "time" + + "github.com/prometheus/client_golang/prometheus/promhttp" + + "git.cactoz.su/cacto/home_automatization/services/health-check-service/internal/config" + "git.cactoz.su/cacto/home_automatization/services/health-check-service/internal/healthcheck" + "git.cactoz.su/cacto/home_automatization/services/health-check-service/internal/metrics" + "git.cactoz.su/cacto/home_automatization/services/health-check-service/internal/rabbitmq" + "git.cactoz.su/cacto/home_automatization/services/health-check-service/internal/shadow" +) + +func main() { + logger := slog.New(slog.NewJSONHandler(os.Stdout, nil)) + + if err := run(logger); err != nil { + logger.Error("health-check-service exited with error", "error", err) + os.Exit(1) + } +} + +func run(logger *slog.Logger) error { + cfg, err := config.Load() + if err != nil { + return err + } + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + store, err := shadow.Connect(ctx, cfg.RedisAddr, cfg.RedisPassword, cfg.RedisDB) + cancel() + if err != nil { + return err + } + defer store.Close() + + publisher, err := rabbitmq.Connect(cfg.RabbitMQURL) + if err != nil { + return err + } + defer publisher.Close() + + checker := healthcheck.New(store, cfg.HealthCheckTimeout, cfg.HealthCheckInterval, func(deviceID string) { + metrics.HealthTransitionsTotal.WithLabelValues("offline").Inc() + + pubCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := publisher.PublishStatusChanged(pubCtx, deviceID, shadow.StatusOffline); err != nil { + logger.Error("publish offline event failed", "device_id", deviceID, "error", err) + } + }, logger) + checker.Start() + defer checker.Stop() + + httpServer := &http.Server{ + Addr: fmt.Sprintf(":%d", cfg.MetricsPort), + Handler: promhttp.Handler(), + } + + go func() { + logger.Info("metrics server started", "metrics_port", cfg.MetricsPort) + if err := httpServer.ListenAndServe(); err != nil && err != http.ErrServerClosed { + logger.Error("metrics server stopped", "error", err) + } + }() + + logger.Info("health-check-service started", + "timeout", cfg.HealthCheckTimeout, "interval", cfg.HealthCheckInterval) + + stop := make(chan os.Signal, 1) + signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM) + <-stop + + logger.Info("shutting down") + + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := httpServer.Shutdown(shutdownCtx); err != nil { + logger.Error("metrics server shutdown failed", "error", err) + } + + return nil +} diff --git a/services/health-check-service/go.mod b/services/health-check-service/go.mod new file mode 100644 index 0000000..d52be7c --- /dev/null +++ b/services/health-check-service/go.mod @@ -0,0 +1,23 @@ +module git.cactoz.su/cacto/home_automatization/services/health-check-service + +go 1.25.0 + +require ( + github.com/alicebob/miniredis/v2 v2.38.0 + github.com/prometheus/client_golang v1.24.1 + github.com/rabbitmq/amqp091-go v1.13.0 + github.com/redis/go-redis/v9 v9.22.0 +) + +require ( + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/prometheus/client_model v0.6.2 // indirect + github.com/prometheus/common v0.70.1 // indirect + github.com/prometheus/procfs v0.21.1 // indirect + github.com/yuin/gopher-lua v1.1.1 // indirect + go.uber.org/atomic v1.11.0 // indirect + golang.org/x/sys v0.47.0 // indirect + google.golang.org/protobuf v1.36.11 // indirect +) diff --git a/services/health-check-service/go.sum b/services/health-check-service/go.sum new file mode 100644 index 0000000..c0027e7 --- /dev/null +++ b/services/health-check-service/go.sum @@ -0,0 +1,54 @@ +github.com/alicebob/miniredis/v2 v2.38.0 h1:nZAzCR+Lj+Vxk4ZXzm2NuKq2O33RXj1XxJ2e2uP9jiw= +github.com/alicebob/miniredis/v2 v2.38.0/go.mod h1:TcL7YfarKPGDAthEtl5NBeHZfeUQj6OXMm/+iu5cLMM= +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/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= +github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= +github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= +github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= +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/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.19.1 h1:VsB4HPswih7mmZ8WleSFQ75c/Ui1M4trX5oAsJnhSlk= +github.com/klauspost/compress v1.19.1/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE= +github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +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/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.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU= +github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE= +github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= +github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= +github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY= +github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc= +github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI= +github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY= +github.com/rabbitmq/amqp091-go v1.13.0 h1:L8NA1WtF76C6KA3LAoufjfLgbist/If1UQYcsOjtxXA= +github.com/rabbitmq/amqp091-go v1.13.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= +github.com/redis/go-redis/v9 v9.22.0 h1:laDvpYXTJtZLloinw1fA5Kqd6HAEH2XKxOkG/PDq2F0= +github.com/redis/go-redis/v9 v9.22.0/go.mod h1:y2g0Wj8rQvuK0ELM+oxSudcLtC09JScs98I/X9gRWY4= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/yuin/gopher-lua v1.1.1 h1:kYKnWBjvbNP4XLT3+bPEwAXJx262OhaHDWDVOPjL46M= +github.com/yuin/gopher-lua v1.1.1/go.mod h1:GBR0iDaNXjAgGg9zfCvksxSRnQx76gclCIb7kdAd1Pw= +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.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= +go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ= +go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= +golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= +golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= +google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +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/health-check-service/internal/config/config.go b/services/health-check-service/internal/config/config.go new file mode 100644 index 0000000..391636e --- /dev/null +++ b/services/health-check-service/internal/config/config.go @@ -0,0 +1,84 @@ +// Package config loads health-check-service settings from environment +// variables, matching the names used in the repo-root .env.example. +package config + +import ( + "fmt" + "os" + "strconv" + "time" +) + +type Config struct { + RedisAddr string + RedisPassword string + RedisDB int + + RabbitMQURL string + + HealthCheckTimeout time.Duration + HealthCheckInterval time.Duration + + MetricsPort int +} + +func Load() (Config, error) { + cfg := Config{ + RedisAddr: fmt.Sprintf("%s:%s", getEnv("REDIS_HOST", "localhost"), getEnv("REDIS_PORT", "6379")), + RedisPassword: getEnv("REDIS_PASSWORD", ""), + + RabbitMQURL: fmt.Sprintf("amqp://%s:%s@%s:%s/", + getEnv("RABBITMQ_USER", "guest"), + getEnv("RABBITMQ_PASSWORD", "guest"), + getEnv("RABBITMQ_HOST", "localhost"), + getEnv("RABBITMQ_PORT", "5672"), + ), + } + + var err error + if cfg.RedisDB, err = getEnvInt("REDIS_DB", 0); err != nil { + return Config{}, err + } + if cfg.HealthCheckTimeout, err = getEnvDuration("HEALTH_CHECK_TIMEOUT", 60*time.Second); err != nil { + return Config{}, err + } + if cfg.HealthCheckInterval, err = getEnvDuration("HEALTH_CHECK_INTERVAL", 15*time.Second); err != nil { + return Config{}, err + } + if cfg.MetricsPort, err = getEnvInt("HEALTH_CHECK_METRICS_PORT", 9103); err != nil { + return Config{}, err + } + + return cfg, nil +} + +func getEnv(key, fallback string) string { + if v := os.Getenv(key); v != "" { + return v + } + return fallback +} + +func getEnvInt(key string, fallback int) (int, error) { + v := os.Getenv(key) + if v == "" { + return fallback, nil + } + n, err := strconv.Atoi(v) + if err != nil { + return 0, fmt.Errorf("%s: %w", key, err) + } + return n, nil +} + +func getEnvDuration(key string, fallback time.Duration) (time.Duration, error) { + v := os.Getenv(key) + if v == "" { + return fallback, nil + } + d, err := time.ParseDuration(v) + if err != nil { + return 0, fmt.Errorf("%s: %w", key, err) + } + return d, nil +} diff --git a/services/device-control-service/internal/healthcheck/checker.go b/services/health-check-service/internal/healthcheck/checker.go similarity index 100% rename from services/device-control-service/internal/healthcheck/checker.go rename to services/health-check-service/internal/healthcheck/checker.go diff --git a/services/device-control-service/internal/healthcheck/checker_test.go b/services/health-check-service/internal/healthcheck/checker_test.go similarity index 100% rename from services/device-control-service/internal/healthcheck/checker_test.go rename to services/health-check-service/internal/healthcheck/checker_test.go diff --git a/services/health-check-service/internal/metrics/metrics.go b/services/health-check-service/internal/metrics/metrics.go new file mode 100644 index 0000000..152b142 --- /dev/null +++ b/services/health-check-service/internal/metrics/metrics.go @@ -0,0 +1,21 @@ +// Package metrics defines health-check-service's Prometheus metrics. +// Registered automatically (via promauto) into the default registry on +// import; served on /metrics. +// +// device_health_transitions_total is intentionally the same metric name +// device-control-service exposes for its "online" transitions (detected via +// MQTT touch) — this service only ever emits direction="offline". Prometheus +// aggregates same-named metrics across scrape targets/jobs transparently, so +// Grafana's `sum by (direction) (...)` panel keeps working across the split +// without a query change. +package metrics + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +var HealthTransitionsTotal = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "device_health_transitions_total", + Help: "Total number of device online/offline transitions detected.", +}, []string{"direction"}) diff --git a/services/health-check-service/internal/rabbitmq/publisher.go b/services/health-check-service/internal/rabbitmq/publisher.go new file mode 100644 index 0000000..7916ef2 --- /dev/null +++ b/services/health-check-service/internal/rabbitmq/publisher.go @@ -0,0 +1,78 @@ +// Package rabbitmq publishes device online/offline transitions for +// notification-service to consume. Delivery reliability matters more than +// latency here, so messages are persistent and the queue is durable. +package rabbitmq + +import ( + "context" + "encoding/json" + "fmt" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +const StatusChangedQueue = "device.status_changed" + +type Publisher struct { + conn *amqp.Connection + ch *amqp.Channel +} + +func Connect(url string) (*Publisher, error) { + conn, err := amqp.Dial(url) + if err != nil { + return nil, fmt.Errorf("dial rabbitmq: %w", err) + } + + ch, err := conn.Channel() + if err != nil { + conn.Close() + return nil, fmt.Errorf("open channel: %w", err) + } + + if _, err := ch.QueueDeclare(StatusChangedQueue, true, false, false, false, nil); err != nil { + ch.Close() + conn.Close() + return nil, fmt.Errorf("declare queue %s: %w", StatusChangedQueue, err) + } + + return &Publisher{conn: conn, ch: ch}, nil +} + +func (p *Publisher) Close() error { + if err := p.ch.Close(); err != nil { + p.conn.Close() + return err + } + return p.conn.Close() +} + +type statusChangedEvent struct { + DeviceID string `json:"device_id"` + Status string `json:"status"` + At time.Time `json:"at"` +} + +// PublishStatusChanged emits one event when a device transitions online or offline. +func (p *Publisher) PublishStatusChanged(ctx context.Context, deviceID, status string) error { + body, err := json.Marshal(statusChangedEvent{ + DeviceID: deviceID, + Status: status, + At: time.Now(), + }) + if err != nil { + return fmt.Errorf("marshal event: %w", err) + } + + err = p.ch.PublishWithContext(ctx, "", StatusChangedQueue, false, false, amqp.Publishing{ + ContentType: "application/json", + DeliveryMode: amqp.Persistent, + Timestamp: time.Now(), + Body: body, + }) + if err != nil { + return fmt.Errorf("publish status change for device %s: %w", deviceID, err) + } + return nil +} diff --git a/services/health-check-service/internal/shadow/store.go b/services/health-check-service/internal/shadow/store.go new file mode 100644 index 0000000..5a5d12b --- /dev/null +++ b/services/health-check-service/internal/shadow/store.go @@ -0,0 +1,91 @@ +// Package shadow gives health-check-service read/offline-detection access +// to the Device Shadow keyspace in Redis. device-control-service owns the +// full Shadow API (desired/reported state patches, online-touch on +// telemetry/ack) — this is deliberately the narrow subset a periodic +// offline-scanner needs, kept as its own small copy rather than a shared +// module: the two services only overlap on reading/writing `status` and +// `last_seen`, and duplicating ~2 key-naming functions is cheaper than a +// fourth shared Go module for that. +package shadow + +import ( + "context" + "fmt" + "strconv" + "time" + + "github.com/redis/go-redis/v9" +) + +const knownDevicesKey = "devices:known" + +func statusKey(deviceID string) string { return "device:" + deviceID + ":status" } +func lastSeenKey(deviceID string) string { return "device:" + deviceID + ":last_seen" } + +const ( + StatusOnline = "online" + StatusOffline = "offline" +) + +type Store struct { + rdb *redis.Client +} + +func New(rdb *redis.Client) *Store { + return &Store{rdb: rdb} +} + +func Connect(ctx context.Context, addr, password string, db int) (*Store, error) { + rdb := redis.NewClient(&redis.Options{Addr: addr, Password: password, DB: db}) + if err := rdb.Ping(ctx).Err(); err != nil { + return nil, fmt.Errorf("ping redis: %w", err) + } + return New(rdb), nil +} + +func (s *Store) Close() error { + return s.rdb.Close() +} + +// KnownDevices lists every device ID device-control-service has heard from. +func (s *Store) KnownDevices(ctx context.Context) ([]string, error) { + ids, err := s.rdb.SMembers(ctx, knownDevicesKey).Result() + if err != nil { + return nil, fmt.Errorf("list known devices: %w", err) + } + return ids, nil +} + +// MarkOfflineIfStale flips deviceID to offline if it is currently online and +// its last_seen is older than timeout. It reports whether a change was made. +func (s *Store) MarkOfflineIfStale(ctx context.Context, deviceID string, timeout time.Duration) (changed bool, err error) { + status, err := s.rdb.Get(ctx, statusKey(deviceID)).Result() + if err == redis.Nil || status != StatusOnline { + return false, nil + } + if err != nil { + return false, fmt.Errorf("get status for %s: %w", deviceID, err) + } + + lastSeenRaw, err := s.rdb.Get(ctx, lastSeenKey(deviceID)).Result() + if err != nil && err != redis.Nil { + return false, fmt.Errorf("get last_seen for %s: %w", deviceID, err) + } + + var lastSeen time.Time + if lastSeenRaw != "" { + sec, err := strconv.ParseInt(lastSeenRaw, 10, 64) + if err != nil { + return false, fmt.Errorf("parse last_seen for %s: %w", deviceID, err) + } + lastSeen = time.Unix(sec, 0) + } + + if time.Since(lastSeen) <= timeout { + return false, nil + } + if err := s.rdb.Set(ctx, statusKey(deviceID), StatusOffline, 0).Err(); err != nil { + return false, fmt.Errorf("set status offline for %s: %w", deviceID, err) + } + return true, nil +} diff --git a/services/health-check-service/internal/shadow/store_test.go b/services/health-check-service/internal/shadow/store_test.go new file mode 100644 index 0000000..12a1df1 --- /dev/null +++ b/services/health-check-service/internal/shadow/store_test.go @@ -0,0 +1,93 @@ +package shadow + +import ( + "context" + "testing" + "time" + + "github.com/alicebob/miniredis/v2" + "github.com/redis/go-redis/v9" +) + +func newTestStore(t *testing.T) *Store { + t.Helper() + mr := miniredis.RunT(t) + rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()}) + t.Cleanup(func() { rdb.Close() }) + return New(rdb) +} + +// markOnline seeds Redis the way device-control-service's Touch() would, +// without depending on that service's code. +func markOnline(t *testing.T, s *Store, deviceID string, lastSeen time.Time) { + t.Helper() + ctx := context.Background() + if err := s.rdb.SAdd(ctx, knownDevicesKey, deviceID).Err(); err != nil { + t.Fatalf("seed known device: %v", err) + } + if err := s.rdb.Set(ctx, statusKey(deviceID), StatusOnline, 0).Err(); err != nil { + t.Fatalf("seed status: %v", err) + } + if err := s.rdb.Set(ctx, lastSeenKey(deviceID), lastSeen.Unix(), 0).Err(); err != nil { + t.Fatalf("seed last_seen: %v", err) + } +} + +func TestKnownDevices_ListsSeededDevices(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + markOnline(t, s, "d1", time.Now()) + markOnline(t, s, "d2", time.Now()) + + got, err := s.KnownDevices(ctx) + if err != nil { + t.Fatalf("known devices: %v", err) + } + if len(got) != 2 { + t.Fatalf("got %d known devices, want 2 (%v)", len(got), got) + } +} + +func TestMarkOfflineIfStale(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + markOnline(t, s, "d1", time.Now()) + + changed, err := s.MarkOfflineIfStale(ctx, "d1", time.Hour) + if err != nil { + t.Fatalf("mark offline (not stale): %v", err) + } + if changed { + t.Fatal("device just touched should not be considered stale") + } + + changed, err = s.MarkOfflineIfStale(ctx, "d1", -time.Second) // any age counts as stale + if err != nil { + t.Fatalf("mark offline (stale): %v", err) + } + if !changed { + t.Fatal("expected status to flip to offline") + } + + // Already offline: a second call should report no further change. + changed, err = s.MarkOfflineIfStale(ctx, "d1", -time.Second) + if err != nil { + t.Fatalf("mark offline (already offline): %v", err) + } + if changed { + t.Fatal("expected no change once already offline") + } +} + +func TestMarkOfflineIfStale_UnknownDeviceIsNoop(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + + changed, err := s.MarkOfflineIfStale(ctx, "ghost", time.Second) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if changed { + t.Fatal("unknown device should never be reported as changed") + } +}