FireLeave_tool/reporter/reporter.go
2026-05-12 15:45:50 +08:00

871 lines
31 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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"
)
// SaveForbidTimeConfig 保存禁止时间配置
func SaveForbidTimeConfig(minutr string) error {
config := struct {
Minutr string `json:"minutr"`
}{
Minutr: minutr,
}
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)
}
// 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 {
// 读取缓存文件
data, err := os.ReadFile(ForbidTimeFile)
if err != nil {
// 缓存中没有值时默认返回0
return "0"
}
var config struct {
Minutr string `json:"minutr"`
}
if err := json.Unmarshal(data, &config); err != nil {
// 解析失败时默认返回0
return "0"
}
return config.Minutr
}
// 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)
}
// 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 := ""
if err != nil {
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)
os.Remove(mergedFile)
return
}
}
if !r.wsClient.IsConnected() {
logger.Printf("WS未就绪延迟上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
os.Remove(mergedFile)
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)
}
// 清理临时文件
logger.Printf("清理临时文件: %s (DeviceUUID=%s)", mergedFile, deviceData.DeviceUUID)
os.Remove(mergedFile)
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 := ""
if err != nil {
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)
os.Remove(mergedFile)
return
}
}
if !r.wsClient.IsConnected() {
logger.Printf("WS未就绪延迟上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
os.Remove(mergedFile)
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)
}
// 清理临时文件
logger.Printf("清理临时文件: %s (DeviceUUID=%s)", mergedFile, deviceData.DeviceUUID)
os.Remove(mergedFile)
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 := ""
if err != nil {
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)
os.Remove(mergedFile)
return
}
}
if !r.wsClient.IsConnected() {
logger.Printf("WS未就绪延迟上报 (DeviceUUID=%s)", deviceData.DeviceUUID)
os.Remove(mergedFile)
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)
}
// 清理临时文件
logger.Printf("清理临时文件: %s (DeviceUUID=%s)", mergedFile, deviceData.DeviceUUID)
os.Remove(mergedFile)
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
}
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()
}