device-control-service: HTTP/JSON API (internal/httpapi) рядом с gRPC —
POST /devices/{id}/turn-on|turn-off|set-level, тот же server.Server внутри,
без дублирования логики. Решение вместо gRPC-клиента на PHP: grpc/grpc
через PECL компилируется в Alpine 10-20+ минут и утяжеляет образ, а
диаграмма архитектуры в ТЗ и так допускала HTTP для Laravel→device-control.
Laravel: DeviceShadow (чтение Device Shadow из Redis, MGET одним запросом
для списка устройств), ClickHouseClient (HTTP-интерфейс ClickHouse,
параметризованные {name:Type}-запросы), DeviceControlClient (HTTP-вызовы
к новому Go-эндпоинту). Redis-клиент — predis (чистый PHP), а не phpredis,
по той же причине, что и решение по gRPC — не добавлять ещё одну
C-компиляцию в образ.
DeviceController::show — страница устройства: live-статус/last_seen/
desired-reported state из Redis, история показаний из ClickHouse (для
сенсоров), кнопки ручного управления (для актуаторов, только owner,
с проверкой capability устройства). В devices/index — бейдж online/
offline/unknown. Дашборд дополнен счётчиком онлайн-устройств.
Два реальных бага найдены и исправлены при сквозной проверке:
1. Пустой action_params сериализовался в JSON-массив "[]" (PHP не
различает пустой список и пустой объект), а Go ждёт объект —
rule-engine-service падал на unmarshal. Фикс — JsonObjectCast
(JSON_FORCE_OBJECT) на AutomationRule::action_params.
2. Redis-ключи device shadow — общее пространство имён с Go-сервisами
(сырые ключи без префикса), а Laravel по умолчанию добавляет ко всем
ключам префикс "app-name-database-" — Laravel никогда не видел
реальные данные. Фикс — REDIS_PREFIX="" в окружении контейнера
(важно: пустое значение в docker-compose YAML нужно задавать явно
через "", просто "KEY:" означает "взять из окружения хоста").
Проверено сквозным тестом через docker compose: полный цикл телеметрия →
правило → команда воспроизведён вживую с реальным исправлением на лету;
ручное управление (turn_on/turn_off/set_level) из Laravel UI подтверждено
через браузер — HTTP-вызов к device-control-service, обновление
desired_state (merge-patch), реальная MQTT-команда поймана мониторингом
топика. 38/38 тестов Laravel, все Go-тесты device-control-service зелёные.
196 lines
6.2 KiB
Go
196 lines
6.2 KiB
Go
// 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.
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log/slog"
|
|
"net"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"strings"
|
|
"syscall"
|
|
"time"
|
|
|
|
"google.golang.org/grpc"
|
|
|
|
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/mqttclient"
|
|
"git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/rabbitmq"
|
|
"git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/server"
|
|
"git.cactoz.su/cacto/home_automatization/services/device-control-service/internal/shadow"
|
|
)
|
|
|
|
func main() {
|
|
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
|
|
|
|
if err := run(logger); err != nil {
|
|
logger.Error("device-control-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)
|
|
shadowStore, err := shadow.Connect(ctx, cfg.RedisAddr, cfg.RedisPassword, cfg.RedisDB)
|
|
cancel()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer shadowStore.Close()
|
|
|
|
publisher, err := rabbitmq.Connect(cfg.RabbitMQURL)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer publisher.Close()
|
|
|
|
mqttClient, err := mqttclient.Connect(mqttclient.Config{
|
|
BrokerURL: cfg.MQTTBrokerURL,
|
|
ClientID: cfg.MQTTClientID,
|
|
}, logger)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer mqttClient.Close()
|
|
|
|
if err := mqttClient.Subscribe("devices/+/telemetry", 1, telemetryHandler(shadowStore, publisher, logger)); err != nil {
|
|
return fmt.Errorf("subscribe telemetry: %w", err)
|
|
}
|
|
if err := mqttClient.Subscribe("devices/+/ack", 1, ackHandler(shadowStore, publisher, logger)); err != nil {
|
|
return fmt.Errorf("subscribe ack: %w", err)
|
|
}
|
|
|
|
checker := healthcheck.New(shadowStore, cfg.HealthCheckTimeout, cfg.HealthCheckInterval, func(deviceID string) {
|
|
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))
|
|
if err != nil {
|
|
return fmt.Errorf("listen on grpc port %d: %w", cfg.GRPCPort, err)
|
|
}
|
|
grpcServer := grpc.NewServer()
|
|
devicecontrol.RegisterDeviceControlServer(grpcServer, srv)
|
|
|
|
// Same server instance, second transport — for callers where a full
|
|
// gRPC client is impractical (Laravel/PHP; see README).
|
|
httpServer := &http.Server{
|
|
Addr: fmt.Sprintf(":%d", cfg.HTTPPort),
|
|
Handler: httpapi.NewRouter(srv, logger),
|
|
}
|
|
|
|
go func() {
|
|
logger.Info("grpc server started", "grpc_port", cfg.GRPCPort)
|
|
if err := grpcServer.Serve(lis); err != nil {
|
|
logger.Error("grpc server stopped", "error", err)
|
|
}
|
|
}()
|
|
|
|
go func() {
|
|
logger.Info("http server started", "http_port", cfg.HTTPPort)
|
|
if err := httpServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
logger.Error("http server stopped", "error", err)
|
|
}
|
|
}()
|
|
|
|
stop := make(chan os.Signal, 1)
|
|
signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM)
|
|
<-stop
|
|
|
|
logger.Info("shutting down")
|
|
grpcServer.GracefulStop()
|
|
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if err := httpServer.Shutdown(shutdownCtx); err != nil {
|
|
logger.Error("http server shutdown failed", "error", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// deviceIDFromTopic extracts {device_id} from "devices/{device_id}/<suffix>".
|
|
func deviceIDFromTopic(topic string) (string, bool) {
|
|
parts := strings.Split(topic, "/")
|
|
if len(parts) != 3 || parts[0] != "devices" || parts[1] == "" {
|
|
return "", false
|
|
}
|
|
return parts[1], true
|
|
}
|
|
|
|
// telemetryHandler only tracks liveness: ingest-service owns storing and
|
|
// interpreting telemetry, this service just needs to know the device is alive.
|
|
func telemetryHandler(store *shadow.Store, publisher *rabbitmq.Publisher, logger *slog.Logger) mqttclient.Handler {
|
|
return func(topic string, _ []byte, _ time.Time) {
|
|
deviceID, ok := deviceIDFromTopic(topic)
|
|
if !ok {
|
|
logger.Warn("dropping telemetry from unparseable topic", "topic", topic)
|
|
return
|
|
}
|
|
touch(context.Background(), store, publisher, deviceID, logger)
|
|
}
|
|
}
|
|
|
|
type ackPayload struct {
|
|
State map[string]any `json:"state"`
|
|
}
|
|
|
|
func ackHandler(store *shadow.Store, publisher *rabbitmq.Publisher, logger *slog.Logger) mqttclient.Handler {
|
|
return func(topic string, payload []byte, _ time.Time) {
|
|
deviceID, ok := deviceIDFromTopic(topic)
|
|
if !ok {
|
|
logger.Warn("dropping ack from unparseable topic", "topic", topic)
|
|
return
|
|
}
|
|
|
|
var ack ackPayload
|
|
if err := json.Unmarshal(payload, &ack); err != nil {
|
|
logger.Warn("dropping invalid ack payload", "device_id", deviceID, "error", err)
|
|
return
|
|
}
|
|
|
|
ctx := context.Background()
|
|
if len(ack.State) > 0 {
|
|
if err := store.PatchReportedState(ctx, deviceID, ack.State); err != nil {
|
|
logger.Error("patch reported state failed", "device_id", deviceID, "error", err)
|
|
}
|
|
}
|
|
touch(ctx, store, publisher, deviceID, logger)
|
|
}
|
|
}
|
|
|
|
func touch(ctx context.Context, store *shadow.Store, publisher *rabbitmq.Publisher, deviceID string, logger *slog.Logger) {
|
|
becameOnline, err := store.Touch(ctx, deviceID)
|
|
if err != nil {
|
|
logger.Error("touch device failed", "device_id", deviceID, "error", err)
|
|
return
|
|
}
|
|
if becameOnline {
|
|
pubCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
|
defer cancel()
|
|
if err := publisher.PublishStatusChanged(pubCtx, deviceID, shadow.StatusOnline); err != nil {
|
|
logger.Error("publish online event failed", "device_id", deviceID, "error", err)
|
|
}
|
|
}
|
|
}
|