Добавлен gRPC-контракт proto/device_control (TurnOn/TurnOff/SetLevel), общий Go-модуль proto/ для переиспользования сгенерированного кода. device-control-service принимает команды по gRPC, обновляет desired_state в Redis и публикует их в MQTT; слушает devices/+/telemetry и devices/+/ack для отметки живости устройства и обновления reported_state; горутина health-check переводит устройство в offline по таймауту и публикует событие в RabbitMQ (device.status_changed) для будущего notification-service. Проверено сквозным тестом через docker compose: grpcurl TurnOn/SetLevel меняет desired_state и уходит в MQTT, mosquitto_pub с ack обновляет reported_state и статус, health-check автоматически переводит устройство в offline по истечении таймаута.
120 lines
2.8 KiB
Go
120 lines
2.8 KiB
Go
package shadow
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"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)
|
|
}
|
|
|
|
func TestPatchDesiredState_MergesFields(t *testing.T) {
|
|
ctx := context.Background()
|
|
s := newTestStore(t)
|
|
|
|
if err := s.PatchDesiredState(ctx, "d1", map[string]any{"power": "on"}); err != nil {
|
|
t.Fatalf("first patch: %v", err)
|
|
}
|
|
if err := s.PatchDesiredState(ctx, "d1", map[string]any{"level": 42.0}); err != nil {
|
|
t.Fatalf("second patch: %v", err)
|
|
}
|
|
|
|
raw, err := s.rdb.Get(ctx, desiredKey("d1")).Result()
|
|
if err != nil {
|
|
t.Fatalf("get: %v", err)
|
|
}
|
|
var got map[string]any
|
|
if err := json.Unmarshal([]byte(raw), &got); err != nil {
|
|
t.Fatalf("unmarshal: %v", err)
|
|
}
|
|
if got["power"] != "on" || got["level"] != 42.0 {
|
|
t.Fatalf("got %+v, want power=on and level=42", got)
|
|
}
|
|
}
|
|
|
|
func TestTouch_FirstCallReportsOnlineTransition(t *testing.T) {
|
|
ctx := context.Background()
|
|
s := newTestStore(t)
|
|
|
|
becameOnline, err := s.Touch(ctx, "d1")
|
|
if err != nil {
|
|
t.Fatalf("touch: %v", err)
|
|
}
|
|
if !becameOnline {
|
|
t.Fatal("expected first touch to report an online transition")
|
|
}
|
|
|
|
becameOnline, err = s.Touch(ctx, "d1")
|
|
if err != nil {
|
|
t.Fatalf("second touch: %v", err)
|
|
}
|
|
if becameOnline {
|
|
t.Fatal("expected second touch (already online) to report no transition")
|
|
}
|
|
|
|
ids, err := s.KnownDevices(ctx)
|
|
if err != nil {
|
|
t.Fatalf("known devices: %v", err)
|
|
}
|
|
if len(ids) != 1 || ids[0] != "d1" {
|
|
t.Fatalf("got known devices %+v, want [d1]", ids)
|
|
}
|
|
}
|
|
|
|
func TestMarkOfflineIfStale(t *testing.T) {
|
|
ctx := context.Background()
|
|
s := newTestStore(t)
|
|
|
|
if _, err := s.Touch(ctx, "d1"); err != nil {
|
|
t.Fatalf("touch: %v", err)
|
|
}
|
|
|
|
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")
|
|
}
|
|
}
|