FireLeave_tool/reporter/reporter.go

778 lines
27 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 sync.Map
eventCycleMutex sync.RWMutex
currentEventCycle EventCycle
cycleTimeout time.Time
)
type EventCycle string
const (
EventCycleFireLeave EventCycle = "fire_leave"
EventCycleFireCheck EventCycle = "fire_check"
CycleTimeoutFile = "/data/cache/event_cycle_timeout.json"
)
type Reporter struct {
wsClient *connect.WSClient
uploadURL string
logger *Logger
}
type Logger struct {
*log.Logger
}
func NewLogger() *Logger {
return &Logger{
Logger: logger.GlobalLoggerManager.GetLogger("reporter"),
}
}
func NewReporter() *Reporter {
eventCycleMutex.Lock()
currentEventCycle = EventCycleFireLeave
eventCycleMutex.Unlock()
return &Reporter{
wsClient: connect.GetWSClient(),
uploadURL: "",
logger: NewLogger(),
}
}
func (r *Reporter) getLogger() *log.Logger {
if r.logger != nil && r.logger.Logger != nil {
return r.logger.Logger
}
return logger.Logger
}
func GetCurrentEventCycle() EventCycle {
eventCycleMutex.RLock()
defer eventCycleMutex.RUnlock()
return currentEventCycle
}
func IsCycleTimeout() bool {
eventCycleMutex.RLock()
defer eventCycleMutex.RUnlock()
return time.Now().After(cycleTimeout)
}
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("Event cycle switched: %s -> %s, timeout=%v", oldCycle, cycle, timeout)
}
return nil
}
type CycleTimeoutConfig struct {
FireLeaveTimeout time.Duration `json:"fire_leave_timeout"`
FireCheckTimeout time.Duration `json:"fire_check_timeout"`
}
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
}
}
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)
}
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
}
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 {
// 缂傛挸鐡ㄦ稉顓熺梾閺堝鈧吋妞傞敍宀勭帛鐠併倛绻戦崶?
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
}
func ShouldReportBothEvents() bool {
return false
}
// jsonString
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("Upload failed, keep temp merged file for troubleshooting: %s (DeviceUUID=%s)", mergedFile, deviceUUID)
return
}
logger.Printf("Remove temp merged file: %s (DeviceUUID=%s)", mergedFile, deviceUUID)
if err := os.Remove(mergedFile); err != nil {
logger.Printf("Remove temp merged file failed: %v (DeviceUUID=%s, file=%s)", err, deviceUUID, mergedFile)
}
}
func (r *Reporter) UploadVideoAndReport(deviceData connect.DeviceData, alarmTime int64, mergedFile string) error {
go func() {
logger := r.getLogger()
defer func() {
if recovered := recover(); recovered != nil {
logger.Printf("upload/report panic: %v (DeviceUUID=%s)", recovered, deviceData.DeviceUUID)
_ = os.Remove(mergedFile)
}
}()
if IsCycleTimeout() {
if err := SwitchEventCycle(EventCycleFireLeave); err != nil {
logger.Printf("switch event cycle failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
}
}
currentCycle := GetCurrentEventCycle()
var reportLevel ReportLevel
switch currentCycle {
case EventCycleFireLeave:
reportLevel = ReportLevelFireLeave
case EventCycleFireCheck:
reportLevel = ReportLevelFireCheck
default:
reportLevel = ReportLevelFireLeave
}
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, reportLevel)
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
if err != nil {
logger.Printf("upload video failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, true)
return
}
if r.wsClient == nil {
r.wsClient = connect.GetWSClient()
}
if r.wsClient == nil || !r.wsClient.IsConnected() {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
return
}
if err := r.ReportEventSync(deviceData, alarmTime, videoURL, string(reportLevel)); err != nil {
logger.Printf("report event failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
return
}
if err := r.ReportDataSnapshotSync(deviceData, alarmTime); err != nil {
logger.Printf("report snapshot failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
return
}
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
}()
return nil
}
// UploadVideoAndReportFireLeave uploads and reports a fire_leave event asynchronously.
func (r *Reporter) UploadVideoAndReportFireLeave(deviceData connect.DeviceData, alarmTime int64, mergedFile string) error {
go func() {
logger := r.getLogger()
defer func() {
if recovered := recover(); recovered != nil {
logger.Printf("fire_leave upload/report panic: %v (DeviceUUID=%s)", recovered, deviceData.DeviceUUID)
_ = os.Remove(mergedFile)
}
}()
if IsCycleTimeout() {
if err := SwitchEventCycle(EventCycleFireLeave); err != nil {
logger.Printf("switch event cycle failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
}
}
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, ReportLevelFireLeave)
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
if err != nil {
logger.Printf("upload fire_leave video failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, true)
return
}
if r.wsClient == nil {
r.wsClient = connect.GetWSClient()
}
if r.wsClient == nil || !r.wsClient.IsConnected() {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
return
}
if err := r.ReportEventSync(deviceData, alarmTime, videoURL, string(ReportLevelFireLeave)); err != nil {
logger.Printf("report fire_leave event failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
return
}
if err := r.ReportDataSnapshotSync(deviceData, alarmTime); err != nil {
logger.Printf("report fire_leave snapshot failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
return
}
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
}()
return nil
}
// UploadVideoAndReportFireCheck uploads and reports a fire_check event synchronously.
// UploadVideoAndReportFireCheck uploads and reports a fire_check event synchronously.
func (r *Reporter) UploadVideoAndReportFireCheck(deviceData connect.DeviceData, alarmTime int64, mergedFile string) (err error) {
logger := r.getLogger()
logger.Printf("Start fire_check upload/report (DeviceUUID=%s, file=%s)", deviceData.DeviceUUID, mergedFile)
defer func() {
if recovered := recover(); recovered != nil {
logger.Printf("fire_check upload/report panic: %v (DeviceUUID=%s)", recovered, deviceData.DeviceUUID)
err = fmt.Errorf("fire_check upload/report panic: %v", recovered)
}
}()
if IsCycleTimeout() {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
if err := SwitchEventCycle(EventCycleFireCheck); err != nil {
logger.Printf("Switch event cycle to fire_check failed: %v", err)
}
}
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, ReportLevelFireCheck)
ensureNginxReportLevel(nginxParams, ReportLevelFireCheck)
logger.Printf("Upload fire_check video to Nginx (DeviceUUID=%s, file=%s, URL=%s)", deviceData.DeviceUUID, mergedFile, deviceData.UploadURL)
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
if err != nil {
logger.Printf("Upload fire_check video to Nginx failed: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
return fmt.Errorf("upload fire_check video failed: %w", err)
}
logger.Printf("Upload fire_check video to Nginx succeeded: %s (DeviceUUID=%s)", videoURL, deviceData.DeviceUUID)
if r.wsClient == nil {
r.wsClient = connect.GetWSClient()
}
if r.wsClient == nil {
return fmt.Errorf("websocket client not initialized")
}
if !r.wsClient.IsConnected() {
return fmt.Errorf("websocket client not connected")
}
if err := r.ReportEventSync(deviceData, alarmTime, videoURL, string(ReportLevelFireCheck)); err != nil {
return err
}
if err := r.ReportDataSnapshotSync(deviceData, alarmTime); err != nil {
return err
}
cleanupMergedFileAfterUpload(logger, mergedFile, deviceData.DeviceUUID, false)
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return nil
}
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)
}
}()
nginxParams := r.GenerateNginxParams(deviceData, alarmTime, reportLevel)
videoURL, err := video_server.UploadToNginx(mergedFile, nginxParams, deviceData.UploadURL)
videoURLStr := ""
if err != nil {
logger.Printf("Nginx涓婁紶澶辫触: %v", err)
} else {
videoURLStr = videoURL
logger.Printf("Nginx upload succeeded: %s", videoURL)
}
if r.wsClient == nil || !r.wsClient.IsConnected() {
logger.Printf("Report data snapshot sync succeeded (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("Report event completed (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, string(reportLevel))
}()
return nil
}
func (r *Reporter) ReportEvent(deviceData connect.DeviceData, alarmTime int64, videoURL string, level string) error {
logger := r.getLogger()
if r.wsClient == nil {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not initialized")
}
if !r.wsClient.IsConnected() {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not connected")
}
eventMsg := r.GenerateEventMessage(deviceData, alarmTime, videoURL, level)
r.wsClient.SendAsync(eventMsg)
logger.Printf("Report event sync succeeded (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, level)
return nil
}
func (r *Reporter) ReportEventSync(deviceData connect.DeviceData, alarmTime int64, videoURL string, level string) error {
logger := r.getLogger()
if r.wsClient == nil {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not initialized")
}
if !r.wsClient.IsConnected() {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not connected")
}
if err := r.wsClient.SendRawMessage(r.GenerateEventMessage(deviceData, alarmTime, videoURL, level)); err != nil {
logger.Printf("浜嬩欢鍚屾涓婃姤澶辫触: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
return err
}
logger.Printf("Report event sync succeeded (DeviceUUID=%s, Level=%s)", deviceData.DeviceUUID, level)
return nil
}
func (r *Reporter) ReportDataSnapshot(deviceData connect.DeviceData, alarmTime int64) error {
logger := r.getLogger()
if r.wsClient == nil {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not initialized")
}
if !r.wsClient.IsConnected() {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not connected")
}
dataMsg := r.GenerateDataMessage(deviceData, alarmTime)
r.wsClient.SendAsync(dataMsg)
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return nil
}
func (r *Reporter) ReportDataSnapshotSync(deviceData connect.DeviceData, alarmTime int64) error {
logger := r.getLogger()
if r.wsClient == nil {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not initialized")
}
if !r.wsClient.IsConnected() {
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return fmt.Errorf("websocket client not connected")
}
if err := r.wsClient.SendRawMessage(r.GenerateDataMessage(deviceData, alarmTime)); err != nil {
logger.Printf("鏁版嵁蹇収鍚屾涓婃姤澶辫触: %v (DeviceUUID=%s)", err, deviceData.DeviceUUID)
return err
}
logger.Printf("Report data snapshot sync succeeded (DeviceUUID=%s)", deviceData.DeviceUUID)
return nil
}
func (r *Reporter) GenerateEventMessage(deviceData connect.DeviceData, alarmTime int64, videoURL string, level string) map[string]interface{} {
hostUUID := r.getHostUUID(deviceData.DeviceUUID)
description := "No person in work area and temperature over threshold"
if level == string(ReportLevelFireCheck) {
description = "Abnormal fire detected"
}
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,
},
}
return map[string]interface{}{
"method": "metric_data",
"params": map[string]interface{}{
"route_key": fmt.Sprintf("/dhlr/device/%s/event", deviceData.DeviceUUID),
"metric": jsonString(eventPayload),
},
}
}
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,
},
}
return map[string]interface{}{
"method": "metric_data",
"params": map[string]interface{}{
"route_key": fmt.Sprintf("/dhlr/device/%s/data", deviceData.DeviceUUID),
"metric": jsonString(dataPayload),
},
}
}
func (r *Reporter) GenerateNginxParams(deviceData connect.DeviceData, alarmTime int64, level ReportLevel) map[string]interface{} {
hostUUID := r.getHostUUID(deviceData.DeviceUUID)
description := "No person in work area and temperature over threshold"
if level == ReportLevelFireCheck {
description = "Abnormal fire detected"
}
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
}
func ensureNginxReportLevel(params map[string]interface{}, level ReportLevel) {
params["level"] = string(level)
if _, ok := params["description"]; !ok {
if level == ReportLevelFireCheck {
params["description"] = "Abnormal fire detected"
} else {
params["description"] = "No person in work area and temperature over threshold"
}
}
}
// TriggerAlarmWithDelay 瀵ゆ儼绻滅憴锕€褰傞幎銉劅楠炴湹绗傞幎銉礄濞撯晛瀹抽幎銉劅閿?// TriggerAlarmWithDelay handles a fire_leave temperature alarm.
func (r *Reporter) TriggerAlarmWithDelay(deviceData connect.DeviceData, alarmTime int64, cooldown int) {
startTime := time.Now()
logger := r.getLogger()
logger.Printf("[%.3f] Start fire_leave alarm handling (DeviceUUID=%s, alarmTime=%d, cooldown=%d)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, alarmTime, cooldown)
defer func() {
if err := recover(); err != nil {
logger.Printf("[%.3f] fire_leave alarm handling 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] Check fire_leave alarm lock (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
if _, loaded := mergingAlarms.LoadOrStore(key, true); loaded {
logger.Printf("[%.3f] fire_leave alarm already handling, skip duplicate (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
return
}
defer mergingAlarms.Delete(key)
logger.Printf("[%.3f] fire_leave alarm lock acquired (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
logger.Printf("[%.3f] Start merging fire_leave video (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] Merge fire_leave video failed (DeviceUUID=%s, elapsed=%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds(), err)
return
}
logger.Printf("[%.3f] Merge fire_leave video succeeded: %s (DeviceUUID=%s, elapsed=%.3fs)", time.Since(startTime).Seconds(), mergedFile, deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds())
logger.Printf("[%.3f] Start upload/report for fire_leave alarm (DeviceUUID=%s, file=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, mergedFile)
uploadStartTime := time.Now()
if err := r.UploadVideoAndReportFireLeave(deviceData, alarmTime, mergedFile); err != nil {
logger.Printf("[%.3f] fire_leave alarm upload/report failed (DeviceUUID=%s, elapsed=%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds(), err)
if err := os.Remove(mergedFile); err != nil {
logger.Printf("[%.3f] Remove merged file after fire_leave failure failed: %v (DeviceUUID=%s)", time.Since(startTime).Seconds(), err, deviceData.DeviceUUID)
}
return
}
logger.Printf("[%.3f] fire_leave alarm upload/report succeeded (DeviceUUID=%s, elapsed=%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds())
logger.Printf("[%.3f] fire_leave alarm handling completed (DeviceUUID=%s, total=%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(startTime).Seconds())
}
func (r *Reporter) TriggerFireCheckAlarm(deviceData connect.DeviceData, alarmTime int64) bool {
startTime := time.Now()
logger := r.getLogger()
logger.Printf("[%.3f] Start fire_check alarm handling (DeviceUUID=%s, alarmTime=%d)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, alarmTime)
defer func() {
if err := recover(); err != nil {
logger.Printf("[%.3f] fire_check alarm handling panic: %v (DeviceUUID=%s)", time.Since(startTime).Seconds(), err, deviceData.DeviceUUID)
key := fmt.Sprintf("%s_fire", deviceData.DeviceUUID)
mergingAlarms.Delete(key)
}
}()
key := fmt.Sprintf("%s_fire", deviceData.DeviceUUID)
logger.Printf("[%.3f] Check fire_check alarm lock (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
if _, loaded := mergingAlarms.LoadOrStore(key, true); loaded {
logger.Printf("[%.3f] fire_check alarm already handling, skip duplicate (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
return false
}
defer mergingAlarms.Delete(key)
logger.Printf("[%.3f] fire_check alarm lock acquired (DeviceUUID=%s, key=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, key)
const fireCheckVideoFlushDelay = 21 * time.Second
logger.Printf("[%.3f] Wait %v for fire_check video cache flush before merge (DeviceUUID=%s)",
time.Since(startTime).Seconds(), fireCheckVideoFlushDelay, deviceData.DeviceUUID)
time.Sleep(fireCheckVideoFlushDelay)
logger.Printf("[%.3f] Start merging fire_check video (DeviceUUID=%s, alarmTime=%d)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, alarmTime)
mergeStartTime := time.Now()
mergedFile, err := video_server.MergeVideoFromDirectory(
deviceData.DeviceUUID+"_fire_check",
deviceData.DeviceUUID,
alarmTime,
)
if err != nil {
logger.Printf("[%.3f] Merge fire_check video failed (DeviceUUID=%s, elapsed=%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds(), err)
dumpFireCheckVideoDirectory(logger, deviceData.DeviceUUID)
return false
}
logger.Printf("[%.3f] Merge fire_check video succeeded: %s (DeviceUUID=%s, elapsed=%.3fs)", time.Since(startTime).Seconds(), mergedFile, deviceData.DeviceUUID, time.Since(mergeStartTime).Seconds())
logger.Printf("[%.3f] Start upload/report for fire_check alarm (DeviceUUID=%s, file=%s)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, mergedFile)
uploadStartTime := time.Now()
if err := r.UploadVideoAndReportFireCheck(deviceData, alarmTime, mergedFile); err != nil {
logger.Printf("[%.3f] fire_check alarm upload/report failed (DeviceUUID=%s, elapsed=%.3fs): %v", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds(), err)
if err := os.Remove(mergedFile); err != nil {
logger.Printf("[%.3f] Remove merged file after fire_check failure failed: %v (DeviceUUID=%s)", time.Since(startTime).Seconds(), err, deviceData.DeviceUUID)
}
return false
}
logger.Printf("[%.3f] fire_check alarm upload/report succeeded (DeviceUUID=%s, elapsed=%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(uploadStartTime).Seconds())
logger.Printf("[%.3f] fire_check alarm handling completed (DeviceUUID=%s, total=%.3fs)", time.Since(startTime).Seconds(), deviceData.DeviceUUID, time.Since(startTime).Seconds())
return true
}
func dumpFireCheckVideoDirectory(logger *log.Logger, deviceUUID string) {
videoDir := fmt.Sprintf("/usr/data/camera/%s_fire_check", deviceUUID)
files, err := os.ReadDir(videoDir)
if err != nil {
logger.Printf("fire_check video directory snapshot failed: dir=%s, err=%v", videoDir, err)
return
}
logger.Printf("fire_check video directory snapshot: dir=%s, files=%d", videoDir, len(files))
for _, file := range files {
info, err := file.Info()
if err != nil {
logger.Printf(" file=%s, stat_error=%v", file.Name(), err)
continue
}
if file.IsDir() {
logger.Printf(" dir=%s, mod=%s", file.Name(), info.ModTime().Format(time.RFC3339))
continue
}
logger.Printf(" file=%s, size=%d, mod=%s", file.Name(), info.Size(), info.ModTime().Format(time.RFC3339))
}
}
func (r *Reporter) getHostUUID(deviceUUID string) string {
configs, err := connect.LoadServiceConfig()
if err != nil {
r.logger.Printf("Load service config failed: %v", err)
return ""
}
for _, cfg := range configs {
if cfg.DeviceUUID == deviceUUID {
return cfg.HostUUID
}
}
return ""
}
// SetWSClient
func (r *Reporter) SetWSClient(client *connect.WSClient) {
r.wsClient = client
}
// SetUploadURL 鐠佸墽鐤嗘稉濠佺炊URL
func (r *Reporter) SetUploadURL(url string) {
r.uploadURL = url
}
// IsWSReady 濡偓閺岊櫇ebSocket閺勵垰鎯佺亸杈╁崕
func (r *Reporter) IsWSReady() bool {
return r.wsClient != nil && r.wsClient.IsConnected()
}