124 lines
3.4 KiB
Go
124 lines
3.4 KiB
Go
// iot-service — монолит IoT на платформе Nubes (контейнер «Простой HTTP»).
|
|
//
|
|
// Три роли в одном процессе:
|
|
// - REST API + MQTT auth/acl (:9090, устройства в PostgreSQL);
|
|
// - bridge — MQTT-подписка EMQX → shared-SQS;
|
|
// - consumer — shared-SQS → per-tenant PostgreSQL.
|
|
//
|
|
// Наследие старого IoT (k8s, оператор, CRD) НЕ используется.
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"log/slog"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"sync"
|
|
"syscall"
|
|
"time"
|
|
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/api"
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/api/handler"
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/bridge"
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/config"
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/consumer"
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/sqsclient"
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/service/store"
|
|
"gitea.services.ngcloud.ru/Nail/IoT/internal/storage/iotpg"
|
|
)
|
|
|
|
// version — подставляется через ldflags при сборке.
|
|
var version = "dev"
|
|
|
|
func main() {
|
|
cfg, err := config.Load()
|
|
if err != nil {
|
|
slog.Error("config", "err", err)
|
|
os.Exit(1)
|
|
}
|
|
|
|
log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: cfg.LogLevel}))
|
|
slog.SetDefault(log)
|
|
|
|
log.Info("iot-service starting", "version", version, "port", cfg.APIPort)
|
|
|
|
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
|
|
defer cancel()
|
|
|
|
// PostgreSQL — устройства и телеметрия.
|
|
storeCtx, storeCancel := context.WithTimeout(ctx, 15*time.Second)
|
|
devices, err := store.Open(storeCtx, cfg.IOTPGDSN)
|
|
storeCancel()
|
|
if err != nil {
|
|
log.Error("devices store", "err", err)
|
|
os.Exit(1)
|
|
}
|
|
defer devices.Close()
|
|
log.Info("devices store ready")
|
|
|
|
iotStore, err := iotpg.New(cfg.IOTPGDSN, log)
|
|
if err != nil {
|
|
log.Error("iotpg store", "err", err)
|
|
os.Exit(1)
|
|
}
|
|
defer iotStore.Close()
|
|
log.Info("telemetry store ready")
|
|
|
|
sqsClient := sqsclient.New(cfg)
|
|
|
|
h := &handler.Handler{
|
|
Devices: devices,
|
|
IoTPG: iotStore,
|
|
SQS: sqsClient,
|
|
Cfg: cfg,
|
|
Log: log,
|
|
Version: version,
|
|
StartedAt: time.Now(),
|
|
}
|
|
|
|
// HTTP API — в основном потоке.
|
|
router := api.NewRouter(h, log, cfg.AuthTestMode, version)
|
|
server := &http.Server{
|
|
Addr: ":" + cfg.APIPort,
|
|
Handler: router,
|
|
ReadHeaderTimeout: 10 * time.Second,
|
|
}
|
|
|
|
// bridge + consumer — в фоне.
|
|
var wg sync.WaitGroup
|
|
wg.Add(2)
|
|
|
|
go func() {
|
|
defer wg.Done()
|
|
if err := bridge.Run(ctx, cfg, sqsClient, log); err != nil && !errors.Is(err, context.Canceled) {
|
|
log.Error("bridge stopped", "err", err)
|
|
}
|
|
}()
|
|
|
|
go func() {
|
|
defer wg.Done()
|
|
if err := consumer.Run(ctx, cfg, sqsClient, iotStore, log); err != nil && !errors.Is(err, context.Canceled) {
|
|
log.Error("consumer stopped", "err", err)
|
|
}
|
|
}()
|
|
|
|
// Сервер: при завершении ctx — гасим.
|
|
go func() {
|
|
<-ctx.Done()
|
|
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer shutdownCancel()
|
|
_ = server.Shutdown(shutdownCtx)
|
|
}()
|
|
|
|
log.Info("iot-service listening", "addr", server.Addr)
|
|
if err := server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
|
log.Error("http server", "err", err)
|
|
os.Exit(1)
|
|
}
|
|
|
|
wg.Wait()
|
|
log.Info("iot-service stopped")
|
|
}
|