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, db) 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) handler.RegisterTraceRoutes(api, db) handler.RegisterHealthRoutes(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() } }