Browse Source

杨邦控制卡dome。 雷达、语音合成模块、redis解释摄像头车牌号数据

xu 7 tháng trước cách đây
mục cha
commit
1b328a7c5e
15 tập tin đã thay đổi với 1068 bổ sung373 xóa
  1. 1 1
      bx/BxAreaDynamic.go
  2. 43 31
      lc/IDevice.go
  3. 27 31
      lc/camera_event.go
  4. 69 0
      lc/command_queue.go
  5. 1 3
      lc/model/camera.go
  6. 1 3
      lc/model/radar.go
  7. 3 4
      lc/model/screen.go
  8. 1 6
      lc/model/speaker.go
  9. 218 35
      lc/radar_event.go
  10. 135 0
      lc/redis_process.go
  11. 159 58
      lc/screen.go
  12. 111 126
      lc/server.go
  13. 223 48
      lc/speaker.go
  14. 16 11
      util/config.go
  15. 60 16
      util/logrus.go

+ 1 - 1
bx/BxAreaDynamic.go

@@ -110,7 +110,7 @@ func NewBxAreaProgram(id, runMode, dispMode byte, x uint16, y uint16, w uint16,
 		singleLine:  0x02,
 		autoNewLine: 0x01,
 		dispMode:    dispMode,
-		speed:       0x0a,
+		speed:       0x08,
 		holdTime:    0x08,
 	}
 }

+ 43 - 31
lc/IDevice.go

@@ -1,10 +1,11 @@
 package lc
 
+// 通知接口
 type Notifier interface {
-	Notify(id string)
+	Notify(text string, isProgram, speed int)
 }
 
-// CorrectTimer 校时接口
+// 校时接口
 type CorrectTimer interface {
 	Correct()
 }
@@ -15,7 +16,7 @@ func DoCorrectTime(a any) {
 	}
 }
 
-// ReConnector 重连接口
+// 重连接口
 type ReConnector interface {
 	Reconnect()
 }
@@ -26,48 +27,59 @@ func DoReconnect(a any) {
 	}
 }
 
-// IDevice 抽象设备,包含屏和扬声器
+// 单个设备接口
 type IDevice interface {
-	// Call 通知输出设备,输出信息。屏显示来车警告信息,扬声器语音提示
-	Call()
-	// Rollback 屏回滚到初始状态
-	Rollback()
-	ReConnector
-	CorrectTimer
+	Call(text string, isProgram, speed int) // 触发警告
+	Rollback()                              // 恢复初始状态
+	ReConnector                             // 重连
+	CorrectTimer                            // 校时
 }
 
-type OutputDeviceInfo struct {
-	Name   string
-	Ip     string
-	Port   string
-	Branch byte
-	Audio  string
+// 设备信息
+type DeviceInfo struct {
+	Name  string
+	Ip    string
+	Port  string
+	Audio string // 扬声器音频
 }
 
-type IntersectionDevice struct {
-	Info    OutputDeviceInfo
+// 单个哨兵设备组合(屏幕+喇叭)
+type SentinelDevice struct {
+	Info    DeviceInfo
 	Screen  Screener
 	Speaker Loudspeaker
 }
 
-func (id *IntersectionDevice) Call() {
-	if id.Screen != nil {
-		id.Screen.Display(1)
-	}
-	if id.Speaker != nil {
-		id.Speaker.Speak("支路来车")
+// 触发警告(屏幕显示+喇叭发声)
+func (s *SentinelDevice) Call(text string, isProgram, speed int) {
+	if s.Screen != nil {
+		if isProgram == 1 || isProgram == 0 {
+			s.Screen.Display(Red, DefaultRunMode, DefaultDisplayMode, isProgram, speed, text) // 显示警告
+		} else if isProgram == 2 {
+			s.Screen.Display(Red, DefaultRunMode, MoveLeft, isProgram, speed, text) // 显示警告
+		}
+
 	}
+	//if s.Speaker != nil {
+	//	s.Speaker.Speak("主路来车,请减速") // 语音提示
+	//}
 }
-func (id *IntersectionDevice) Rollback() {
-	if id.Screen != nil {
-		id.Screen.Display(0)
+
+// 恢复初始状态
+func (s *SentinelDevice) Rollback() {
+	if s.Screen != nil {
+		s.Screen.Display(Red, DefaultRunMode, DefaultDisplayMode, 0, 0, "")        //恢复默认
+		s.Screen.Display(Yellow, DefaultRunMode, DefaultDisplayMode, 1, 0, "减速慢行") // 恢复默认显示
 	}
 }
 
-func (id *IntersectionDevice) Reconnect() {
-	DoReconnect(id.Screen)
+// 重连设备
+func (s *SentinelDevice) Reconnect() {
+	DoReconnect(s.Screen)
+	//DoReconnect(s.Speaker)
 }
 
-func (id *IntersectionDevice) Correct() {
-	DoCorrectTime(id.Screen)
+// 校时
+func (s *SentinelDevice) Correct() {
+	DoCorrectTime(s.Screen)
 }

+ 27 - 31
lc/camera_event.go

@@ -15,48 +15,51 @@ import (
 	"time"
 )
 
+// 单摄像头服务器
 func NewCameraEventServer() *CameraServer {
-	server := &CameraServer{Cameras: util.Config.Cameras, Notifiers: make(map[string]Notifier, 4)}
+	camera := util.Config.Cameras[0] // 取第一个摄像头配置
+	server := &CameraServer{
+		Camera:   camera,
+		Notifier: nil,
+	}
 	return server
 }
 
 type CameraServer struct {
-	Cameras   []model.CameraInfo
-	Notifiers map[string]Notifier
+	Camera   model.CameraInfo // 单个摄像头
+	Notifier Notifier         // 单个回调
 }
 
+// 启动HTTP服务监听摄像头事件
 func (s *CameraServer) Start() {
-	http.HandleFunc(util.Config.HikServer.Path, s.Handler)
-	logrus.Fatal("事件监听服务启动失败:", util.Config.HikServer.Addr, ", error:", http.ListenAndServe(util.Config.HikServer.Addr, nil))
+	http.HandleFunc(util.Config.Server.HikServer.Path, s.Handler)
+	logrus.Fatalf("摄像头事件服务启动失败 %s: %v", util.Config.Server.HikServer.Addr,
+		http.ListenAndServe(util.Config.Server.HikServer.Addr, nil))
 }
 
-func (s *CameraServer) RegisterCallback(branch byte, notifier Notifier) {
-	for _, camera := range s.Cameras {
-		//关联主路led屏和支路摄像头;关联支路led屏和主路摄像头
-		if branch == 0 && camera.Branch == 1 || branch == 1 && camera.Branch == 0 {
-			s.Notifiers[camera.IP] = notifier
-		}
-	}
+// 注册回调
+func (s *CameraServer) RegisterCallback(notifier Notifier) {
+	s.Notifier = notifier
 }
 
-func (s *CameraServer) Callback(ip, id string) {
-	notifier, ok := s.Notifiers[ip]
-	if !ok {
-		logrus.Errorf("回调函数注册表没有该ip:%s", ip)
+// 触发回调
+func (s *CameraServer) Callback(text string) {
+	if s.Notifier == nil {
+		logrus.Error("摄像头回调未注册")
 		return
 	}
-	notifier.Notify(id)
-	logrus.Debugf("camera [%s] Callback", ip)
+	s.Notifier.Notify(text, 1, 0)
+	logrus.Debugf("摄像头 %s 触发事件", s.Camera.IP)
 }
 
+// HTTP请求处理(保持原有逻辑,适配单设备)
 func (s *CameraServer) Handler(w http.ResponseWriter, r *http.Request) {
 	gopool.Go(func() {
-		//监听主机应答固定,直接先应答
 		w.WriteHeader(200)
 		w.Header().Add("Date", time.Now().String())
 		w.Header().Add("Connection", "keep-alive")
 	})
-	//1.
+
 	contentType := r.Header.Get("Content-Type")
 	if strings.Contains(contentType, "application/xml") {
 		bytes, err := io.ReadAll(r.Body)
@@ -64,36 +67,30 @@ func (s *CameraServer) Handler(w http.ResponseWriter, r *http.Request) {
 			return
 		}
 		var event model.EventNotificationAlert
-		err = xml.Unmarshal(bytes, &event)
-		if err != nil {
+		if err := xml.Unmarshal(bytes, &event); err != nil {
 			return
 		}
-		//发送事件通知
 		for _, v := range event.DetectionRegionList.DetectionRegionEntry {
-			s.Callback(event.IpAddress, v.RegionID)
+			s.Callback(v.RegionID)
 		}
 	} else if strings.Contains(contentType, "multipart/form-data") {
 		s.HandleMultipart(r)
 	}
 }
 
-// 处理多文件事件
 func (s *CameraServer) HandleMultipart(r *http.Request) {
 	multipartReader := multipart.NewReader(r.Body, "boundary")
-	// 循环读取每个 part
 	var event model.EventNotificationAlert
 	for {
 		part, err := multipartReader.NextPart()
-		//defer part.Close()
 		if err == io.EOF {
 			break
 		}
 		if err != nil {
-			log.Println("Failed to read part:", err)
+			log.Println("读取multipart错误:", err)
 			return
 		}
 		if part.FormName() != "linedetectionImage" || !strings.Contains(part.FormName(), "Image") {
-			//不含图片的xml部分数据
 			xmlData, err := ioutil.ReadAll(part)
 			if err != nil {
 				return
@@ -102,6 +99,5 @@ func (s *CameraServer) HandleMultipart(r *http.Request) {
 			continue
 		}
 	}
-	//发送事件通知
-	s.Callback(event.IpAddress, "")
+	s.Callback("")
 }

+ 69 - 0
lc/command_queue.go

@@ -0,0 +1,69 @@
+package lc
+
+import (
+	"sync"
+	"time"
+)
+
+// CommandQueue 命令队列,用于串行处理Send操作(带150ms间隔)
+type CommandQueue struct {
+	tasks    chan []byte    // 存储待发送的数据
+	running  bool           // 队列运行状态
+	mu       sync.Mutex     // 保护running状态的互斥锁
+	wg       sync.WaitGroup // 用于等待队列结束
+	interval time.Duration  // 任务执行间隔(毫秒)
+}
+
+// NewCommandQueue 创建新的命令队列(指定间隔,默认150ms)
+func NewCommandQueue(bufferSize int, interval ...time.Duration) *CommandQueue {
+	// 默认间隔150毫秒,支持自定义传入
+	intervalMs := 200 * time.Millisecond
+	if len(interval) > 0 {
+		intervalMs = interval[0]
+	}
+	return &CommandQueue{
+		tasks:    make(chan []byte, bufferSize),
+		interval: intervalMs,
+	}
+}
+
+// Start 启动队列处理循环
+func (q *CommandQueue) Start(handle func([]byte)) {
+	q.mu.Lock()
+	defer q.mu.Unlock()
+	if q.running {
+		return
+	}
+	q.running = true
+	q.wg.Add(1)
+	go func() {
+		defer q.wg.Done()
+		// 循环处理队列中的发送任务
+		for data := range q.tasks {
+			handle(data) // 执行实际发送操作
+			time.Sleep(q.interval)
+		}
+	}()
+}
+
+// AddTask 添加发送任务到队列
+func (q *CommandQueue) AddTask(data []byte) {
+	q.mu.Lock()
+	defer q.mu.Unlock()
+	if !q.running {
+		return
+	}
+	q.tasks <- data
+}
+
+// Stop 停止队列,等待所有任务处理完毕
+func (q *CommandQueue) Stop() {
+	q.mu.Lock()
+	defer q.mu.Unlock()
+	if !q.running {
+		return
+	}
+	close(q.tasks) // 关闭通道,让处理goroutine退出
+	q.wg.Wait()    // 等待所有任务处理完毕
+	q.running = false
+}

+ 1 - 3
lc/model/camera.go

@@ -1,7 +1,5 @@
 package model
 
 type CameraInfo struct {
-	IP     string `yaml:"ip"`
-	Name   string `yaml:"name"`
-	Branch byte   `yaml:"branch"`
+	IP string `yaml:"ip"`
 }

+ 1 - 3
lc/model/radar.go

@@ -1,7 +1,5 @@
 package model
 
 type RadarInfo struct {
-	Port   string `yaml:"port"` //对应哪个485口
-	Name   string `yaml:"name"`
-	Branch byte   `yaml:"branch"`
+	Port string `yaml:"port"` //对应哪个485口
 }

+ 3 - 4
lc/model/screen.go

@@ -1,8 +1,7 @@
 package model
 
 type ScreenInfo struct {
-	Name   string `yaml:"name"`
-	Ip     string `yaml:"ip"`
-	Port   string `yaml:"port"`
-	Branch byte   `yaml:"branch"`
+	Name string `yaml:"name"`
+	Ip   string `yaml:"ip"`
+	Port string `yaml:"port"`
 }

+ 1 - 6
lc/model/speaker.go

@@ -1,12 +1,7 @@
 package model
 
 type SpeakerInfo struct {
-	Name   string `yaml:"name"`
-	Ip     string `yaml:"ip"`
-	Branch byte   `yaml:"branch"`
-	Speed  byte   `yaml:"speed"`
-	Volume byte   `yaml:"volume"`
-	Audio  string `yaml:"audio"`
+	Port string `yaml:"port"`
 }
 
 type PlayReq struct {

+ 218 - 35
lc/radar_event.go

@@ -5,76 +5,259 @@ import (
 	"github.com/sirupsen/logrus"
 	"lc-smartX/lc/model"
 	"lc-smartX/util"
+	"math"
+	"strconv"
+	"time"
 )
 
+// 雷达数据结构体(存储解析后的有效数据)
+type RadarData struct {
+	Direction string  // 方向:in(来向)/out(去向)
+	Speed     float64 // 速度值(单位:根据雷达手册,如km/h)
+	Valid     bool    // 数据是否有效
+}
+
+// 单雷达服务器
 func NewRadarEventServer() *RadarServer {
-	s := &RadarServer{Radars: util.Config.Radars, Notifiers: make(map[string]Notifier, 4)}
+	// 只取第一个雷达配置(单个设备)
+	radar := util.Config.Radars[0]
+	s := &RadarServer{
+		Radar:    radar,
+		Notifier: nil, // 单个回调
+	}
 	return s
 }
 
 type RadarServer struct {
-	Radars    []model.RadarInfo
-	Notifiers map[string]Notifier //485通道名
+	Radar    model.RadarInfo // 单个雷达信息
+	Notifier Notifier        // 单个回调
 }
 
+// 启动单个雷达串口监听
 func (s *RadarServer) Start() {
-	for _, radar := range s.Radars {
-		go s.OpenSerial(radar.Port)
-	}
+	go s.OpenSerial(s.Radar.Port) // 直接启动当前雷达
 }
 
-func (s *RadarServer) RegisterCallback(branch byte, notifier Notifier) {
-	for _, radar := range s.Radars {
-		//关联主路led屏和支路雷达;关联支路led屏和主路雷达
-		if branch == 0 && radar.Branch == 1 || branch == 1 && radar.Branch == 0 {
-			s.Notifiers[radar.Port] = notifier
-		}
-	}
+// 注册回调(无需分支判断)
+func (s *RadarServer) RegisterCallback(notifier Notifier) {
+	s.Notifier = notifier
 }
 
-func (s *RadarServer) Callback(port string) {
-	notifier, ok := s.Notifiers[port]
-	if !ok {
-		logrus.Errorf("回调函数注册表没有该ip:%s", port)
+// 兼容原空参数Callback方法(保持接口兼容)
+func (s *RadarServer) Callback() {
+	s.CallbackWithData(RadarData{
+		Direction: "unknown",
+		Speed:     0,
+		Valid:     false,
+	})
+}
+
+// 带数据的回调方法(核心回调逻辑)
+func (s *RadarServer) CallbackWithData(radarData RadarData) {
+	if s.Notifier == nil {
+		logrus.Error("雷达回调未注册")
 		return
 	}
-	notifier.Notify("")
+	// 自定义触发条件:速度>0时触发警告(可根据业务调整阈值,如>5)
+	if radarData.Valid && radarData.Speed > 10 {
+		// 拼接方向和速度作为通知消息
+		s.Notifier.Notify("主路来车", 0, int(math.Round(radarData.Speed)))
+	}
 }
 
+// 打开串口(树莓派USB转485)+ 解析9字节雷达协议
 func (s *RadarServer) OpenSerial(portName string) {
-	// 配置串口参数
+	// 串口配置(严格匹配雷达手册:9600bps、8数据位、1停止位、无校验)
 	options := serial.OpenOptions{
-		PortName:        portName, // /dev/ttymxc4 6 3
-		BaudRate:        9600,
+		PortName:        portName, // 树莓派USB转485通常为/dev/ttyUSB0
+		BaudRate:        9600,     // 根据雷达手册调整波特率
 		DataBits:        8,
 		StopBits:        1,
-		MinimumReadSize: 4,
+		MinimumReadSize: 9,                  // 每帧固定9字节,设置最小读取长度
+		ParityMode:      serial.PARITY_NONE, // 无校验(根据雷达手册调整)
 	}
 
 	// 打开串口
 	port, err := serial.Open(options)
 	if err != nil {
-		panic(err.Error())
+		logrus.Fatalf("雷达串口打开失败: %v", err)
 	}
-
-	// 关闭串口
 	defer port.Close()
+	//logrus.Infof("雷达串口 %s 已打开,开始解析9字节V+/V-协议数据", portName)
+
+	// ======================== 新增:发送雷达配置指令 ========================
+	// 配置指令:来向探测、≈4帧/秒(0x03)、千米/小时(0x00)
+	radarCmd := []byte{0x43, 0x46, 0x02, 0x01, 0x03, 0x00, 0x0d, 0x0a}
+
+	// 发送指令
+	sendLen, err := port.Write(radarCmd)
+	if err != nil {
+		logrus.Errorf("雷达配置指令发送失败: %v", err)
+	} else if sendLen != len(radarCmd) {
+		logrus.Warnf("雷达配置指令发送不完整,发送%d字节,预期%d字节", sendLen, len(radarCmd))
+	} else {
+		logrus.Infof("雷达配置指令发送成功,共发送%d字节", sendLen)
+	}
+
+	// 延迟100ms,确保指令完全发送并被雷达接收(根据雷达响应调整)
+	time.Sleep(100 * time.Millisecond)
+	// ======================== 新增结束 ========================
+
+	// 缓存读取的字节(处理可能的粘包/拆包)
+	readBuf := make([]byte, 0, 32)
+	buf := make([]byte, 9) // 单次读取9字节
+
+	// 持续读取雷达数据
 	for {
-		// 读取数据
-		buf := make([]byte, 128)
+		// 读取9字节数据
 		n, err := port.Read(buf)
 		if err != nil {
-			break
+			logrus.Errorf("雷达读取错误: %v,3秒后尝试重连", err)
+			time.Sleep(3 * time.Second)
+			s.OpenSerial(portName) // 重连串口
+			return
 		}
-		if n < 8 {
-			continue
+
+		// 拼接读取到的字节(处理拆包)
+		readBuf = append(readBuf, buf[:n]...)
+
+		// 确保缓存中有完整的9字节帧才解析
+		for len(readBuf) > 0 {
+			// 查找帧分隔符(\r 或 \n),定位帧结束位置
+			frameEndIdx := -1
+			for i, b := range readBuf {
+				if b == 0x0D || b == 0x0A {
+					frameEndIdx = i
+					break
+				}
+			}
+			if frameEndIdx == -1 {
+				// 无分隔符,若缓存超过32字节则清空(防止垃圾数据堆积)
+				if len(readBuf) > 32 {
+					readBuf = readBuf[len(readBuf)-16:] // 保留最后16字节,避免丢失有效数据
+				}
+				break
+			}
+
+			// 提取一帧数据(从开头到分隔符)
+			frame := readBuf[:frameEndIdx+1]
+			readBuf = readBuf[frameEndIdx+1:] // 剩余数据留到下次解析
+
+			// 解析帧
+			radarData := s.parseRadarFrame(frame)
+			if radarData.Valid {
+				//logrus.Debugf("解析到有效雷达数据:方向=%s,速度=%.1f", radarData.Direction, radarData.Speed)
+				// 触发带数据的回调
+				s.CallbackWithData(radarData)
+			} else {
+				// 无效帧仅打印调试日志(不触发回调)
+				//logrus.Debugf("无效雷达帧,原始字节:%x", frame)
+			}
 		}
-		result := false
-		if buf[0] == 'x' && (buf[1] > '0' || buf[2] > '0' || buf[3] > '0') {
-			result = true
+	}
+}
+
+// 解析单帧雷达数据(优化:兼容前缀空白、帧尾不完整)
+func (s *RadarServer) parseRadarFrame(frame []byte) RadarData {
+	var data RadarData
+
+	// ========== 新增:预处理帧数据 - 清理前缀空白字符(\n/\r/空格/制表符) ==========
+	// 定义需要清理的空白字符:0x00(空)、0x0a(\n)、0x0d(\r)、0x20(空格)、0x09(制表符)
+	trimPrefixChars := []byte{0x00, 0x0a, 0x0d, 0x20, 0x09}
+	// 循环清理开头的空白字符
+	trimmedFrame := frame
+	for len(trimmedFrame) > 0 {
+		isTrimChar := false
+		for _, c := range trimPrefixChars {
+			if trimmedFrame[0] == c {
+				trimmedFrame = trimmedFrame[1:] // 移除首字符
+				isTrimChar = true
+				break
+			}
 		}
-		if result {
-			s.Callback(portName)
+		if !isTrimChar {
+			break // 无空白字符,停止清理
 		}
 	}
+	// 调试日志:打印预处理前后的帧数据
+	//logrus.Debugf("帧预处理:原始=%q (0x%x) → 清理后=%q (0x%x)", frame, frame, trimmedFrame, trimmedFrame)
+
+	// ========== 1. 基础校验:清理后的数据至少保留核心7字节(V+001.9) ==========
+	if len(trimmedFrame) < 7 {
+		//logrus.Warnf("雷达帧长度异常,清理后不足7字节,实际%d字节", len(trimmedFrame))
+		return data
+	}
+
+	// ========== 2. 校验帧头(首字节必须是ASCII的'V') ==========
+	if trimmedFrame[0] != 'V' {
+		//logrus.Debugf("雷达帧头错误,期望'V'(0x56),实际0x%02x(字符:%q)", trimmedFrame[0], trimmedFrame[0])
+		return data
+	}
+
+	// ========== 3. 适配帧尾:兼容仅含\r(0x0D)或完整\r\n(0x0D+0x0A) ==========
+	// 核心数据长度:V + +/- + 百位+十位+个位 + . + 小数位 = 7字节(索引0-6)
+	// 帧尾允许两种情况:7字节后是\r(0x0D) 或 \r\n(0x0D+0x0A)
+	var frameTailValid bool
+	if len(trimmedFrame) >= 8 && trimmedFrame[7] == 0x0D {
+		// 情况1:7字节核心数据 + \r(共8字节)
+		frameTailValid = true
+	} else if len(trimmedFrame) >= 9 && trimmedFrame[7] == 0x0D && trimmedFrame[8] == 0x0A {
+		// 情况2:7字节核心数据 + \r\n(共9字节,标准协议)
+		frameTailValid = true
+	}
+	if !frameTailValid {
+		//logrus.Debugf("雷达帧尾错误,清理后帧尾:%q (0x%x),期望\\r 或 \\r\\n", trimmedFrame[7:], trimmedFrame[7:])
+		return data
+	}
+
+	// ========== 4. 解析方向(第2字节:+ 来向 / - 去向) ==========
+	switch trimmedFrame[1] {
+	case '+':
+		data.Direction = "in" // 来向目标
+	case '-':
+		data.Direction = "out" // 去向目标
+	default:
+		//logrus.Debugf("雷达方向标识错误,期望'+'/'-',实际0x%02x(字符:%q)", trimmedFrame[1], trimmedFrame[1])
+		return data
+	}
+
+	// ========== 5. 校验小数点位置(第6字节必须是'.') ==========
+	if trimmedFrame[5] != '.' {
+		//logrus.Debugf("雷达小数点位置错误,期望0x2E('.'),实际0x%02x(字符:%q)", trimmedFrame[5], trimmedFrame[5])
+		return data
+	}
+
+	// ========== 6. 字符转数字(ASCII字符转数值,如'1'→1) ==========
+	// 百位:trimmedFrame[2] → 0-9的ASCII字符
+	hundreds, err := strconv.Atoi(string(trimmedFrame[2]))
+	if err != nil {
+		//logrus.Debugf("雷达百位数值错误:%v(字符:%q)", err, trimmedFrame[2])
+		return data
+	}
+	// 十位:trimmedFrame[3]
+	tens, err := strconv.Atoi(string(trimmedFrame[3]))
+	if err != nil {
+		//logrus.Debugf("雷达十位数值错误:%v(字符:%q)", err, trimmedFrame[3])
+		return data
+	}
+	// 个位:trimmedFrame[4]
+	units, err := strconv.Atoi(string(trimmedFrame[4]))
+	if err != nil {
+		//logrus.Debugf("雷达个位数值错误:%v(字符:%q)", err, trimmedFrame[4])
+		return data
+	}
+	// 小数位:trimmedFrame[6]
+	decimal, err := strconv.Atoi(string(trimmedFrame[6]))
+	if err != nil {
+		//logrus.Debugf("雷达小数位数值错误:%v(字符:%q)", err, trimmedFrame[6])
+		return data
+	}
+
+	// ========== 7. 计算最终速度值(xxx.x格式) ==========
+	data.Speed = float64(hundreds*100+tens*10+units) + float64(decimal)/10.0
+	// 标记数据有效
+	data.Valid = true
+
+	//logrus.Debugf("解析成功:方向=%s,速度=%.1f(原始数据:%q)", data.Direction, data.Speed, frame)
+	return data
 }

+ 135 - 0
lc/redis_process.go

@@ -0,0 +1,135 @@
+package lc
+
+import (
+	"encoding/json"
+	"fmt"
+	"github.com/redis/go-redis/v9"
+	"github.com/sirupsen/logrus"
+	"regexp"
+	"strings"
+	"time"
+)
+
+// 启动定时任务:动态调整轮询间隔(未拿到车牌号3秒轮询,拿到则5秒轮询)
+func (s *SentinelServer) StartFetchingLatestPlate() {
+	go func() {
+		for {
+			// 1. 执行获取最新车牌号逻辑
+			latestData, err := s.GetLatestPlateData(1)
+			if err != nil {
+				//logrus.Warn("获取最新车牌号失败", "错误", err)
+				time.Sleep(3 * time.Second)
+			} else if len(latestData) == 0 {
+				logrus.Info("未获取到车牌号数据")
+				time.Sleep(3 * time.Second)
+			} else {
+				// 2. 成功获取到车牌号,处理业务逻辑
+				latestPlateNo := latestData[0].PlateNo
+				logrus.Info("最新车牌号", "plate_no", latestPlateNo)
+				s.Notify(latestPlateNo, 2, 0)
+				time.Sleep(12 * time.Second)
+			}
+		}
+	}()
+}
+
+type RedisProcess struct{}
+
+type PlateData struct {
+	PlateNo     string `json:"plate_no"`     // 车牌号
+	PlateColor  string `json:"plate_color"`  // 车牌颜色
+	DetectConf  string `json:"detect_conf"`  // 检测置信度
+	ColorConf   string `json:"color_conf"`   // 颜色置信度
+	RecAvg      string `json:"rec_avg"`      // 识别平均置信度
+	Timestamp   string `json:"timestamp"`    // 秒级时间戳
+	TimestampMs string `json:"timestamp_ms"` // 毫秒级时间戳
+	Datetime    string `json:"datetime"`     // 格式化时间
+	Source      string `json:"source"`       // 数据来源
+}
+
+func (s *SentinelServer) GetLatestPlateData(limit int) ([]PlateData, error) {
+	if s.rdb == nil {
+		return nil, fmt.Errorf("Redis客户端未初始化")
+	}
+
+	// 步骤1:从ZSet获取最新的limit个timestamp(按Score倒序)
+	zsetKey := fmt.Sprintf("%s:sorted", RedisKeyPrefix)
+	// ZREVRANGE:倒序取0到limit-1的member(最新的N个)
+	timestamps, err := s.rdb.ZRevRange(s.ctx, zsetKey, 0, int64(limit-1)).Result()
+	if err != nil {
+		return nil, fmt.Errorf("读取ZSet失败: %w", err)
+	}
+	if len(timestamps) == 0 {
+		return nil, fmt.Errorf("Redis中无车牌数据")
+	}
+
+	// 步骤2:批量从Hash读取对应数据
+	hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix)
+	pipe := s.rdb.Pipeline()
+	for _, ts := range timestamps {
+		pipe.HGet(s.ctx, hashKey, ts)
+	}
+	results, err := pipe.Exec(s.ctx)
+	if err != nil {
+		return nil, fmt.Errorf("批量读取Hash失败: %w", err)
+	}
+
+	// 步骤3:解析Python格式的字符串为Go结构体
+	var plateDataList []PlateData
+	for i, res := range results {
+		if res == nil {
+			continue
+		}
+		// 获取Hash的value(Python的str(dict)字符串)
+		pyStr, err := res.(*redis.StringCmd).Result()
+		if err != nil {
+			logrus.Warn("解析单条数据失败", "timestamp", timestamps[i], "err", err)
+			continue
+		}
+		// 转换Python字符串为标准JSON,再解析
+		plateData, err := parsePythonDictStr(pyStr)
+		if err != nil {
+			logrus.Warn("转换Python字符串失败", "str", pyStr, "err", err)
+			continue
+		}
+		plateDataList = append(plateDataList, plateData)
+	}
+
+	return plateDataList, nil
+}
+
+// 2. 读取指定时间戳的车牌数据(精准读取)
+func (s *SentinelServer) GetPlateDataByTimestamp(timestamp string) (*PlateData, error) {
+	if s.rdb == nil {
+		return nil, fmt.Errorf("Redis客户端未初始化")
+	}
+
+	hashKey := fmt.Sprintf("%s:data", RedisKeyPrefix)
+	pyStr, err := s.rdb.HGet(s.ctx, hashKey, timestamp).Result()
+	if err != nil {
+		return nil, fmt.Errorf("读取指定时间戳数据失败: %w", err)
+	}
+
+	// 解析为结构体
+	plateData, err := parsePythonDictStr(pyStr)
+	if err != nil {
+		return nil, fmt.Errorf("解析数据失败: %w", err)
+	}
+	return &plateData, nil
+}
+
+// ========== 辅助函数:解析Python的dict字符串为Go结构体 ==========
+// Python的str(dict)是单引号,需替换为双引号,且处理特殊字符
+func parsePythonDictStr(pyStr string) (PlateData, error) {
+	var plateData PlateData
+	// 1. 替换单引号为双引号(Python -> JSON)
+	jsonStr := strings.ReplaceAll(pyStr, "'", "\"")
+	// 2. 处理可能的空格/换行(可选,视Python字符串格式)
+	jsonStr = regexp.MustCompile(`\s+`).ReplaceAllString(jsonStr, "")
+	// 3. 解析为JSON
+	err := json.Unmarshal([]byte(jsonStr), &plateData)
+	if err != nil {
+		return plateData, fmt.Errorf("JSON解析失败: %w, raw_str: %s", err, jsonStr)
+	}
+	return plateData, nil
+}

+ 159 - 58
lc/screen.go

@@ -1,6 +1,7 @@
 package lc
 
 import (
+	"encoding/hex"
 	"fmt"
 	"github.com/sirupsen/logrus"
 	"golang.org/x/text/encoding/simplifiedchinese"
@@ -8,11 +9,12 @@ import (
 	"net"
 	"strconv"
 	"time"
+	"unicode/utf8"
 )
 
 // Screener 屏接口
 type Screener interface {
-	Display(int)
+	Display(Color, RunMode, DisplayMode, int, int, string)
 }
 
 type Screen struct {
@@ -22,24 +24,66 @@ type Screen struct {
 	liveState bool
 	StateInfo *bx.StateInfo //状态信息
 	Params    *bx.Params    //屏参
+	sendQueue *CommandQueue // 发送队列,用于串行处理Send操作
 }
 
+// 初始化屏幕时创建并启动发送队列
 func NewScreen(name string, ip, port string) *Screen {
 	s := &Screen{
 		Name:      name,
 		Addr:      fmt.Sprintf("%s:%s", ip, port),
 		StateInfo: &bx.StateInfo{},
 		Params:    &bx.Params{},
+		sendQueue: NewCommandQueue(100), // 缓冲区大小100,可调整
 	}
+	// 启动队列,传入实际发送处理函数
+	s.sendQueue.Start(s.handleSend)
 	s.Reconnect()
 	return s
 }
 
-func (s *Screen) Display(id int) {
+// 实际执行发送的处理函数(由队列goroutine调用)
+func (s *Screen) handleSend(data []byte) {
+	if !s.getLiveState() || s.conn == nil {
+		return
+	}
+	_, err := s.conn.Write(data)
+	if err != nil {
+		logrus.WithFields(map[string]interface{}{"设备名": s.Name}).Error("tcp write error:", err)
+		s.liveState = false
+	}
+
+	//resp := s.ReadResp()
+	//if !resp.IsAck() {
+	//	logrus.Error("设备拒绝写文件! error:", resp.Error().Description)
+	//	return
+	//}
+}
+
+// 修改Send方法,将数据添加到队列而非直接发送
+func (s *Screen) Send(data []byte) {
 	if !s.getLiveState() {
 		return
 	}
-	s.SendRam(id)
+	// 将发送任务加入队列,由队列串行处理
+	s.sendQueue.AddTask(data)
+}
+
+// 添加Close方法,用于释放队列资源
+func (s *Screen) Close() {
+	if s.sendQueue != nil {
+		s.sendQueue.Stop()
+	}
+	if s.conn != nil {
+		s.conn.Close()
+	}
+}
+
+func (s *Screen) Display(color Color, runMode RunMode, displayMode DisplayMode, isProgram, speed int, text string) {
+	if !s.getLiveState() {
+		return
+	}
+	s.SendRam(color, runMode, displayMode, isProgram, speed, text)
 }
 
 // Correct 校正时间
@@ -53,11 +97,15 @@ func (s *Screen) Correct() {
 	s.Send(data.Pack())
 }
 
-// Reconnect 重连
+// Reconnect方法中添加连接关闭时的队列处理(可选增强)
 func (s *Screen) Reconnect() {
 	if s.getLiveState() {
 		return
 	}
+	// 关闭旧连接(如果存在)
+	if s.conn != nil {
+		s.conn.Close()
+	}
 	conn, err := net.DialTimeout("tcp", s.Addr, 5*time.Second)
 	if err != nil {
 		logrus.Error(s.Name, "-", s.Addr, "[屏]重连接失败! error:", err)
@@ -81,17 +129,17 @@ func (s *Screen) setConn(conn net.Conn) {
 	s.liveState = true
 }
 
-// 给屏发送数据
-func (s *Screen) Send(data []byte) {
-	if !s.getLiveState() {
-		return
-	}
-	_, err := s.conn.Write(data)
-	if err != nil {
-		logrus.WithFields(map[string]interface{}{"设备名": s.Name}).Error("tcp write error:", err)
-		s.liveState = false
-	}
-}
+//// 给屏发送数据
+//func (s *Screen) Send(data []byte) {
+//	if !s.getLiveState() {
+//		return
+//	}
+//	_, err := s.conn.Write(data)
+//	if err != nil {
+//		logrus.WithFields(map[string]interface{}{"设备名": s.Name}).Error("tcp write error:", err)
+//		s.liveState = false
+//	}
+//}
 
 //以下对协议进行封装
 //↓↓↓↓↓↓↓↓↓↓↓↓↓↓↓
@@ -115,60 +163,92 @@ const (
 )
 
 // 发送动态区节目 0正常页面 1红色提醒页面
-func (s *Screen) SendRam(id int) {
-	file := FlashFile{}
-	if id == 0 {
-		file.SetMsg("减速慢行", Yellow)
-		file.SetMode(DefaultRunMode, DefaultDisplayMode)
-		file.SetOrigin(0, true, 0)
-		file.SetArea(64, true, 16)
-		s.TextRam(file, false)
-	} else {
-		file.SetMsg("支路来车", Red)
-		file.SetSoundData("支路来车,请减速")
-		file.SetMode(DefaultRunMode, DefaultDisplayMode)
-		file.SetOrigin(0, true, 0)
-		file.SetArea(64, true, 16)
-		s.TextRam(file, true)
+func (s *Screen) SendRam(color Color, runMode RunMode, displayMode DisplayMode, isProgram, speed int, text string) {
+
+	if isProgram == 0 {
+		file1 := FlashFile{}
+		file1.SetMsg(FormatSpeed(speed), color)
+		file1.SetMode(runMode, displayMode)
+		file1.SetOrigin(8, true, 0)
+		file1.SetArea(16, true, 32)
+		s.TextRam(file1, false)
+	} else if isProgram == 1 {
+		file1 := FlashFile{}
+		file1.SetMsg(text, color)
+		file1.SetMode(runMode, displayMode)
+		file1.SetOrigin(32, true, 4)
+		file1.SetArea(96, true, 32)
+		s.TextFlash([]FlashFile{file1}, "P000", false, false)
+	} else if isProgram == 2 {
+		file1 := FlashFile{}
+		file1.SetMsg(text, color)
+		file1.SetMode(runMode, displayMode)
+		file1.SetOrigin(32, true, 4)
+		file1.SetArea(96, true, 32)
+		s.TextRam(file1, true)
+	}
+
+}
+
+func FormatSpeed(speed int) string {
+	// 校验输入范围(仅处理0-99的整数)
+	if speed < 0 || speed > 99 {
+		return "0\\n0"
+	}
+
+	// 处理个位数(0-9)
+	if speed < 10 {
+		speedStr := strconv.Itoa(speed)
+		return fmt.Sprintf("%s\\n0", speedStr)
 	}
+
+	// 处理十位数(10-99)
+	tens := speed / 10  // 提取十位数字(如21→2)
+	units := speed % 10 // 提取个位数字(如21→1)
+	return fmt.Sprintf("%d\\n%d", units, tens)
 }
 
 // TextRam 发送动态区节目实现
-func (s *Screen) TextRam(ff FlashFile, needSpeak bool) {
+func (s *Screen) TextRam(ff FlashFile, isNumber bool) {
 	if !s.getLiveState() {
 		return
 	}
 	var areas []bx.BxArea
 	encoder := simplifiedchinese.GB18030.NewEncoder()
 	var bytes []byte
-	if ff.color == Default {
+
+	var id = 0
+
+	if ff.color == Default && isNumber == false {
 		bytes, _ = encoder.Bytes([]byte(ff.msg))
-	} else {
-		bytes, _ = encoder.Bytes([]byte("\\C" + strconv.Itoa(int(ff.color)) + ff.msg))
-	}
-	if needSpeak {
-		soundData, _ := encoder.Bytes([]byte(ff.soundData))
-		area := bx.NewBxAreaDynamic(0, 1, byte(ff.dispMode), ff.originX, ff.originY, ff.width,
-			ff.height, bytes, soundData, false)
-		areas = append(areas, area)
-	} else {
-		area := bx.NewBxAreaProgram(0, 1, byte(ff.dispMode), ff.originX, ff.originY, ff.width,
-			ff.height, bytes, false)
-		areas = append(areas, area)
+	} else if isNumber == true {
+		province, number := ParsePlateUTF8(ff.msg)
+		logrus.Info("车牌号", ff.msg)
+		logrus.Info("车牌号", "号码", number, "省份", province)
+		bytes, _ = encoder.Bytes([]byte("\\FO000\\C" + strconv.Itoa(int(ff.color)) + province + "\\FE001\\C" + strconv.Itoa(int(Green)) + number))
+		id = 1
+	} else if isNumber == false {
+		bytes, _ = encoder.Bytes([]byte("\\FE000\\C" + strconv.Itoa(int(ff.color)) + ff.msg))
 	}
+	area := bx.NewBxAreaProgram(byte(id), 1, byte(ff.dispMode), ff.originX, ff.originY, ff.width,
+		ff.height, bytes, false)
+	areas = append(areas, area)
 
 	//
 	cmd := bx.NewBxCmdSendDynamicArea(areas)
 	pack := bx.NewBxDataPackCmd(cmd)
-	pack.SetDisplayType(1) //动态显示模式
+	pack.SetDisplayType(0) //动态显示模式
 	d := pack.Pack()
+	if isNumber {
+		logrus.Info(hex.EncodeToString(d))
+	}
 	s.Send(d)
 	resp := s.ReadResp()
 	if !resp.IsAck() {
-		println("设备拒绝写文件! error:", resp.Error().Description)
+		logrus.Error("设备拒绝写文件! error:", resp.Error().Description)
 		return
 	}
-	//s.StateInfo.DynaAreaNum++
+	s.StateInfo.DynaAreaNum++
 }
 
 // DelRamText 删除动态区,不传删除所有
@@ -258,7 +338,7 @@ func (ft *FlashFile) SetArea(w uint16, yIsPixel bool, h uint16) {
 }
 
 // TextFlash 发送静态文件节目,掉电保存,文件名格式"P000","P001"
-func (s *Screen) TextFlash(ft []FlashFile, name string, isLogo bool) {
+func (s *Screen) TextFlash(ft []FlashFile, name string, isNumber, isLogo bool) {
 	if !s.getLiveState() {
 		return
 	}
@@ -269,10 +349,15 @@ func (s *Screen) TextFlash(ft []FlashFile, name string, isLogo bool) {
 	var areas []bx.BxArea
 	for _, i := range ft {
 		var bytes []byte
-		if i.color == Default {
+		if i.color == Default && isNumber == false {
 			bytes, _ = encoder.Bytes([]byte(i.msg))
-		} else {
-			bytes, _ = encoder.Bytes([]byte("\\C" + strconv.Itoa(int(i.color)) + i.msg))
+		} else if isNumber == true {
+			province, number := ParsePlateUTF8(i.msg)
+			logrus.Info("车牌号", i.msg)
+			logrus.Info("车牌号", "号码", number, "省份", province)
+			bytes, _ = encoder.Bytes([]byte("\\FO000\\C" + strconv.Itoa(int(i.color)) + province + "\\FE001\\C" + strconv.Itoa(int(Green)) + number))
+		} else if isNumber == false {
+			bytes, _ = encoder.Bytes([]byte("\\FO000\\C" + strconv.Itoa(int(i.color)) + i.msg))
 		}
 		area := bx.NewBxAreaProgram(0xff, byte(i.runMode), byte(i.dispMode), i.originX, i.originY, i.width, i.height, bytes, false)
 		areas = append(areas, area)
@@ -282,19 +367,29 @@ func (s *Screen) TextFlash(ft []FlashFile, name string, isLogo bool) {
 	cmd := file.NewCmdWriteFile()
 	pack := bx.NewBxDataPackCmd(cmd)
 	data := pack.Pack()
+	if isNumber {
+		logrus.Info(hex.EncodeToString(data))
+	}
 	s.Send(data)
 	resp := s.ReadResp()
 	if !resp.IsAck() {
 		logrus.Error("设备拒绝写文件! error:", resp.Error().Description)
 		return
 	}
-	pack1 := bx.NewBxDataPackCmd(cmd)
-	data1 := pack1.Pack()
-	s.Send(data1)
-	resp1 := s.ReadResp()
-	if resp1.NoError() {
-		s.StateInfo.ProgramNum++
-	}
+	//pack1 := bx.NewBxDataPackCmd(cmd)
+	//data1 := pack1.Pack()
+	//s.Send(data1)
+	//resp1 := s.ReadResp()
+	//if resp1.NoError() {
+	//	s.StateInfo.ProgramNum++
+	//}
+}
+
+func ParsePlateUTF8(rawPlate string) (province, number string) {
+	firstRune, firstLen := utf8.DecodeRuneInString(rawPlate)
+	province = string(firstRune)
+	number = rawPlate[firstLen:]
+	return
 }
 
 // Bitmap 发送自定义位图节目
@@ -320,7 +415,13 @@ func (s *Screen) Bitmap(name string, bitmap []byte) {
 func (s *Screen) Lock(flag byte, name string) {
 	cmd := bx.NewCmdLock(flag, name)
 	pack := bx.NewBxDataPackCmd(&cmd)
+	fmt.Println(hex.EncodeToString(pack.Pack()))
 	s.Send(pack.Pack())
+	resp := s.ReadResp()
+	if !resp.IsAck() {
+		logrus.Error("设备拒绝写文件! error:", resp.Error().Description)
+		return
+	}
 }
 
 // DelFile 删除静态文件节目

+ 111 - 126
lc/server.go

@@ -1,155 +1,140 @@
 package lc
 
 import (
+	"context"
+	"github.com/sirupsen/logrus"
 	"lc-smartX/util"
 	"lc-smartX/util/gopool"
 	"time"
+
+	// 仅引入Redis客户端核心依赖
+	"github.com/redis/go-redis/v9"
 )
 
-type SmartXServer interface {
-	Serve()
-}
+// Redis基础配置(极简版,仅用于初始化连接)
+const (
+	RedisInitKey   = "sentinel:init:flag" // 仅用于标记Redis初始化完成的空Key
+	RedisKeyPrefix = "plate_results"
+)
 
-type IntersectionServer struct {
-	RadarEventServer  *RadarServer
-	CameraEventServer *CameraServer
-	MainState         byte
-	SubState          byte
-	MainDevices       []IDevice
-	SubDevices        []IDevice
-	ReTicker          *time.Ticker
-	Main              *time.Ticker
-	Sub               *time.Ticker
+// 单个哨兵服务(仅新增Redis客户端和上下文字段)
+type SentinelServer struct {
+	RadarServer  *RadarServer
+	CameraServer *CameraServer
+	Device       IDevice      // 本哨兵的设备(屏幕+喇叭)
+	StateTicker  *time.Ticker // 状态回滚定时器
+	ReTicker     *time.Ticker // 重连定时器
+	// === 仅新增Redis核心字段 ===
+	rdb         *redis.Client   // Redis客户端
+	ctx         context.Context // 上下文(用于Redis操作)
+	fetchTicker *time.Ticker    // 定时任务ticker
+	stopChan    chan struct{}   // 停止信号通道
 }
 
-func StartIntersectionServer() {
-	is := &IntersectionServer{
-		Main:     time.NewTicker(5 * time.Second),  //主路状态回滚
-		Sub:      time.NewTicker(5 * time.Second),  //支路状态回滚
-		ReTicker: time.NewTicker(30 * time.Second), //重连
-	}
-	if util.Config.Server.SupportCamera {
-		is.CameraEventServer = NewCameraEventServer()
-		gopool.Go(is.CameraEventServer.Start)
+// 启动哨兵服务
+func StartSentinelServer() {
+	s := &SentinelServer{
+		StateTicker: time.NewTicker(5 * time.Second),  // 5秒后自动回滚
+		ReTicker:    time.NewTicker(30 * time.Second), // 30秒重连检查
+		ctx:         context.Background(),             // 初始化Redis上下文
 	}
+
+	// === 第一步:优先初始化Redis(极简版) ===
+	s.initRedis()
+
+	// 原有逻辑:初始化设备
+	s.initDevice()
+
+	time.Sleep(2 * time.Second)
+
+	// 原有逻辑:启动雷达服务
 	if util.Config.Server.SupportRadar {
-		is.RadarEventServer = NewRadarEventServer()
-		gopool.Go(is.RadarEventServer.Start)
+		s.RadarServer = NewRadarEventServer()
+		s.RadarServer.RegisterCallback(s) // 注册雷达回调
+		gopool.Go(s.RadarServer.Start)
 	}
-	//等事件服务先启动
-	time.Sleep(1 * time.Second)
-	is.Serve()
-}
 
-type MainNotifier struct{ s *IntersectionServer }
-
-// Notify 主路来车,通知支路设备
-func (m MainNotifier) Notify(id string) {
-	m.s.Main.Reset(5 * time.Second)
-	if m.s.MainState != 1 {
-		m.s.MainState = 1
-		if id == "1" {
-			for _, v := range m.s.SubDevices {
-				gopool.Go(v.Call)
-			}
-		} else if id == "2" {
-			for _, v := range m.s.MainDevices {
-				gopool.Go(v.Call)
-			}
-		}
+	// 原有逻辑:启动摄像头服务
+	if util.Config.Server.SupportCamera {
+		s.CameraServer = NewCameraEventServer()
+		s.CameraServer.RegisterCallback(s) // 注册摄像头回调
+		gopool.Go(s.CameraServer.Start)
 	}
-}
+	// 启动定时任务:每3秒获取最新一条车牌号
+	s.StartFetchingLatestPlate()
 
-type SubNotifier struct{ s *IntersectionServer }
-
-// Notify 支路来车,通知主路设备
-func (sub SubNotifier) Notify(id string) {
-	sub.s.Sub.Reset(5 * time.Second)
-	if sub.s.SubState != 1 {
-		sub.s.SubState = 1
-		if id == "1" {
-			for _, v := range sub.s.SubDevices {
-				gopool.Go(v.Call)
-			}
-		} else if id == "2" {
-			for _, v := range sub.s.MainDevices {
-				gopool.Go(v.Call)
-			}
-		}
-	}
+	// 原有逻辑:启动主循环
+	s.Serve()
 }
 
-func (is *IntersectionServer) Serve() {
-	if util.Config.Server.SupportCamera {
-		is.CameraEventServer.RegisterCallback(1, &SubNotifier{is})
-		is.CameraEventServer.RegisterCallback(0, &MainNotifier{is})
-	}
-	if util.Config.Server.SupportRadar {
-		is.RadarEventServer.RegisterCallback(1, &SubNotifier{is})
-		is.RadarEventServer.RegisterCallback(0, &MainNotifier{is})
+// === 核心:极简Redis初始化(仅连接+标记初始化) ===
+func (s *SentinelServer) initRedis() {
+	// 1. 创建Redis客户端
+	s.rdb = redis.NewClient(&redis.Options{
+		Addr:     util.Config.RedisConfig.Addr,
+		Password: util.Config.RedisConfig.Password,
+		DB:       util.Config.RedisConfig.DB,
+	})
+
+	// 2. 测试Redis连接
+	if err := s.rdb.Ping(s.ctx).Err(); err != nil {
+		logrus.Warn("⚠️ Redis连接失败,哨兵服务正常启动(无Redis支持)", "err", err)
+		s.rdb = nil // 标记Redis不可用
+		return
 	}
 
-	//先创建响应设备
-	for _, c := range util.Config.Screens {
-		iDevice := &IntersectionDevice{
-			Info: OutputDeviceInfo{
-				Name:   c.Name,
-				Ip:     c.Ip,
-				Port:   c.Port,
-				Branch: c.Branch,
-			},
-			Screen: NewScreen(c.Name, c.Ip, c.Port),
-		}
-		if c.Branch == 1 {
-			is.MainDevices = append(is.MainDevices, iDevice)
-		} else {
-			is.SubDevices = append(is.SubDevices, iDevice)
-		}
+	// 3. 写入初始化标记Key(仅标记,无业务数据)
+	if err := s.rdb.Set(s.ctx, RedisInitKey, "initialized", 0).Err(); err != nil {
+		logrus.Warn("⚠️ Redis初始化标记写入失败", "err", err)
+	} else {
+		logrus.Info("✅ Redis初始化完成(仅连接+标记)")
 	}
-	for _, c := range util.Config.Speakers {
-		iDevice := &IntersectionDevice{
-			Info: OutputDeviceInfo{
-				Name:   c.Name,
-				Ip:     c.Ip,
-				Branch: c.Branch,
-				Audio:  c.Audio,
-			},
-			Speaker: NewIpCast(c.Name, c.Ip, c.Audio, c.Speed, c.Volume),
-		}
-		if c.Branch == 1 {
-			is.MainDevices = append(is.MainDevices, iDevice)
-		} else {
-			is.SubDevices = append(is.SubDevices, iDevice)
-		}
+}
+
+// 原有逻辑:初始化本哨兵的设备(屏幕+喇叭)
+func (s *SentinelServer) initDevice() {
+	// 屏幕配置(取第一个屏幕)
+	screenCfg := util.Config.Screens[0]
+	screen := NewScreen(screenCfg.Name, screenCfg.Ip, screenCfg.Port)
+
+	// 喇叭配置(取第一个喇叭)
+	speakerCfg := util.Config.Speakers[0]
+	speaker := NewIpCast(
+		speakerCfg.Port,
+	)
+
+	// 组合设备
+	s.Device = &SentinelDevice{
+		Info: DeviceInfo{
+			Name:  screenCfg.Name,
+			Ip:    screenCfg.Ip,
+			Port:  screenCfg.Port,
+			Audio: "支路来车",
+		},
+		Screen:  screen,
+		Speaker: speaker,
 	}
+}
+
+// 原有逻辑:实现Notifier接口,处理雷达/摄像头事件
+func (s *SentinelServer) Notify(text string, isProgram, speed int) {
+	s.StateTicker.Reset(5 * time.Second) // 重置回滚定时器
+	go s.Device.Call(text, isProgram, speed)
+	//gopool.Go(func() {
+	//	s.Device.Call(text, isProgram, speed)
+	//}) // 触发设备警告
+}
 
+// 原有逻辑:主循环:处理状态回滚和重连
+func (s *SentinelServer) Serve() {
 	for {
 		select {
-		case <-is.Main.C: //检查主路状态->支路输出设备回到初始状态
-			for _, v := range is.SubDevices {
-				if is.MainState == 1 {
-					gopool.Go(v.Rollback)
-				}
-			}
-			is.MainState = 0
-		case <-is.Sub.C: //检查支路状态->主路输出设备作出响应
-			for _, v := range is.MainDevices {
-				if is.SubState == 1 {
-					gopool.Go(v.Rollback)
-				}
-			}
-			is.SubState = 0
-		case <-is.ReTicker.C: //每19s检查并尝试重连
-			gopool.Go(func() {
-				for _, v := range is.MainDevices {
-					gopool.Go(v.Reconnect)
-				}
-			})
-			gopool.Go(func() {
-				for _, v := range is.SubDevices {
-					gopool.Go(v.Reconnect)
-				}
-			})
+		case <-s.StateTicker.C:
+			// 定时器触发,回滚设备状态
+			gopool.Go(s.Device.Rollback)
+		case <-s.ReTicker.C:
+			// 定时重连设备
+			gopool.Go(s.Device.Reconnect)
 		}
 	}
 }

+ 223 - 48
lc/speaker.go

@@ -2,12 +2,17 @@ package lc
 
 import (
 	"bytes"
-	"encoding/json"
+	"encoding/binary"
+	"encoding/hex"
+	"errors"
 	"fmt"
-	"github.com/sirupsen/logrus"
 	"io"
-	"lc-smartX/lc/model"
-	"net/http"
+	"sync"
+	"time"
+
+	"github.com/tarm/serial"
+	"golang.org/x/text/encoding/simplifiedchinese"
+	"golang.org/x/text/transform"
 )
 
 // Loudspeaker 扬声器接口
@@ -15,70 +20,240 @@ type Loudspeaker interface {
 	Speak(txt string)
 }
 
+// IpCast 串口语音播报实现(核心:播放中丢弃请求+仅空闲时接收+无缓存)
 type IpCast struct {
-	Name      string
-	Ip        string
-	Speed     byte
-	Volume    byte
-	Audio     string
-	liveState bool
+	Port       string        // 串口端口(外部传入)
+	Baud       int           // 波特率(固定为115200)
+	serialPort *serial.Port  // 串口连接实例
+	mu         sync.Mutex    // 全局互斥锁(保护isPlaying/串口操作)
+	idleChan   chan struct{} // 芯片空闲通知通道
+	isClosed   bool          // 串口是否已关闭
+	isPlaying  bool          // 标记是否正在播放语音(核心控制字段)
 }
 
-func NewIpCast(name, ip, audio string, speed, volume byte) *IpCast {
+// NewIpCast 初始化串口语音播报实例
+// prot: 串口端口(如"/dev/ttyUSB0")
+func NewIpCast(prot string) *IpCast {
 	s := &IpCast{
-		Name:      name,
-		Ip:        ip,
-		Speed:     speed,
-		Volume:    volume,
-		Audio:     audio,
-		liveState: false,
-	}
-	s.Reconnect()
+		Port:      prot,
+		Baud:      115200,
+		idleChan:  make(chan struct{}, 1), // 缓冲通道避免阻塞
+		isClosed:  false,
+		isPlaying: false, // 初始为未播放状态
+	}
+	// 初始化串口
+	if err := s.Reconnect(); err != nil {
+		fmt.Printf("串口初始化失败: %v\n", err)
+	}
+	// 启动后台监听协程
+	go s.listenSerialResponse()
 	return s
 }
 
-func (ip *IpCast) Speak(txt string) {
-	data := &model.PlayReq{
-		Text:   txt,
-		Vcn:    "xiaoyan",
-		Speed:  ip.Speed,
-		Volume: ip.Volume,
-		Rdn:    "2",
-		Rcn:    "1",
-		Reg:    0,
-		Sync:   false,
-		Queue:  false,
-	}
-	data.Loop.Times = 1
-	body, err := json.Marshal(data)
+// Reconnect 重连串口(加锁保护+标记串口状态)
+func (ip *IpCast) Reconnect() error {
+	ip.mu.Lock()
+	defer ip.mu.Unlock()
+
+	// 标记串口为关闭状态
+	ip.isClosed = true
+	if ip.serialPort != nil {
+		_ = ip.serialPort.Close()
+		ip.serialPort = nil
+	}
+
+	// 配置串口参数
+	cfg := &serial.Config{
+		Name:        ip.Port,
+		Baud:        ip.Baud,
+		Size:        8,
+		Parity:      serial.ParityNone,
+		StopBits:    1,
+		ReadTimeout: time.Millisecond * 200, // 减少EOF报错频率
+	}
+
+	port, err := serial.OpenPort(cfg)
 	if err != nil {
-		logrus.Errorf("IpCast Marshal err : %s", err.Error())
+		return fmt.Errorf("打开串口失败: %w", err)
+	}
+
+	ip.serialPort = port
+	ip.isClosed = false // 标记串口可用
+	fmt.Println("串口重连成功")
+	return nil
+}
+
+// clearIdleChan 清空空闲通道所有残留信号(避免旧信号干扰)
+func (ip *IpCast) clearIdleChan() {
+	for {
+		select {
+		case <-ip.idleChan:
+		default:
+			return
+		}
+	}
+}
+
+// listenSerialResponse 后台监听串口返回(仅处理0x41/0x4F,清空残留信号)
+func (ip *IpCast) listenSerialResponse() {
+	buf := make([]byte, 64)
+	for {
+		// 串口未就绪时低频轮询
+		if ip.isClosed || ip.serialPort == nil {
+			time.Sleep(time.Second * 1)
+			continue
+		}
+
+		n, err := ip.serialPort.Read(buf)
+		// 忽略EOF和读取超时(正常无数据场景)
+		if err != nil {
+			switch {
+			case err == io.EOF:
+				continue
+			case err.Error() == "serial: read timed out":
+				continue
+			default:
+				fmt.Printf("串口读取异常: %v\n", err)
+				ip.mu.Lock()
+				ip.isClosed = true
+				ip.mu.Unlock()
+			}
+			continue
+		}
+
+		if n == 0 {
+			continue
+		}
+
+		// 解析返回字节
+		recvBytes := buf[:n]
+		fmt.Printf("串口返回原始字节(十六进制): %s\n", hex.EncodeToString(recvBytes))
+		for _, b := range recvBytes {
+			switch b {
+			case 0x41:
+				fmt.Println("<---- 41  接收成功")
+			case 0x4F:
+				fmt.Println("<---- 4F  芯片空闲")
+				// 清空残留信号后发送新的空闲标记
+				ip.clearIdleChan()
+				ip.idleChan <- struct{}{}
+			}
+		}
+	}
+}
+
+// Speak 发送语音指令(核心逻辑:播放中丢弃请求+仅空闲时接收+无缓存)
+func (ip *IpCast) Speak(txt string) {
+	// 1. 加锁检查播放状态,核心控制逻辑
+	ip.mu.Lock()
+	// 若正在播放,直接丢弃当前雷达触发请求
+	if ip.isPlaying {
+		fmt.Printf("语音正在播放中,丢弃雷达触发请求:%s\n", txt)
+		ip.mu.Unlock()
 		return
 	}
+	// 标记为正在播放(后续只有收到0x4F才会重置)
+	ip.isPlaying = true
+	ip.mu.Unlock()
+
+	// 2. 函数退出时保证重置播放状态(无论成功/失败/超时)
+	defer func() {
+		ip.mu.Lock()
+		ip.isPlaying = false
+		ip.mu.Unlock()
+		fmt.Println("语音播放流程结束,恢复接收雷达信号")
+	}()
+
+	// 3. 校验串口状态
+	ip.mu.Lock()
+	isClosed := ip.isClosed
+	serialPort := ip.serialPort
+	ip.mu.Unlock()
+
+	if isClosed || serialPort == nil {
+		fmt.Println("串口未连接,尝试重连...")
+		if err := ip.Reconnect(); err != nil {
+			fmt.Printf("串口重连失败,无法发送指令: %v\n", err)
+			return
+		}
+		// 重连后重新获取串口实例
+		ip.mu.Lock()
+		serialPort = ip.serialPort
+		ip.mu.Unlock()
+	}
+
+	// 4. 清空空闲通道残留信号(避免旧信号干扰)
+	ip.clearIdleChan()
 
-	req, _ := http.NewRequest("POST", fmt.Sprintf("http://%s/v1/speech", ip.Ip), bytes.NewReader(body))
-	req.Header.Set("Content-Type", "application/json")
-	rsp, err := http.DefaultClient.Do(req)
+	// 5. 文本转GBK编码
+	GBKBytes, err := convertToGBK(txt)
 	if err != nil {
-		logrus.Errorf("IpCast Speak Do err : %s", err.Error())
+		fmt.Printf("文本转GBK失败: %v\n", err)
 		return
 	}
-	rspData, err := io.ReadAll(rsp.Body)
-	logrus.Debugf("IpCast Speak rsp : %+v", string(rspData))
 
-}
+	// 6. 构造语音帧数据(匹配官方格式)
+	dataAreaLen := uint16(1 + 1 + len(GBKBytes)) // 命令字+编码格式+文本长度
+	frameBuf := bytes.NewBuffer([]byte{0xFD})    // 帧头FD
+	_ = binary.Write(frameBuf, binary.BigEndian, dataAreaLen)
+	frameBuf.Write([]byte{0x01, 0x01}) // 命令字01 + 编码格式01(GBK)
+	frameBuf.Write(GBKBytes)
+
+	// 调试打印帧数据
+	//frameHex := hex.EncodeToString(frameBuf.Bytes())
 
-func (ip *IpCast) Reconnect() {
-	c := http.DefaultClient
-	req, _ := http.NewRequest("GET", fmt.Sprintf("http://%s/v1/check_alive", ip.Ip), nil)
-	rsp, err := c.Do(req)
+	// 7. 发送帧数据到串口
+	_, err = serialPort.Write(frameBuf.Bytes())
 	if err != nil {
-		logrus.Errorf("IpCast Reconnect err : %s", err.Error())
+		fmt.Printf("串口发送失败: %v\n", err)
+		ip.mu.Lock()
+		ip.isClosed = true
+		ip.mu.Unlock()
+		_ = ip.Reconnect()
 		return
 	}
-	if rsp.StatusCode == http.StatusOK {
-		ip.liveState = true
+
+	// 8. 等待芯片空闲(0x4F),仅收到空闲信号才允许下一次播放
+	select {
+	case <-ip.idleChan:
+	case <-time.After(time.Second * 10): // 超时延长至10秒,适配长语音
 	}
 }
 
+// CorrectTime 保留原有方法(空实现)
 func (ip IpCast) CorrectTime() {}
+
+// convertToGBK 文本转GBK编码(保留原有逻辑)
+func convertToGBK(s string) ([]byte, error) {
+	if s == "" {
+		return nil, errors.New("文本不能为空")
+	}
+
+	encoder := simplifiedchinese.GBK.NewEncoder()
+	reader := transform.NewReader(bytes.NewReader([]byte(s)), encoder)
+	result, err := io.ReadAll(reader)
+	if err != nil {
+		return nil, fmt.Errorf("编码转换失败: %w", err)
+	}
+
+	fmt.Printf("文本「%s」的GBK编码(十六进制): %s\n", s, hex.EncodeToString(result))
+	if len(result) == 0 {
+		return nil, errors.New("GBK转换结果为空")
+	}
+	return result, nil
+}
+
+// Close 手动关闭串口(清理资源)
+func (ip *IpCast) Close() {
+	ip.mu.Lock()
+	defer ip.mu.Unlock()
+
+	ip.isClosed = true
+	ip.isPlaying = false // 强制重置播放状态
+	ip.clearIdleChan()   // 清空通道
+	if ip.serialPort != nil {
+		_ = ip.serialPort.Close()
+		ip.serialPort = nil
+	}
+	fmt.Println("串口已手动关闭")
+}

+ 16 - 11
util/config.go

@@ -21,22 +21,27 @@ var Config = func() config {
 }()
 
 type config struct {
-	HikServer hikServer           `yaml:"hikServer"`
-	Cameras   []model.CameraInfo  `yaml:"cameras"`
-	Radars    []model.RadarInfo   `yaml:"radars"`
-	Screens   []model.ScreenInfo  `yaml:"screens"`
-	Speakers  []model.SpeakerInfo `yaml:"speakers"`
-	Server    service             `yaml:"service"`
+	Cameras     []model.CameraInfo  `yaml:"cameras"`
+	Radars      []model.RadarInfo   `yaml:"radars"`
+	Screens     []model.ScreenInfo  `yaml:"screens"`
+	Speakers    []model.SpeakerInfo `yaml:"speakers"`
+	Server      server              `yaml:"server"`
+	RedisConfig RedisConfig         `yaml:"redis"`
 }
 
-type service struct {
-	SupportRadar   bool `yaml:"support_radar"`
-	SupportCamera  bool `yaml:"support_camera"`
-	SupportSpeaker bool `yaml:"support_speaker"`
-	SupportLed     bool `yaml:"support_led"`
+type server struct {
+	SupportRadar  bool      `yaml:"supportRadar"`
+	SupportCamera bool      `yaml:"supportCamera"`
+	HikServer     hikServer `yaml:"hikServer"`
 }
 
 type hikServer struct {
 	Addr string `yaml:"addr"`
 	Path string `yaml:"path"`
 }
+
+type RedisConfig struct {
+	Addr     string `yaml:"addr"`
+	Password string `yaml:"password"`
+	DB       int    `yaml:"db"`
+}

+ 60 - 16
util/logrus.go

@@ -1,29 +1,73 @@
 package util
 
 import (
-	rotatelogs "github.com/lestrrat/go-file-rotatelogs"
-	"github.com/sirupsen/logrus"
+	"io"
 	"os"
 	"path"
 	"time"
+
+	rotatelogs "github.com/lestrrat/go-file-rotatelogs"
+	"github.com/sirupsen/logrus"
 )
 
-var _ = func() error {
-	err := os.MkdirAll("./log", os.ModeDir)
+func init() {
+	// ========== 核心配置 ==========
+	logDir := "./log"
+	dirPerm := os.FileMode(0777)  // 目录权限:pi可读写执行,同组可读执行
+	filePerm := os.FileMode(0777) // 文件权限:pi可读写,其他可读
+	today := time.Now().Format("20060102")
+	todayLogFile := path.Join(logDir, "info."+today+".log")
+	logFileTemplate := path.Join(logDir, "info.%Y%m%d.log")
+
+	// ========== 1. 强制创建目录 + 修复目录归属/权限 ==========
+	if err := os.MkdirAll(logDir, dirPerm); err != nil {
+		logrus.Fatalf("[日志初始化] 创建日志目录失败: %v", err)
+	}
+	// 关键:强制将目录归属当前运行用户(pi),无论之前归属谁
+	if err := os.Chown(logDir, os.Getuid(), os.Getgid()); err != nil {
+		logrus.Warnf("[日志初始化] 修正目录归属失败(非致命): %v", err)
+	}
+	// 二次校验目录权限(确保不是只读)
+	if err := os.Chmod(logDir, dirPerm); err != nil {
+		logrus.Warnf("[日志初始化] 修正目录权限失败(非致命): %v", err)
+	}
+
+	// ========== 2. 强制创建今日日志文件 + 修复文件归属/权限 ==========
+	// 打开文件(创建+追加+写入,确保文件存在且权限正确)
+	f, err := os.OpenFile(todayLogFile, os.O_CREATE|os.O_WRONLY|os.O_APPEND, filePerm)
 	if err != nil {
-		logrus.Error("创建日志目录失败!", err)
-		return err
+		logrus.Fatalf("[日志初始化] 创建日志文件失败: %v", err)
+	}
+	f.Close() // 关闭临时句柄
+
+	// 强制将文件归属当前用户(核心:解决sudo残留的root归属问题)
+	if err := os.Chown(todayLogFile, os.Getuid(), os.Getgid()); err != nil {
+		logrus.Warnf("[日志初始化] 修正文件归属失败(非致命): %v", err)
 	}
-	fileName := path.Join("./log", "info")
-	writer, _ := rotatelogs.New(
-		fileName+".%Y%m%d.log",
-		rotatelogs.WithMaxAge(5*24*time.Hour),     // 文件最大保存时间
-		rotatelogs.WithRotationTime(24*time.Hour), // 日志切割时间间隔
+	// 二次校验文件权限
+	if err := os.Chmod(todayLogFile, filePerm); err != nil {
+		logrus.Warnf("[日志初始化] 修正文件权限失败(非致命): %v", err)
+	}
+
+	// ========== 3. 初始化日志切割(适配旧版本rotatelogs) ==========
+	writer, err := rotatelogs.New(
+		logFileTemplate,
+		rotatelogs.WithMaxAge(5*24*time.Hour),
+		rotatelogs.WithRotationTime(24*time.Hour),
+		rotatelogs.WithLinkName(path.Join(logDir, "info.log")),
 	)
-	logrus.SetFormatter(&logrus.JSONFormatter{})
+	if err != nil {
+		logrus.Fatalf("[日志初始化] 创建日志切割器失败: %v", err)
+	}
+
+	// ========== 4. 日志格式配置 ==========
+	logrus.SetFormatter(&logrus.JSONFormatter{
+		TimestampFormat: "2006-01-02 15:04:05", // 标准Go时间模板
+	})
 	logrus.SetLevel(logrus.DebugLevel)
-	logrus.SetOutput(os.Stdout)
 	logrus.SetReportCaller(true)
-	logrus.SetOutput(writer)
-	return nil
-}()
+	// 同时输出到控制台+文件
+	logrus.SetOutput(io.MultiWriter(os.Stdout, writer))
+
+	logrus.Info("[日志初始化] 日志配置完成,权限已强制适配pi用户")
+}