937 lines
33 KiB
Go
937 lines
33 KiB
Go
package reporter
|
||
|
||
import (
|
||
"encoding/json"
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"sync"
|
||
"time"
|
||
|
||
"FireLeave_tool/connect"
|
||
"FireLeave_tool/logger"
|
||
"FireLeave_tool/video_server"
|
||
)
|
||
|
||
var (
|
||
// mergingAlarms 用于合并报警,避免重复处理
|
||
mergingAlarms sync.Map
|
||
// 事件周期状态
|
||
eventCycleMutex sync.RWMutex
|
||
currentEventCycle EventCycle
|
||
cycleTimeout time.Time
|
||
)
|
||
|
||
// EventCycle 事件周期类型
|
||
type EventCycle string
|
||
|
||
const (
|
||
// EventCycleFireLeave 动火离人上报周期
|
||
EventCycleFireLeave EventCycle = "fire_leave"
|
||
// EventCycleFireCheck 火焰检测上报周期
|
||
EventCycleFireCheck EventCycle = "fire_check"
|
||
// 周期时限缓存文件
|
||
CycleTimeoutFile = "/data/cache/event_cycle_timeout.json"
|
||
)
|
||
|
||
// Reporter 上报管理器
|
||
type Reporter struct {
|
||
wsClient *connect.WSClient
|
||
uploadURL string
|
||
logger *Logger
|
||
}
|
||
|
||
// Logger 上报管理器专用日志
|
||
type Logger struct {
|
||
*log.Logger
|
||
}
|
||
|
||
// NewLogger 创建上报管理器专用日志
|
||
func NewLogger() *Logger {
|
||
return &Logger{
|
||
Logger: logger.GlobalLoggerManager.GetLogger("reporter"),
|
||
}
|
||
}
|
||
|
||
// NewReporter 创建上报管理器
|
||
func NewReporter() *Reporter {
|
||
// 初始化默认事件周期
|
||
eventCycleMutex.Lock()
|
||
currentEventCycle = EventCycleFireLeave
|
||
eventCycleMutex.Unlock()
|
||
|
||
return &Reporter{
|
||
wsClient: connect.GetWSClient(),
|
||
uploadURL: "",
|
||
logger: NewLogger(),
|
||
}
|
||
}
|
||
|
||
// getLogger 获取日志实例,确保不为nil
|
||
func (r *Reporter) getLogger() *log.Logger {
|
||
if r.logger != nil && r.logger.Logger != nil {
|
||
return r.logger.Logger
|
||
}
|
||
// 如果logger未初始化,使用全局Logger
|
||
return logger.Logger
|
||
}
|
||
|
||
// GetCurrentEventCycle 获取当前事件周期
|
||
func GetCurrentEventCycle() EventCycle {
|
||
eventCycleMutex.RLock()
|
||
defer eventCycleMutex.RUnlock()
|
||
return currentEventCycle
|
||
}
|
||
|
||
// IsCycleTimeout 检查事件周期是否超时
|
||
func IsCycleTimeout() bool {
|
||
eventCycleMutex.RLock()
|
||
defer eventCycleMutex.RUnlock()
|
||
return time.Now().After(cycleTimeout)
|
||
}
|
||
|
||
// SwitchEventCycle 切换事件周期
|
||
func SwitchEventCycle(cycle EventCycle) error {
|
||
eventCycleMutex.Lock()
|
||
defer eventCycleMutex.Unlock()
|
||
|
||
// 读取周期时限
|
||
timeout, err := loadCycleTimeout()
|
||
if err != nil {
|
||
// 使用默认值
|
||
timeout = 30 * time.Minute
|
||
}
|
||
|
||
// 更新事件周期和超时时间
|
||
oldCycle := currentEventCycle
|
||
currentEventCycle = cycle
|
||
cycleTimeout = time.Now().Add(timeout)
|
||
|
||
// 记录日志
|
||
if logger := logger.GlobalLoggerManager.GetLogger("reporter"); logger != nil {
|
||
logger.Printf("事件周期切换: %s → %s, 时限: %v", oldCycle, cycle, timeout)
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// CycleTimeoutConfig 周期时限配置
|
||
type CycleTimeoutConfig struct {
|
||
FireLeaveTimeout time.Duration `json:"fire_leave_timeout"`
|
||
FireCheckTimeout time.Duration `json:"fire_check_timeout"`
|
||
}
|
||
|
||
// loadCycleTimeout 加载周期时限
|
||
func loadCycleTimeout() (time.Duration, error) {
|
||
data, err := os.ReadFile(CycleTimeoutFile)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
|
||
var config CycleTimeoutConfig
|
||
if err := json.Unmarshal(data, &config); err != nil {
|
||
return 0, err
|
||
}
|
||
|
||
// 根据当前周期返回对应时限
|
||
switch currentEventCycle {
|
||
case EventCycleFireLeave:
|
||
return config.FireLeaveTimeout, nil
|
||
case EventCycleFireCheck:
|
||
return config.FireCheckTimeout, nil
|
||
default:
|
||
return 30 * time.Minute, nil
|
||
}
|
||
}
|
||
|
||
// SaveCycleTimeout 保存周期时限
|
||
func SaveCycleTimeout(config CycleTimeoutConfig) error {
|
||
data, err := json.Marshal(config)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 确保目录存在
|
||
if err := os.MkdirAll("/data/cache", 0755); err != nil {
|
||
return err
|
||
}
|
||
|
||
return os.WriteFile(CycleTimeoutFile, data, 0644)
|
||
}
|
||
|
||
// ForbidTimeConfig 禁止时间配置
|
||
var (
|
||
// ForbidTimeFile 禁止时间缓存文件
|
||
ForbidTimeFile = "/data/cache/forbid_time_config.json"
|
||
)
|
||
|
||
type ForbidTimeConfig struct {
|
||
Minute string `json:"minute"`
|
||
Minutr string `json:"minutr,omitempty"`
|
||
StartTime int64 `json:"start_time"`
|
||
EndTime int64 `json:"end_time"`
|
||
}
|
||
|
||
// SaveForbidTimeConfig 保存禁止时间配置
|
||
func SaveForbidTimeConfig(minute string) error {
|
||
return SaveForbidTimeWindow(minute, 0, 0)
|
||
}
|
||
|
||
func SaveForbidTimeWindow(minute string, startTime, endTime int64) error {
|
||
config := ForbidTimeConfig{
|
||
Minute: minute,
|
||
StartTime: startTime,
|
||
EndTime: endTime,
|
||
}
|
||
|
||
data, err := json.Marshal(config)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 确保目录存在
|
||
if err := os.MkdirAll("/data/cache", 0755); err != nil {
|
||
return err
|
||
}
|
||
|
||
return os.WriteFile(ForbidTimeFile, data, 0644)
|
||
}
|
||
|
||
func LoadForbidTimeConfig() (ForbidTimeConfig, error) {
|
||
data, err := os.ReadFile(ForbidTimeFile)
|
||
if err != nil {
|
||
return ForbidTimeConfig{}, err
|
||
}
|
||
|
||
var config ForbidTimeConfig
|
||
if err := json.Unmarshal(data, &config); err != nil {
|
||
return ForbidTimeConfig{}, err
|
||
}
|
||
if config.Minute == "" {
|
||
config.Minute = config.Minutr
|
||
}
|
||
return config, nil
|
||
}
|
||
|
||
// SaveEventCycleTimeoutConfig 保存事件周期时限配置
|
||
func SaveEventCycleTimeoutConfig(fireLeaveTimeout, fireCheckTimeout time.Duration) error {
|
||
config := CycleTimeoutConfig{
|
||
FireLeaveTimeout: fireLeaveTimeout,
|
||
FireCheckTimeout: fireCheckTimeout,
|
||
}
|
||
|
||
data, err := json.Marshal(config)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 确保目录存在
|
||
if err := os.MkdirAll("/data/cache", 0755); err != nil {
|
||
return err
|
||
}
|
||
|
||
return os.WriteFile(CycleTimeoutFile, data, 0644)
|
||
}
|
||
|
||
// GetForbidTimeConfig 获取禁止时间配置
|
||
func GetForbidTimeConfig() string {
|
||
config, err := LoadForbidTimeConfig()
|
||
if err != nil {
|
||
// 缓存中没有值时,默认返回0
|
||
return "0"
|
||
}
|
||
|
||
if config.Minute != "" {
|
||
return config.Minute
|
||
}
|
||
return config.Minutr
|
||
}
|
||
|
||
func GetForbidDuration(defaultDuration time.Duration) time.Duration {
|
||
return ParseForbidDuration(GetForbidTimeConfig(), defaultDuration)
|
||
}
|
||
|
||
func ParseForbidDuration(value string, defaultDuration time.Duration) time.Duration {
|
||
if value == "" || value == "0" {
|
||
return defaultDuration
|
||
}
|
||
if duration, err := time.ParseDuration(value); err == nil && duration > 0 {
|
||
return duration
|
||
}
|
||
var minutes int
|
||
if _, err := fmt.Sscanf(value, "%d", &minutes); err == nil && minutes > 0 {
|
||
return time.Duration(minutes) * time.Minute
|
||
}
|
||
return defaultDuration
|
||
}
|
||
|
||
// ShouldReportBothEvents 判断是否应该同时上报两种事件
|
||
func ShouldReportBothEvents() bool {
|
||
// 注意:此逻辑已被修改,现在只在特定情况下使用
|
||
// 正常情况下,应根据当前事件周期上报对应类型的事件
|
||
return false
|
||
}
|
||
|
||
// jsonString 将任意类型转换为JSON字符串
|
||
func jsonString(v interface{}) string {
|
||
b, err := json.Marshal(v)
|
||
if err != nil {
|
||
if logger := logger.GlobalLoggerManager.GetLogger("reporter"); logger != nil {
|
||
logger.Printf("JSON编码失败: %v", err)
|
||
}
|
||
return ""
|
||
}
|
||
return string(b)
|
||
}
|
||
|
||
func logJSON(logger *log.Logger, label string, v interface{}) {
|
||
if logger == nil {
|
||
return
|
||
}
|
||
b, err := json.Marshal(v)
|
||
if err != nil {
|
||
logger.Printf("%s JSON编码失败: %v", label, err)
|
||
return
|
||
}
|
||
logger.Printf("%s %s", label, string(b))
|
||
}
|
||
|
||
func cleanupMergedFileAfterUpload(logger *log.Logger, mergedFile, deviceUUID string, uploadFailed bool) {
|
||
if uploadFailed {
|
||
logger.Printf("Nginx上传失败,保留临时文件用于排查: %s (DeviceUUID=%s)", mergedFile, deviceUUID)
|
||
return
|
||
}
|
||
logger.Printf("清理临时文件: %s (DeviceUUID=%s)", mergedFile, deviceUUID)
|
||
if err := os.Remove(mergedFile); err != nil {
|
||
logger.Printf("清理临时文件失败: %v (DeviceUUID=%s, 文件=%s)", err, deviceUUID, mergedFile)
|
||
}
|
||
}
|
||
|
||
// UploadVideoAndReport 上传视频并上报事件
|
||
func (r *Reporter) UploadVideoAndReport(deviceData connect.DeviceData, alarmTime int64, mergedFile string) error {
|
||
// 异步执行视频上传和上报,避免阻塞主线程
|
||
go func() {
|
||
logger := r.getLogger()
|
||
logger.Printf("开始异步处理视频上传和事件上报 (DeviceUUID=%s, 文件=%s)", deviceData.DeviceUUID, mergedFile)
|
||
defer func() {
|
||
if err := recover(); err != nil {
|
||
logger.Printf("上传视频和上报事件发生 panic: %v", err)
|
||
// 确保临时文件被清理
|
||
os.Remove(mergedFile)
|
||
}
|
||
}()
|
||
|
||
// 检查事件周期是否超时
|
||
logger.Printf("检查事件周期状态 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if IsCycleTimeout() {
|
||
// 超时后切换回默认周期
|
||
logger.Printf("事件周期超时,切换回默认周期 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if err := SwitchEventCycle(EventCycleFireLeave); err != nil {
|
||
logger.Printf("事件周期超时切换失败: %v", err)
|
||
}
|
||
}
|
||
|
||
// 获取当前事件周期
|
||
currentCycle := GetCurrentEventCycle()
|
||
logger.Printf("当前事件周期: %s (DeviceUUID=%s)", currentCycle, deviceData.DeviceUUID)
|
||
|
||
// 检查是否应该同时上报两种事件
|
||
if ShouldReportBothEvents() {
|
||
logger.Printf("Forbid time config is 0, will report both event types (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
|
||
// 复制文件路径,避免在goroutine中使用同一个变量
|
||
fileToClean := mergedFile
|
||
|
||
// 上报fire_leave事件
|
||
logger.Printf("开始上报fire_leave事件 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
r.reportEventWithLevel(deviceData, alarmTime, mergedFile, ReportLevelFireLeave)
|
||
|
||
// 上报fire_check事件
|
||
logger.Printf("开始上报fire_check事件 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
r.reportEventWithLevel(deviceData, alarmTime, mergedFile, ReportLevelFireCheck)
|
||
|
||
// 延迟清理临时文件,确保上传完成
|
||
go func() {
|
||
// 等待一段时间,确保上传操作有足够时间完成
|
||
time.Sleep(5 * time.Second)
|
||
logger.Printf("清理临时文件: %s (DeviceUUID=%s)", fileToClean, deviceData.DeviceUUID)
|
||
os.Remove(fileToClean)
|
||
}()
|
||
|
||
logger.Printf("双事件上报完成 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
return
|
||
}
|
||
|
||
// 根据当前事件周期确定上报级别
|
||
var reportLevel ReportLevel
|
||
switch currentCycle {
|
||
case EventCycleFireLeave:
|
||
reportLevel = ReportLevelFireLeave
|
||
case EventCycleFireCheck:
|
||
reportLevel = ReportLevelFireCheck
|
||
default:
|
||
reportLevel = ReportLevelFireLeave
|
||
}
|
||
logger.Printf("确定上报级别: %s (DeviceUUID=%s)", reportLevel, deviceData.DeviceUUID)
|
||
|
||
// 生成Nginx上传参数
|
||
logger.Printf("生成Nginx上传参数 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, reportLevel)
|
||
|
||
// 上传视频到Nginx
|
||
logger.Printf("开始上传视频到Nginx (DeviceUUID=%s, 文件=%s, URL=%s)", deviceData.DeviceUUID, mergedFile, deviceData.UploadURL)
|
||
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
|
||
videoURLStr := ""
|
||
uploadFailed := false
|
||
if err != nil {
|
||
uploadFailed = true
|
||
logger.Printf("Nginx上传失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
videoURLStr = videoURL
|
||
logger.Printf("Nginx上传成功: %s (DeviceUUID=%s)", videoURL, deviceData.DeviceUUID)
|
||
}
|
||
|
||
// 检查WebSocket客户端
|
||
logger.Printf("检查WebSocket客户端连接状态 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端未初始化,尝试重新获取 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
// 尝试重新获取WebSocket客户端
|
||
r.wsClient = connect.GetWSClient()
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端仍然为nil,无法上报事件 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
return
|
||
}
|
||
}
|
||
|
||
if !r.wsClient.IsConnected() {
|
||
logger.Printf("WS未就绪,延迟上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
return
|
||
}
|
||
logger.Printf("WebSocket客户端连接状态正常 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
|
||
// 上报事件
|
||
logger.Printf("开始上报事件 (DeviceUUID=%s, Level=%s, VideoURL=%s)", deviceData.DeviceUUID, reportLevel, videoURLStr)
|
||
if err := r.ReportEvent(deviceData, alarmTime, videoURLStr, string(reportLevel)); err != nil {
|
||
logger.Printf("事件上报失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
logger.Printf("事件上报成功 (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, reportLevel)
|
||
}
|
||
|
||
// 上报数据快照
|
||
logger.Printf("开始上报数据快照 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if err := r.ReportDataSnapshot(deviceData, alarmTime); err != nil {
|
||
logger.Printf("数据快照上报失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
logger.Printf("数据快照上报成功 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
}
|
||
|
||
// 清理临时文件
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
logger.Printf("视频上传和事件上报完成 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
}()
|
||
|
||
return nil
|
||
}
|
||
|
||
// UploadVideoAndReportFireLeave 上传视频并上报温度报警事件(只上报fire_leave类型)
|
||
func (r *Reporter) UploadVideoAndReportFireLeave(deviceData connect.DeviceData, alarmTime int64, mergedFile string) error {
|
||
// 异步执行视频上传和上报,避免阻塞主线程
|
||
go func() {
|
||
logger := r.getLogger()
|
||
logger.Printf("开始异步处理温度报警视频上传和事件上报 (DeviceUUID=%s, 文件=%s)", deviceData.DeviceUUID, mergedFile)
|
||
defer func() {
|
||
if err := recover(); err != nil {
|
||
logger.Printf("上传温度报警视频和上报事件发生 panic: %v", err)
|
||
// 确保临时文件被清理
|
||
os.Remove(mergedFile)
|
||
}
|
||
}()
|
||
|
||
// 检查事件周期是否超时
|
||
logger.Printf("检查事件周期状态 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if IsCycleTimeout() {
|
||
// 超时后切换回默认周期
|
||
logger.Printf("事件周期超时,切换回默认周期 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if err := SwitchEventCycle(EventCycleFireLeave); err != nil {
|
||
logger.Printf("事件周期超时切换失败: %v", err)
|
||
}
|
||
}
|
||
|
||
// 生成Nginx上传参数
|
||
logger.Printf("生成Nginx上传参数 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, ReportLevelFireLeave)
|
||
|
||
// 上传视频到Nginx
|
||
logger.Printf("开始上传视频到Nginx (DeviceUUID=%s, 文件=%s, URL=%s)", deviceData.DeviceUUID, mergedFile, deviceData.UploadURL)
|
||
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
|
||
videoURLStr := ""
|
||
uploadFailed := false
|
||
if err != nil {
|
||
uploadFailed = true
|
||
logger.Printf("Nginx上传失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
videoURLStr = videoURL
|
||
logger.Printf("Nginx上传成功: %s (DeviceUUID=%s)", videoURL, deviceData.DeviceUUID)
|
||
}
|
||
|
||
// 检查WebSocket客户端
|
||
logger.Printf("检查WebSocket客户端连接状态 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端未初始化,尝试重新获取 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
// 尝试重新获取WebSocket客户端
|
||
r.wsClient = connect.GetWSClient()
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端仍然为nil,无法上报事件 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
return
|
||
}
|
||
}
|
||
|
||
if !r.wsClient.IsConnected() {
|
||
logger.Printf("WS未就绪,延迟上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
return
|
||
}
|
||
logger.Printf("WebSocket客户端连接状态正常 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
|
||
// 上报事件(只上报fire_leave类型)
|
||
logger.Printf("开始上报温度报警事件 (DeviceUUID=%s, Level=%s, VideoURL=%s)", deviceData.DeviceUUID, ReportLevelFireLeave, videoURLStr)
|
||
if err := r.ReportEvent(deviceData, alarmTime, videoURLStr, string(ReportLevelFireLeave)); err != nil {
|
||
logger.Printf("温度报警事件上报失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
logger.Printf("温度报警事件上报成功 (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, ReportLevelFireLeave)
|
||
}
|
||
|
||
// 上报数据快照
|
||
logger.Printf("开始上报数据快照 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if err := r.ReportDataSnapshot(deviceData, alarmTime); err != nil {
|
||
logger.Printf("数据快照上报失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
logger.Printf("数据快照上报成功 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
}
|
||
|
||
// 清理临时文件
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
logger.Printf("温度报警视频上传和事件上报完成 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
}()
|
||
|
||
return nil
|
||
}
|
||
|
||
// UploadVideoAndReportFireCheck 上传视频并上报火焰检测报警事件(只上报fire_check类型)
|
||
func (r *Reporter) UploadVideoAndReportFireCheck(deviceData connect.DeviceData, alarmTime int64, mergedFile string) error {
|
||
// 异步执行视频上传和上报,避免阻塞主线程
|
||
go func() {
|
||
logger := r.getLogger()
|
||
logger.Printf("开始异步处理火焰检测报警视频上传和事件上报 (DeviceUUID=%s, 文件=%s)", deviceData.DeviceUUID, mergedFile)
|
||
defer func() {
|
||
if err := recover(); err != nil {
|
||
logger.Printf("上传火焰检测报警视频和上报事件发生 panic: %v", err)
|
||
// 确保临时文件被清理
|
||
os.Remove(mergedFile)
|
||
}
|
||
}()
|
||
|
||
// 检查事件周期是否超时
|
||
logger.Printf("检查事件周期状态 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if IsCycleTimeout() {
|
||
// 超时后切换到fire_check周期
|
||
logger.Printf("事件周期超时,切换到fire_check周期 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if err := SwitchEventCycle(EventCycleFireCheck); err != nil {
|
||
logger.Printf("事件周期超时切换失败: %v", err)
|
||
}
|
||
}
|
||
|
||
// 生成Nginx上传参数
|
||
logger.Printf("生成Nginx上传参数 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, ReportLevelFireCheck)
|
||
|
||
// 上传视频到Nginx
|
||
logger.Printf("开始上传视频到Nginx (DeviceUUID=%s, 文件=%s, URL=%s)", deviceData.DeviceUUID, mergedFile, deviceData.UploadURL)
|
||
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
|
||
videoURLStr := ""
|
||
uploadFailed := false
|
||
if err != nil {
|
||
uploadFailed = true
|
||
logger.Printf("Nginx上传失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
videoURLStr = videoURL
|
||
logger.Printf("Nginx上传成功: %s (DeviceUUID=%s)", videoURL, deviceData.DeviceUUID)
|
||
}
|
||
|
||
// 检查WebSocket客户端
|
||
logger.Printf("检查WebSocket客户端连接状态 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端未初始化,尝试重新获取 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
// 尝试重新获取WebSocket客户端
|
||
r.wsClient = connect.GetWSClient()
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端仍然为nil,无法上报事件 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
return
|
||
}
|
||
}
|
||
|
||
if !r.wsClient.IsConnected() {
|
||
logger.Printf("WS未就绪,延迟上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
return
|
||
}
|
||
logger.Printf("WebSocket客户端连接状态正常 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
|
||
// 上报事件(只上报fire_check类型)
|
||
logger.Printf("开始上报火焰检测报警事件 (DeviceUUID=%s, Level=%s, VideoURL=%s)", deviceData.DeviceUUID, ReportLevelFireCheck, videoURLStr)
|
||
if err := r.ReportEvent(deviceData, alarmTime, videoURLStr, string(ReportLevelFireCheck)); err != nil {
|
||
logger.Printf("火焰检测报警事件上报失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
logger.Printf("火焰检测报警事件上报成功 (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, ReportLevelFireCheck)
|
||
}
|
||
|
||
// 上报数据快照
|
||
logger.Printf("开始上报数据快照 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
if err := r.ReportDataSnapshot(deviceData, alarmTime); err != nil {
|
||
logger.Printf("数据快照上报失败: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
|
||
} else {
|
||
logger.Printf("数据快照上报成功 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
}
|
||
|
||
// 清理临时文件
|
||
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, uploadFailed)
|
||
logger.Printf("火焰检测报警视频上传和事件上报完成 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
}()
|
||
|
||
return nil
|
||
}
|
||
|
||
// reportEventWithLevel 按指定级别上报事件
|
||
func (r *Reporter) reportEventWithLevel(deviceData connect.DeviceData, alarmTime int64, mergedFile string, reportLevel ReportLevel) error {
|
||
// 异步执行视频上传和上报,避免阻塞主线程
|
||
go func() {
|
||
logger := r.getLogger()
|
||
defer func() {
|
||
if err := recover(); err != nil {
|
||
logger.Printf("按指定级别上报事件发生 panic: %v", err)
|
||
}
|
||
}()
|
||
|
||
// 生成Nginx上传参数
|
||
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, reportLevel)
|
||
|
||
// 上传视频到Nginx
|
||
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
|
||
videoURLStr := ""
|
||
if err != nil {
|
||
logger.Printf("Nginx上传失败: %v", err)
|
||
} else {
|
||
videoURLStr = videoURL
|
||
logger.Printf("Nginx上传成功: %s", videoURL)
|
||
}
|
||
|
||
// 检查WebSocket客户端
|
||
if r.wsClient == nil || !r.wsClient.IsConnected() {
|
||
logger.Printf("WS未就绪,延迟上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
return
|
||
}
|
||
|
||
// 上报事件
|
||
if err := r.ReportEvent(deviceData, alarmTime, videoURLStr, string(reportLevel)); err != nil {
|
||
logger.Printf("事件上报失败: %v", err)
|
||
}
|
||
|
||
// 上报数据快照
|
||
if err := r.ReportDataSnapshot(deviceData, alarmTime); err != nil {
|
||
logger.Printf("数据快照上报失败: %v", err)
|
||
}
|
||
|
||
logger.Printf("按指定级别上报事件完成 (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, string(reportLevel))
|
||
}()
|
||
|
||
return nil
|
||
}
|
||
|
||
// ReportEvent 上报事件
|
||
func (r *Reporter) ReportEvent(deviceData connect.DeviceData, alarmTime int64, videoURL string, level string) error {
|
||
logger := r.getLogger()
|
||
|
||
// 检查WebSocket客户端
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端未初始化,无法上报事件 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
return fmt.Errorf("websocket client not initialized")
|
||
}
|
||
|
||
if !r.wsClient.IsConnected() {
|
||
logger.Printf("WebSocket客户端未连接,无法上报事件 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
return fmt.Errorf("websocket client not connected")
|
||
}
|
||
|
||
// 生成事件消息
|
||
eventMsg := r.GenerateEventMessage(deviceData, alarmTime, videoURL, level)
|
||
|
||
// 发送事件消息
|
||
r.wsClient.SendAsync(eventMsg)
|
||
logger.Printf("报警事件已上报 (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, level)
|
||
|
||
return nil
|
||
}
|
||
|
||
// ReportDataSnapshot 上报数据快照
|
||
func (r *Reporter) ReportDataSnapshot(deviceData connect.DeviceData, alarmTime int64) error {
|
||
logger := r.getLogger()
|
||
|
||
// 检查WebSocket客户端
|
||
if r.wsClient == nil {
|
||
logger.Printf("WebSocket客户端未初始化,无法上报数据快照 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
return fmt.Errorf("websocket client not initialized")
|
||
}
|
||
|
||
if !r.wsClient.IsConnected() {
|
||
logger.Printf("WebSocket客户端未连接,无法上报数据快照 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
return fmt.Errorf("websocket client not connected")
|
||
}
|
||
|
||
// 生成数据快照消息
|
||
dataMsg := r.GenerateDataMessage(deviceData, alarmTime)
|
||
|
||
// 发送数据快照消息
|
||
r.wsClient.SendAsync(dataMsg)
|
||
logger.Printf("数据快照已上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
|
||
|
||
return nil
|
||
}
|
||
|
||
// GenerateEventMessage 生成事件消息
|
||
func (r *Reporter) GenerateEventMessage(deviceData connect.DeviceData, alarmTime int64, videoURL string, level string) map[string]interface{} {
|
||
// 获取HostUUID
|
||
hostUUID := r.getHostUUID(deviceData.DeviceUUID)
|
||
|
||
// 生成描述
|
||
description := "无人工作区且温度超阈值报警"
|
||
if level == string(ReportLevelFireCheck) {
|
||
description = "检测到异常火苗"
|
||
}
|
||
|
||
// 构建事件负载
|
||
eventPayload := EventPayload{
|
||
Type: string(ReportTypeEvent),
|
||
Args: map[string]interface{}{
|
||
"level": level,
|
||
"description": description,
|
||
"device_uuid": deviceData.DeviceUUID,
|
||
"host_uuid": hostUUID,
|
||
"video_url": videoURL,
|
||
"alarm_time": alarmTime,
|
||
"address": deviceData.Address,
|
||
},
|
||
}
|
||
|
||
// 构建上报消息
|
||
eventMsg := map[string]interface{}{
|
||
"method": "metric_data",
|
||
"params": map[string]interface{}{
|
||
"route_key": fmt.Sprintf("/dhlr/device/%s/event", deviceData.DeviceUUID),
|
||
"metric": jsonString(eventPayload),
|
||
},
|
||
}
|
||
|
||
return eventMsg
|
||
}
|
||
|
||
// GenerateDataMessage 生成数据快照消息
|
||
func (r *Reporter) GenerateDataMessage(deviceData connect.DeviceData, alarmTime int64) map[string]interface{} {
|
||
// 构建数据快照负载
|
||
dataPayload := DataPayload{
|
||
Type: string(ReportTypeData),
|
||
Args: map[string]interface{}{
|
||
"person_count": deviceData.PersonCount,
|
||
"temperature": deviceData.Temperature,
|
||
"camera_rtsp": deviceData.CameraRTSP,
|
||
"task_id": deviceData.TaskID,
|
||
"camera_ip": deviceData.CameraIP,
|
||
"confidence": deviceData.Confidence,
|
||
"alarm_time": alarmTime,
|
||
"address": deviceData.Address,
|
||
},
|
||
}
|
||
|
||
// 构建上报消息
|
||
dataMsg := map[string]interface{}{
|
||
"method": "metric_data",
|
||
"params": map[string]interface{}{
|
||
"route_key": fmt.Sprintf("/dhlr/device/%s/data", deviceData.DeviceUUID),
|
||
"metric": jsonString(dataPayload),
|
||
},
|
||
}
|
||
|
||
return dataMsg
|
||
}
|
||
|
||
// GenerateNginxParams 生成Nginx上传参数
|
||
func (r *Reporter) GenerateNginxParams(deviceData connect.DeviceData, alarmTime int64, level ReportLevel) map[string]interface{} {
|
||
// 获取HostUUID
|
||
hostUUID := r.getHostUUID(deviceData.DeviceUUID)
|
||
|
||
// 生成描述
|
||
description := "无人工作区且温度超阈值报警"
|
||
if level == ReportLevelFireCheck {
|
||
description = "检测到异常火苗"
|
||
}
|
||
|
||
// 构建参数
|
||
params := map[string]interface{}{
|
||
"alarm_time": alarmTime,
|
||
"description": description,
|
||
"device_UUID": deviceData.DeviceUUID,
|
||
"host_uuid": hostUUID,
|
||
"person_count": deviceData.PersonCount,
|
||
"level": string(level),
|
||
}
|
||
|
||
// 火焰检测不需要温度数据
|
||
if level != ReportLevelFireCheck {
|
||
params["temperature"] = deviceData.Temperature
|
||
}
|
||
|
||
logJSON(r.getLogger(), fmt.Sprintf("[Nginx Upload Params JSON] DeviceUUID=%s Level=%s", deviceData.DeviceUUID, level), params)
|
||
|
||
return params
|
||
}
|
||
|
||
// TriggerAlarmWithDelay 延迟触发报警并上报(温度报警)
|
||
func (r *Reporter) TriggerAlarmWithDelay(deviceData connect.DeviceData, alarmTime int64, cooldown int) {
|
||
startTime := time.Now()
|
||
logger := r.getLogger()
|
||
logger.Printf("[%.3f] 开始处理温度报警事件 (DeviceUUID=%s, alarmTime=%d, cooldown=%d)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, alarmTime, cooldown)
|
||
|
||
// 捕获所有异常,确保函数不会因为panic而崩溃
|
||
defer func() {
|
||
if err := recover(); err != nil {
|
||
logger.Printf("[%.3f] 处理温度报警事件发生panic: %v (DeviceUUID=%s)", time.Since(startTime).Seconds(), err, deviceData.DeviceUUID)
|
||
// 确保锁被释放
|
||
key := fmt.Sprintf("%s_%d", deviceData.DeviceUUID, alarmTime)
|
||
mergingAlarms.Delete(key)
|
||
}
|
||
}()
|
||
|
||
// 防止重复处理
|
||
key := fmt.Sprintf("%s_%d", deviceData.DeviceUUID, alarmTime)
|
||
logger.Printf("[%.3f] 检查温度报警事件是否已在处理中 (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
|
||
if _, loaded := mergingAlarms.LoadOrStore(key, true); loaded {
|
||
logger.Printf("[%.3f] 温度报警事件已在处理中,跳过重复处理 (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
|
||
return
|
||
}
|
||
defer mergingAlarms.Delete(key)
|
||
logger.Printf("[%.3f] 成功获取处理锁 (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
|
||
|
||
// 合并视频
|
||
logger.Printf("[%.3f] 开始合并视频 (DeviceUUID=%s, alarmTime=%d)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, alarmTime)
|
||
mergeStartTime := time.Now()
|
||
mergedFile, err := video_server.MergeVideo(deviceData.DeviceUUID, alarmTime)
|
||
if err != nil {
|
||
logger.Printf("[%.3f] 视频合并失败 (DeviceUUID=%s, 耗时:%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds(), err)
|
||
return
|
||
}
|
||
logger.Printf("[%.3f] 视频合并成功,生成文件: %s (DeviceUUID=%s, 耗时:%.3fs)", time.Since(startTime).Seconds(), mergedFile, deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds())
|
||
|
||
// 上传视频并上报(只上报温度报警,即fire_leave类型)
|
||
logger.Printf("[%.3f] 开始上传视频并上报温度报警事件 (DeviceUUID=%s, 文件=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, mergedFile)
|
||
uploadStartTime := time.Now()
|
||
if err := r.UploadVideoAndReportFireLeave(deviceData, alarmTime, mergedFile); err != nil {
|
||
logger.Printf("[%.3f] 温度报警上报失败 (DeviceUUID=%s, 耗时:%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds(), err)
|
||
// 清理临时文件
|
||
if err := os.Remove(mergedFile); err != nil {
|
||
logger.Printf("[%.3f] 清理临时文件失败: %v (DeviceUUID=%s)", time.Since(startTime).Seconds(), err, deviceData.DeviceUUID)
|
||
}
|
||
return
|
||
}
|
||
logger.Printf("[%.3f] 温度报警上报成功 (DeviceUUID=%s, 耗时:%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds())
|
||
|
||
logger.Printf("[%.3f] 温度报警事件处理完成 (DeviceUUID=%s, 总耗时:%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(startTime).Seconds())
|
||
}
|
||
|
||
// TriggerFireCheckAlarm 触发火焰检测报警并上报
|
||
func (r *Reporter) TriggerFireCheckAlarm(deviceData connect.DeviceData, alarmTime int64) {
|
||
startTime := time.Now()
|
||
logger := r.getLogger()
|
||
logger.Printf("[%.3f] 开始处理火焰检测报警事件 (DeviceUUID=%s, alarmTime=%d)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, alarmTime)
|
||
|
||
// 捕获所有异常,确保函数不会因为panic而崩溃
|
||
defer func() {
|
||
if err := recover(); err != nil {
|
||
logger.Printf("[%.3f] 处理火焰检测报警事件发生panic: %v (DeviceUUID=%s)", time.Since(startTime).Seconds(), err, deviceData.DeviceUUID)
|
||
// 确保锁被释放
|
||
key := fmt.Sprintf("%s_fire_%d", deviceData.DeviceUUID, alarmTime)
|
||
mergingAlarms.Delete(key)
|
||
}
|
||
}()
|
||
|
||
// 防止重复处理
|
||
key := fmt.Sprintf("%s_fire_%d", deviceData.DeviceUUID, alarmTime)
|
||
logger.Printf("[%.3f] 检查火焰检测报警事件是否已在处理中 (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
|
||
if _, loaded := mergingAlarms.LoadOrStore(key, true); loaded {
|
||
logger.Printf("[%.3f] 火焰检测报警事件已在处理中,跳过重复处理 (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
|
||
return
|
||
}
|
||
defer mergingAlarms.Delete(key)
|
||
logger.Printf("[%.3f] 成功获取处理锁 (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
|
||
|
||
// 合并视频(使用fire_check目录的视频)
|
||
logger.Printf("[%.3f] 开始合并火焰检测视频 (DeviceUUID=%s, alarmTime=%d)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, alarmTime)
|
||
mergeStartTime := time.Now()
|
||
// 这里可以修改MergeVideo函数,使其支持指定目录,或者创建一个新的函数用于合并fire_check视频
|
||
mergedFile, err := video_server.MergeVideo(deviceData.DeviceUUID+"_fire_check", alarmTime)
|
||
if err != nil {
|
||
logger.Printf("[%.3f] 火焰检测视频合并失败 (DeviceUUID=%s, 耗时:%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds(), err)
|
||
return
|
||
}
|
||
logger.Printf("[%.3f] 火焰检测视频合并成功,生成文件: %s (DeviceUUID=%s, 耗时:%.3fs)", time.Since(startTime).Seconds(), mergedFile, deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds())
|
||
|
||
// 上传视频并上报(只上报火焰检测,即fire_check类型)
|
||
logger.Printf("[%.3f] 开始上传视频并上报火焰检测报警事件 (DeviceUUID=%s, 文件=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, mergedFile)
|
||
uploadStartTime := time.Now()
|
||
if err := r.UploadVideoAndReportFireCheck(deviceData, alarmTime, mergedFile); err != nil {
|
||
logger.Printf("[%.3f] 火焰检测报警上报失败 (DeviceUUID=%s, 耗时:%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds(), err)
|
||
// 清理临时文件
|
||
if err := os.Remove(mergedFile); err != nil {
|
||
logger.Printf("[%.3f] 清理临时文件失败: %v (DeviceUUID=%s)", time.Since(startTime).Seconds(), err, deviceData.DeviceUUID)
|
||
}
|
||
return
|
||
}
|
||
logger.Printf("[%.3f] 火焰检测报警上报成功 (DeviceUUID=%s, 耗时:%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds())
|
||
|
||
logger.Printf("[%.3f] 火焰检测报警事件处理完成 (DeviceUUID=%s, 总耗时:%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(startTime).Seconds())
|
||
}
|
||
|
||
// getHostUUID 获取设备的HostUUID
|
||
func (r *Reporter) getHostUUID(deviceUUID string) string {
|
||
configs, err := connect.LoadServiceConfig()
|
||
if err != nil {
|
||
r.logger.Printf("加载服务配置失败: %v", err)
|
||
return ""
|
||
}
|
||
|
||
for _, cfg := range configs {
|
||
if cfg.DeviceUUID == deviceUUID {
|
||
return cfg.HostUUID
|
||
}
|
||
}
|
||
|
||
return ""
|
||
}
|
||
|
||
// SetWSClient 设置WebSocket客户端
|
||
func (r *Reporter) SetWSClient(client *connect.WSClient) {
|
||
r.wsClient = client
|
||
}
|
||
|
||
// SetUploadURL 设置上传URL
|
||
func (r *Reporter) SetUploadURL(url string) {
|
||
r.uploadURL = url
|
||
}
|
||
|
||
// IsWSReady 检查WebSocket是否就绪
|
||
func (r *Reporter) IsWSReady() bool {
|
||
return r.wsClient != nil && r.wsClient.IsConnected()
|
||
}
|