Вынос health-check-service в отдельный сервис

Тикер офлайн-детекции (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.
This commit is contained in:
2026-08-13 11:24:25 +05:00
parent be1f238f89
commit 17af76c159
20 changed files with 656 additions and 58 deletions
+12
View File
@@ -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"]
+47 -8
View File
@@ -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-образа идёт из собственного контекста (не корня репозитория).
@@ -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
}
+23
View File
@@ -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
)
+54
View File
@@ -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=
@@ -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
}
@@ -0,0 +1,93 @@
// Package healthcheck periodically flips devices to offline when they stop
// sending telemetry/acks, per the platform's Device Shadow health-check
// design (a goroutine + graceful shutdown, deliberately simple).
package healthcheck
import (
"context"
"log/slog"
"time"
)
// DeviceStore is the subset of shadow.Store the checker needs, kept as an
// interface so tests don't require a real Redis.
type DeviceStore interface {
KnownDevices(ctx context.Context) ([]string, error)
MarkOfflineIfStale(ctx context.Context, deviceID string, timeout time.Duration) (bool, error)
}
// OnOfflineFunc is called once per device that just transitioned to offline.
type OnOfflineFunc func(deviceID string)
type Checker struct {
store DeviceStore
timeout time.Duration
interval time.Duration
onOffline OnOfflineFunc
logger *slog.Logger
stop chan struct{}
done chan struct{}
}
func New(store DeviceStore, timeout, interval time.Duration, onOffline OnOfflineFunc, logger *slog.Logger) *Checker {
return &Checker{
store: store,
timeout: timeout,
interval: interval,
onOffline: onOffline,
logger: logger,
}
}
// Start begins the periodic scan loop. Call once.
func (c *Checker) Start() {
c.stop = make(chan struct{})
c.done = make(chan struct{})
go c.loop()
}
// Stop signals the loop to exit and waits for it to finish.
func (c *Checker) Stop() {
close(c.stop)
<-c.done
}
func (c *Checker) loop() {
defer close(c.done)
ticker := time.NewTicker(c.interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
c.checkAll()
case <-c.stop:
return
}
}
}
func (c *Checker) checkAll() {
ctx, cancel := context.WithTimeout(context.Background(), c.interval)
defer cancel()
ids, err := c.store.KnownDevices(ctx)
if err != nil {
c.logger.Error("healthcheck: list known devices failed", "error", err)
return
}
for _, id := range ids {
changed, err := c.store.MarkOfflineIfStale(ctx, id, c.timeout)
if err != nil {
c.logger.Error("healthcheck: check device failed", "device_id", id, "error", err)
continue
}
if changed {
c.logger.Info("device marked offline", "device_id", id)
if c.onOffline != nil {
c.onOffline(id)
}
}
}
}
@@ -0,0 +1,76 @@
package healthcheck
import (
"context"
"io"
"log/slog"
"sync"
"testing"
"time"
)
type fakeStore struct {
mu sync.Mutex
known []string
stale map[string]bool // deviceID -> whether MarkOfflineIfStale should report a change
changes []string // devices actually marked changed, in call order
}
func (f *fakeStore) KnownDevices(context.Context) ([]string, error) {
f.mu.Lock()
defer f.mu.Unlock()
return append([]string(nil), f.known...), nil
}
func (f *fakeStore) MarkOfflineIfStale(_ context.Context, deviceID string, _ time.Duration) (bool, error) {
f.mu.Lock()
defer f.mu.Unlock()
if !f.stale[deviceID] {
return false, nil
}
// Only report the transition once, like the real store would.
f.stale[deviceID] = false
f.changes = append(f.changes, deviceID)
return true, nil
}
func discardLogger() *slog.Logger {
return slog.New(slog.NewTextHandler(io.Discard, nil))
}
func TestChecker_FlipsStaleDevicesAndCallsOnOffline(t *testing.T) {
store := &fakeStore{
known: []string{"d1", "d2"},
stale: map[string]bool{"d1": true, "d2": false},
}
offline := make(chan string, 2)
c := New(store, time.Minute, 10*time.Millisecond, func(deviceID string) {
offline <- deviceID
}, discardLogger())
c.Start()
defer c.Stop()
select {
case id := <-offline:
if id != "d1" {
t.Fatalf("got offline callback for %q, want d1", id)
}
case <-time.After(time.Second):
t.Fatal("timed out waiting for onOffline callback")
}
select {
case id := <-offline:
t.Fatalf("unexpected second onOffline callback for %q (d2 was never stale)", id)
case <-time.After(50 * time.Millisecond):
// expected: no further callbacks
}
}
func TestChecker_StopWaitsForLoopExit(t *testing.T) {
store := &fakeStore{known: nil}
c := New(store, time.Minute, time.Hour, nil, discardLogger())
c.Start()
c.Stop() // must return without hanging
}
@@ -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"})
@@ -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
}
@@ -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
}
@@ -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")
}
}