| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167 |
- package parking
- import (
- "context"
- "errors"
- "fmt"
- "io"
- "net/http"
- "net/url"
- "os"
- "os/exec"
- "path/filepath"
- "strings"
- "sync"
- "time"
- "github.com/gofrs/uuid/v5"
- "go.uber.org/zap"
- "wails-app/internal/dao"
- "wails-app/internal/global"
- )
- const cameraStreamTokenTTL = 2 * time.Minute
- type cameraStreamToken struct {
- DeviceCode string
- ExpiresAt time.Time
- }
- var cameraStreamTokens sync.Map
- func isSystemMJPEG(protocol string) bool {
- protocol = strings.ToLower(strings.TrimSpace(protocol))
- return protocol == "system-mjpeg" || protocol == "rtsp" || protocol == "system"
- }
- func issueCameraStreamToken(deviceCode string) string {
- for {
- id, err := uuid.NewV4()
- if err != nil {
- continue
- }
- token := id.String()
- cameraStreamTokens.Store(token, cameraStreamToken{
- DeviceCode: deviceCode,
- ExpiresAt: time.Now().Add(cameraStreamTokenTTL),
- })
- return token
- }
- }
- func validateCameraStreamToken(deviceCode, token string) bool {
- value, ok := cameraStreamTokens.Load(token)
- if !ok {
- return false
- }
- entry := value.(cameraStreamToken)
- if entry.DeviceCode != deviceCode || time.Now().After(entry.ExpiresAt) {
- cameraStreamTokens.Delete(token)
- return false
- }
- return true
- }
- // StreamCameraMJPEG 在系统端把摄像头原始 RTSP/HTTP 流转换为 multipart MJPEG。
- // 该输出可由 WebView 的 <img> 直接播放,不要求边缘端实现 WHEP。
- func (s *PassageService) StreamCameraMJPEG(deviceCode, token string, writer http.ResponseWriter, ctx context.Context) error {
- if !validateCameraStreamToken(deviceCode, token) {
- return errors.New("摄像头播放令牌无效或已过期")
- }
- if global.GVA_DB == nil {
- return errors.New("数据库未初始化")
- }
- var camera dao.Camera
- if err := global.GVA_DB.Where("device_code = ? AND is_active = ?", deviceCode, true).First(&camera).Error; err != nil {
- return fmt.Errorf("摄像头不存在或未启用: %w", err)
- }
- if !isSystemMJPEG(camera.StreamProtocol) {
- return errors.New("当前摄像头未启用系统侧转换")
- }
- sourceURL := strings.TrimSpace(camera.SourceURL)
- if sourceURL == "" {
- sourceURL = strings.TrimSpace(camera.StreamURL)
- }
- if sourceURL == "" {
- return errors.New("摄像头未配置原始视频地址")
- }
- ffmpegPath := strings.TrimSpace(global.GVA_CONFIG.Camera.FFmpegPath)
- if ffmpegPath == "" {
- ffmpegPath = "ffmpeg"
- } else if info, statErr := os.Stat(ffmpegPath); statErr == nil && info.IsDir() {
- // Accept a configured FFmpeg bin directory for compatibility with older config files.
- ffmpegPath = filepath.Join(ffmpegPath, "ffmpeg.exe")
- }
- fps := global.GVA_CONFIG.Camera.MJPEGFPS
- if fps < 1 || fps > 25 {
- fps = 8
- }
- quality := global.GVA_CONFIG.Camera.JPEGQuality
- if quality < 2 || quality > 31 {
- quality = 5
- }
- args := []string{"-hide_banner", "-loglevel", "error"}
- if strings.HasPrefix(strings.ToLower(sourceURL), "rtsp://") {
- args = append(args, "-rtsp_transport", "tcp")
- }
- args = append(args, "-i", sourceURL, "-an", "-vf", fmt.Sprintf("fps=%d", fps), "-q:v", fmt.Sprintf("%d", quality), "-f", "mpjpeg", "pipe:1")
- cmd := exec.CommandContext(ctx, ffmpegPath, args...)
- if global.GVA_LOG != nil {
- global.GVA_LOG.Info("启动摄像头 FFmpeg 转码", zap.String("device_code", deviceCode), zap.String("ffmpeg", ffmpegPath), zap.String("source", redactCameraURL(sourceURL)))
- }
- stdout, err := cmd.StdoutPipe()
- if err != nil {
- return fmt.Errorf("创建视频转换输出失败: %w", err)
- }
- cmd.Stderr = io.Discard
- if err := cmd.Start(); err != nil {
- if global.GVA_LOG != nil {
- global.GVA_LOG.Error("摄像头 FFmpeg 启动失败", zap.String("device_code", deviceCode), zap.String("ffmpeg", ffmpegPath), zap.Error(err))
- }
- return fmt.Errorf("启动 FFmpeg 失败(%s): %w", ffmpegPath, err)
- }
- writer.Header().Set("Content-Type", "multipart/x-mixed-replace; boundary=ffmpeg")
- writer.Header().Set("Cache-Control", "no-store, no-cache, must-revalidate")
- writer.Header().Set("X-Content-Type-Options", "nosniff")
- if flusher, ok := writer.(http.Flusher); ok {
- flusher.Flush()
- }
- _, copyErr := io.Copy(&flushWriter{writer: writer}, stdout)
- waitErr := cmd.Wait()
- if copyErr != nil && !errors.Is(copyErr, context.Canceled) {
- return fmt.Errorf("视频流传输失败: %w", copyErr)
- }
- if ctx.Err() != nil {
- return nil
- }
- if waitErr != nil {
- if global.GVA_LOG != nil {
- global.GVA_LOG.Error("摄像头 FFmpeg 转码进程退出", zap.String("device_code", deviceCode), zap.Error(waitErr))
- }
- return fmt.Errorf("FFmpeg 进程退出: %w", waitErr)
- }
- return nil
- }
- func redactCameraURL(raw string) string {
- u, err := url.Parse(raw)
- if err != nil || u.User == nil {
- return raw
- }
- return strings.Replace(raw, u.User.String()+"@", "***@", 1)
- }
- type flushWriter struct {
- writer http.ResponseWriter
- }
- func (w *flushWriter) Write(data []byte) (int, error) {
- n, err := w.writer.Write(data)
- if flusher, ok := w.writer.(http.Flusher); ok {
- flusher.Flush()
- }
- return n, err
- }
|