camera_stream.go 4.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167
  1. package parking
  2. import (
  3. "context"
  4. "errors"
  5. "fmt"
  6. "io"
  7. "net/http"
  8. "net/url"
  9. "os"
  10. "os/exec"
  11. "path/filepath"
  12. "strings"
  13. "sync"
  14. "time"
  15. "github.com/gofrs/uuid/v5"
  16. "go.uber.org/zap"
  17. "wails-app/internal/dao"
  18. "wails-app/internal/global"
  19. )
  20. const cameraStreamTokenTTL = 2 * time.Minute
  21. type cameraStreamToken struct {
  22. DeviceCode string
  23. ExpiresAt time.Time
  24. }
  25. var cameraStreamTokens sync.Map
  26. func isSystemMJPEG(protocol string) bool {
  27. protocol = strings.ToLower(strings.TrimSpace(protocol))
  28. return protocol == "system-mjpeg" || protocol == "rtsp" || protocol == "system"
  29. }
  30. func issueCameraStreamToken(deviceCode string) string {
  31. for {
  32. id, err := uuid.NewV4()
  33. if err != nil {
  34. continue
  35. }
  36. token := id.String()
  37. cameraStreamTokens.Store(token, cameraStreamToken{
  38. DeviceCode: deviceCode,
  39. ExpiresAt: time.Now().Add(cameraStreamTokenTTL),
  40. })
  41. return token
  42. }
  43. }
  44. func validateCameraStreamToken(deviceCode, token string) bool {
  45. value, ok := cameraStreamTokens.Load(token)
  46. if !ok {
  47. return false
  48. }
  49. entry := value.(cameraStreamToken)
  50. if entry.DeviceCode != deviceCode || time.Now().After(entry.ExpiresAt) {
  51. cameraStreamTokens.Delete(token)
  52. return false
  53. }
  54. return true
  55. }
  56. // StreamCameraMJPEG 在系统端把摄像头原始 RTSP/HTTP 流转换为 multipart MJPEG。
  57. // 该输出可由 WebView 的 <img> 直接播放,不要求边缘端实现 WHEP。
  58. func (s *PassageService) StreamCameraMJPEG(deviceCode, token string, writer http.ResponseWriter, ctx context.Context) error {
  59. if !validateCameraStreamToken(deviceCode, token) {
  60. return errors.New("摄像头播放令牌无效或已过期")
  61. }
  62. if global.GVA_DB == nil {
  63. return errors.New("数据库未初始化")
  64. }
  65. var camera dao.Camera
  66. if err := global.GVA_DB.Where("device_code = ? AND is_active = ?", deviceCode, true).First(&camera).Error; err != nil {
  67. return fmt.Errorf("摄像头不存在或未启用: %w", err)
  68. }
  69. if !isSystemMJPEG(camera.StreamProtocol) {
  70. return errors.New("当前摄像头未启用系统侧转换")
  71. }
  72. sourceURL := strings.TrimSpace(camera.SourceURL)
  73. if sourceURL == "" {
  74. sourceURL = strings.TrimSpace(camera.StreamURL)
  75. }
  76. if sourceURL == "" {
  77. return errors.New("摄像头未配置原始视频地址")
  78. }
  79. ffmpegPath := strings.TrimSpace(global.GVA_CONFIG.Camera.FFmpegPath)
  80. if ffmpegPath == "" {
  81. ffmpegPath = "ffmpeg"
  82. } else if info, statErr := os.Stat(ffmpegPath); statErr == nil && info.IsDir() {
  83. // Accept a configured FFmpeg bin directory for compatibility with older config files.
  84. ffmpegPath = filepath.Join(ffmpegPath, "ffmpeg.exe")
  85. }
  86. fps := global.GVA_CONFIG.Camera.MJPEGFPS
  87. if fps < 1 || fps > 25 {
  88. fps = 8
  89. }
  90. quality := global.GVA_CONFIG.Camera.JPEGQuality
  91. if quality < 2 || quality > 31 {
  92. quality = 5
  93. }
  94. args := []string{"-hide_banner", "-loglevel", "error"}
  95. if strings.HasPrefix(strings.ToLower(sourceURL), "rtsp://") {
  96. args = append(args, "-rtsp_transport", "tcp")
  97. }
  98. args = append(args, "-i", sourceURL, "-an", "-vf", fmt.Sprintf("fps=%d", fps), "-q:v", fmt.Sprintf("%d", quality), "-f", "mpjpeg", "pipe:1")
  99. cmd := exec.CommandContext(ctx, ffmpegPath, args...)
  100. if global.GVA_LOG != nil {
  101. global.GVA_LOG.Info("启动摄像头 FFmpeg 转码", zap.String("device_code", deviceCode), zap.String("ffmpeg", ffmpegPath), zap.String("source", redactCameraURL(sourceURL)))
  102. }
  103. stdout, err := cmd.StdoutPipe()
  104. if err != nil {
  105. return fmt.Errorf("创建视频转换输出失败: %w", err)
  106. }
  107. cmd.Stderr = io.Discard
  108. if err := cmd.Start(); err != nil {
  109. if global.GVA_LOG != nil {
  110. global.GVA_LOG.Error("摄像头 FFmpeg 启动失败", zap.String("device_code", deviceCode), zap.String("ffmpeg", ffmpegPath), zap.Error(err))
  111. }
  112. return fmt.Errorf("启动 FFmpeg 失败(%s): %w", ffmpegPath, err)
  113. }
  114. writer.Header().Set("Content-Type", "multipart/x-mixed-replace; boundary=ffmpeg")
  115. writer.Header().Set("Cache-Control", "no-store, no-cache, must-revalidate")
  116. writer.Header().Set("X-Content-Type-Options", "nosniff")
  117. if flusher, ok := writer.(http.Flusher); ok {
  118. flusher.Flush()
  119. }
  120. _, copyErr := io.Copy(&flushWriter{writer: writer}, stdout)
  121. waitErr := cmd.Wait()
  122. if copyErr != nil && !errors.Is(copyErr, context.Canceled) {
  123. return fmt.Errorf("视频流传输失败: %w", copyErr)
  124. }
  125. if ctx.Err() != nil {
  126. return nil
  127. }
  128. if waitErr != nil {
  129. if global.GVA_LOG != nil {
  130. global.GVA_LOG.Error("摄像头 FFmpeg 转码进程退出", zap.String("device_code", deviceCode), zap.Error(waitErr))
  131. }
  132. return fmt.Errorf("FFmpeg 进程退出: %w", waitErr)
  133. }
  134. return nil
  135. }
  136. func redactCameraURL(raw string) string {
  137. u, err := url.Parse(raw)
  138. if err != nil || u.User == nil {
  139. return raw
  140. }
  141. return strings.Replace(raw, u.User.String()+"@", "***@", 1)
  142. }
  143. type flushWriter struct {
  144. writer http.ResponseWriter
  145. }
  146. func (w *flushWriter) Write(data []byte) (int, error) {
  147. n, err := w.writer.Write(data)
  148. if flusher, ok := w.writer.(http.Flusher); ok {
  149. flusher.Flush()
  150. }
  151. return n, err
  152. }