223 lines
6.8 KiB
Go
223 lines
6.8 KiB
Go
package main
|
||
|
||
import (
|
||
"context"
|
||
"log/slog"
|
||
"os"
|
||
"strconv"
|
||
"time"
|
||
|
||
"github.com/gin-gonic/gin"
|
||
"github.com/redis/go-redis/v9"
|
||
"gorm.io/gorm"
|
||
|
||
"silk-server-go/internal/config"
|
||
"silk-server-go/internal/database"
|
||
"silk-server-go/internal/handler"
|
||
"silk-server-go/internal/middleware"
|
||
"silk-server-go/internal/model"
|
||
"silk-server-go/internal/service"
|
||
"silk-server-go/internal/ws"
|
||
)
|
||
|
||
func main() {
|
||
// 1. 加载配置
|
||
cfg, err := config.Load()
|
||
if err != nil {
|
||
slog.Error("配置加载失败", "err", err)
|
||
os.Exit(1)
|
||
}
|
||
|
||
// 2. 连接数据库(GORM 自动迁移)
|
||
if err := database.Init(cfg); err != nil {
|
||
slog.Error("数据库连接失败", "err", err)
|
||
os.Exit(1)
|
||
}
|
||
db := database.DB
|
||
|
||
// 3. 初始化 IoTDB(失败则降级 PostgreSQL)
|
||
iotdb := service.NewIoTDBService(cfg.IoTDBURL)
|
||
if err := iotdb.Init(); err != nil {
|
||
slog.Warn("IoTDB 初始化失败,将降级使用 PostgreSQL", "err", err)
|
||
}
|
||
|
||
// 4. 连接 Redis(失败不阻断启动)
|
||
if opt, err := redis.ParseURL(cfg.Redis); err == nil {
|
||
rdb := redis.NewClient(opt)
|
||
if err := rdb.Ping(context.Background()).Err(); err != nil {
|
||
slog.Warn("Redis 连接失败", "err", err)
|
||
} else {
|
||
slog.Info("Redis 连接成功")
|
||
}
|
||
}
|
||
|
||
// 5. 创建 WebSocket Hub
|
||
hub := ws.NewHub(cfg.JWTSecret)
|
||
|
||
// 6. 创建并启动 MQTT 服务
|
||
mqttSvc := service.NewMQTTService(cfg.MQTT, db, iotdb, hub)
|
||
if err := mqttSvc.Start(); err != nil {
|
||
slog.Warn("MQTT 启动失败", "err", err)
|
||
}
|
||
defer mqttSvc.Stop()
|
||
|
||
// 启动设备离线检测(每 60 秒检查,超过 5 分钟未上报标记离线)
|
||
mqttSvc.StartOfflineChecker(5 * time.Minute)
|
||
|
||
// 7. 注入 MQTT publisher(用于控制命令下发)
|
||
handler.SetMQTTPublisher(mqttSvc)
|
||
handler.SetDeviceCommander(mqttSvc)
|
||
|
||
// 8. 创建 S3/Media/Transcode 服务
|
||
s3Svc := service.NewS3Service(cfg)
|
||
mediaSvc := service.NewMediaService(cfg)
|
||
transcodeSvc := service.NewTranscodeService(db, s3Svc)
|
||
aiSvc := service.NewAIClient(cfg.AIServiceBase)
|
||
wechatSvc := service.NewWechatService(cfg.WechatAppID, cfg.WechatSecret)
|
||
weatherSvc := service.NewWeatherService(cfg.QWeatherAPIKey, cfg.QWeatherLocation)
|
||
|
||
// 9. 创建 Gin 引擎
|
||
gin.SetMode(gin.ReleaseMode)
|
||
r := gin.New()
|
||
r.Use(middleware.Logger())
|
||
r.Use(middleware.SecurityHeaders())
|
||
r.Use(middleware.CORS())
|
||
|
||
// 10. WebSocket 路由(不经过 auth 中间件)
|
||
r.GET("/ws", hub.HandleWebSocket)
|
||
|
||
// 11. API 路由组(经过 JWT auth 中间件,白名单路径自动跳过)
|
||
api := r.Group("/api/v1")
|
||
api.Use(middleware.Auth(cfg))
|
||
|
||
// 公开路由(auth 白名单中跳过鉴权)
|
||
handler.RegisterHealthRoutes(api)
|
||
handler.RegisterVideoStreamRoutes(api, transcodeSvc, db, mediaSvc, cfg)
|
||
|
||
// JWT 保护的业务路由
|
||
handler.RegisterAuthRoutes(api, db, cfg)
|
||
handler.RegisterRoomRoutes(api, db)
|
||
handler.RegisterDeviceRoutes(api, db)
|
||
handler.RegisterSensorRoutes(api, db)
|
||
handler.RegisterThresholdRoutes(api, db)
|
||
handler.RegisterAlarmRoutes(api, db)
|
||
handler.RegisterAlarmClipRoutes(api, db)
|
||
handler.RegisterControlRoutes(api, db)
|
||
handler.RegisterNotificationRoutes(api, db)
|
||
handler.RegisterTelemetryRoutes(api, db, iotdb)
|
||
handler.RegisterVideoCameraRoutes(api, db, mediaSvc)
|
||
handler.RegisterVideoClipRoutes(api, db, cfg)
|
||
handler.RegisterVideoRecordRoutes(api, db, mediaSvc, cfg)
|
||
handler.RegisterStorageRoutes(api, db)
|
||
handler.RegisterKnowledgeRoutes(api, db, s3Svc, cfg.S3BucketImages)
|
||
handler.RegisterInspectionRoutes(api, db, s3Svc, aiSvc, cfg.S3BucketImages, wechatSvc, cfg.WechatTemplateInspection)
|
||
handler.RegisterTrayBatchRoutes(api, db)
|
||
handler.RegisterWechatRoutes(api, db, wechatSvc)
|
||
handler.RegisterWeatherRoutes(api, db, weatherSvc)
|
||
handler.RegisterLampRoutes(api, db, s3Svc, cfg.S3BucketImages)
|
||
handler.RegisterConsumableRoutes(api, db)
|
||
handler.RegisterConsultationRoutes(api, db)
|
||
handler.RegisterDetectionMethodRoutes(api, db)
|
||
|
||
// 启动高发病天气预警定时任务(未配置时跳过)
|
||
go startWeatherAlertLoop(db, weatherSvc, time.Duration(cfg.QWeatherIntervalMin)*time.Minute)
|
||
|
||
// 启动后台设备状态同步(每 30 秒查询 WVP 设备在线状态)
|
||
go startDeviceStatusSync(db, mediaSvc)
|
||
|
||
// 用户管理与审计日志路由(各路由内部按权限码校验)
|
||
handler.RegisterAuditRoutes(api, db)
|
||
handler.RegisterUserRoutes(api, db)
|
||
handler.RegisterPermissionRoutes(api, db)
|
||
|
||
// 12. 启动 HTTP 服务
|
||
addr := ":" + strconv.Itoa(cfg.Port)
|
||
slog.Info("🚀 Silk Go server 启动", "addr", addr, "apiPrefix", "/api/v1")
|
||
if err := r.Run(addr); err != nil {
|
||
slog.Error("服务启动失败", "err", err)
|
||
os.Exit(1)
|
||
}
|
||
}
|
||
|
||
// startDeviceStatusSync 定期从 WVP 同步设备信息到数据库(在线状态 + 共有参数)
|
||
func startDeviceStatusSync(db *gorm.DB, media *service.MediaService) {
|
||
ticker := time.NewTicker(30 * time.Second)
|
||
defer ticker.Stop()
|
||
|
||
for range ticker.C {
|
||
devMap, err := media.SyncWvpDevices()
|
||
if err != nil {
|
||
slog.Debug("同步 WVP 设备信息失败", "error", err)
|
||
continue
|
||
}
|
||
if len(devMap) == 0 {
|
||
continue
|
||
}
|
||
|
||
// 更新所有有 GB 设备 ID 的摄像头
|
||
var cameras []model.Camera
|
||
db.Where("gb_device_id IS NOT NULL AND gb_device_id != ''").Find(&cameras)
|
||
for _, cam := range cameras {
|
||
if cam.GbDeviceID == nil {
|
||
continue
|
||
}
|
||
info, exists := devMap[*cam.GbDeviceID]
|
||
if !exists {
|
||
continue
|
||
}
|
||
updates := map[string]interface{}{}
|
||
if cam.IsOnline != info.OnLine {
|
||
updates["is_online"] = info.OnLine
|
||
}
|
||
if info.Name != "" && cam.Name != info.Name {
|
||
updates["name"] = info.Name
|
||
}
|
||
if info.Manufacturer != "" {
|
||
if cam.GbManufacturer == nil || *cam.GbManufacturer != info.Manufacturer {
|
||
updates["gb_manufacturer"] = info.Manufacturer
|
||
}
|
||
}
|
||
if info.Password != "" {
|
||
if cam.GbAuthPassword == nil || *cam.GbAuthPassword != info.Password {
|
||
updates["gb_auth_password"] = info.Password
|
||
}
|
||
}
|
||
if len(updates) > 0 {
|
||
db.Model(&model.Camera{}).Where("id = ?", cam.ID).Updates(updates)
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// startWeatherAlertLoop 定期拉取天气并按规则写入 weather_alerts
|
||
func startWeatherAlertLoop(db *gorm.DB, weather *service.WeatherService, interval time.Duration) {
|
||
if !weather.Configured() {
|
||
slog.Warn("天气服务未配置(QWEATHER_API_KEY/QWEATHER_LOCATION),跳过定时预警")
|
||
return
|
||
}
|
||
run := func() {
|
||
ctx := context.Background()
|
||
now, err := weather.FetchNow(ctx)
|
||
if err != nil {
|
||
slog.Warn("天气拉取失败", "err", err)
|
||
return
|
||
}
|
||
daily, err := weather.FetchDaily(ctx)
|
||
if err != nil {
|
||
slog.Warn("天气预报拉取失败", "err", err)
|
||
return
|
||
}
|
||
for _, a := range service.EvaluateWeatherRules(*now, daily, "") {
|
||
if err := db.Create(&a).Error; err != nil {
|
||
slog.Warn("天气预警写入失败", "disease", a.Disease, "err", err)
|
||
}
|
||
}
|
||
}
|
||
run()
|
||
ticker := time.NewTicker(interval)
|
||
defer ticker.Stop()
|
||
for range ticker.C {
|
||
run()
|
||
}
|
||
}
|