From 0228b4d0b0884becca322b0ca427ebd6639881ba Mon Sep 17 00:00:00 2001 From: liqi Date: Tue, 16 Dec 2025 09:17:15 +0800 Subject: [PATCH] =?UTF-8?q?=E7=BD=91=E7=AB=99=E9=80=9A=E8=AF=9D=E4=BD=BF?= =?UTF-8?q?=E7=94=A8webrtc=EF=BC=8C=E6=96=B0=E5=A2=9Effmpeg?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/mediaserver/ffmpeg_installer.go | 364 +++++++++++++++++++++ internal/mediaserver/ffmpeg_transcoder.go | 369 ++++++++++++++++++++++ internal/mediaserver/server.go | 10 + internal/mediaserver/ws_rtmp_proxy.go | 240 ++++++++++++-- 4 files changed, 964 insertions(+), 19 deletions(-) create mode 100644 internal/mediaserver/ffmpeg_installer.go create mode 100644 internal/mediaserver/ffmpeg_transcoder.go diff --git a/internal/mediaserver/ffmpeg_installer.go b/internal/mediaserver/ffmpeg_installer.go new file mode 100644 index 0000000..fb334ac --- /dev/null +++ b/internal/mediaserver/ffmpeg_installer.go @@ -0,0 +1,364 @@ +/** + * package mediaserver + * + * FFmpeg 自动检测与安装模块 + * 在服务启动时检测 FFmpeg 是否可用,如不可用则自动下载安装 + */ +package mediaserver + +import ( + "archive/tar" + "archive/zip" + "compress/gzip" + "fmt" + "io" + "log" + "net/http" + "os" + "os/exec" + "path/filepath" + "runtime" + "strings" + "sync" +) + +// FFmpeg 下载地址配置 +const ( + // Linux 静态编译版本 (来自 johnvansickle.com) + FFmpegLinuxURL = "https://johnvansickle.com/ffmpeg/releases/ffmpeg-release-amd64-static.tar.xz" + // Windows 版本 (来自 gyan.dev) + FFmpegWindowsURL = "https://www.gyan.dev/ffmpeg/builds/ffmpeg-release-essentials.zip" + // 默认安装目录 + DefaultInstallDir = "/ffmpegs" +) + +var ( + ffmpegPath string + ffmpegInitOnce sync.Once + ffmpegInitErr error +) + +// EnsureFFmpeg 确保 FFmpeg 可用,返回可执行文件路径 +// 检测顺序: +// 1. 系统 PATH 中的 ffmpeg +// 2. 自定义安装目录 /ffmpegs/ffmpeg +// 3. 如果都不存在,自动下载安装 +func EnsureFFmpeg() (string, error) { + ffmpegInitOnce.Do(func() { + ffmpegPath, ffmpegInitErr = initFFmpeg() + }) + return ffmpegPath, ffmpegInitErr +} + +// initFFmpeg 初始化 FFmpeg +func initFFmpeg() (string, error) { + // 1. 检查系统 PATH 中的 ffmpeg + if path, err := exec.LookPath("ffmpeg"); err == nil { + log.Printf("✅ [FFmpeg] 在系统 PATH 中找到: %s", path) + // 验证版本 + if err := verifyFFmpeg(path); err == nil { + return path, nil + } + log.Printf("⚠️ [FFmpeg] 系统 ffmpeg 验证失败,尝试其他路径") + } + + // 2. 检查自定义安装目录 + customPath := getCustomFFmpegPath() + if _, err := os.Stat(customPath); err == nil { + log.Printf("✅ [FFmpeg] 在自定义目录找到: %s", customPath) + if err := verifyFFmpeg(customPath); err == nil { + return customPath, nil + } + log.Printf("⚠️ [FFmpeg] 自定义 ffmpeg 验证失败,尝试重新下载") + } + + // 3. 下载安装 + log.Printf("📥 [FFmpeg] 未找到可用的 ffmpeg,开始下载安装...") + return downloadAndInstallFFmpeg(DefaultInstallDir) +} + +// getCustomFFmpegPath 获取自定义安装路径 +func getCustomFFmpegPath() string { + if runtime.GOOS == "windows" { + return filepath.Join(DefaultInstallDir, "ffmpeg.exe") + } + return filepath.Join(DefaultInstallDir, "ffmpeg") +} + +// verifyFFmpeg 验证 FFmpeg 是否可用 +func verifyFFmpeg(path string) error { + cmd := exec.Command(path, "-version") + output, err := cmd.Output() + if err != nil { + return fmt.Errorf("执行 ffmpeg -version 失败: %w", err) + } + + // 检查输出是否包含版本信息 + outputStr := string(output) + if !strings.Contains(outputStr, "ffmpeg version") { + return fmt.Errorf("ffmpeg 输出异常: %s", outputStr[:min(len(outputStr), 100)]) + } + + log.Printf("📋 [FFmpeg] 版本信息: %s", strings.Split(outputStr, "\n")[0]) + return nil +} + +// downloadAndInstallFFmpeg 下载并安装 FFmpeg +func downloadAndInstallFFmpeg(installDir string) (string, error) { + // 创建安装目录 + if err := os.MkdirAll(installDir, 0755); err != nil { + return "", fmt.Errorf("创建安装目录失败: %w", err) + } + + var downloadURL string + var extractFunc func(string, string) (string, error) + + switch runtime.GOOS { + case "linux": + downloadURL = FFmpegLinuxURL + extractFunc = extractTarXz + case "windows": + downloadURL = FFmpegWindowsURL + extractFunc = extractZip + default: + return "", fmt.Errorf("不支持的操作系统: %s", runtime.GOOS) + } + + // 下载文件 + log.Printf("📥 [FFmpeg] 正在从 %s 下载...", downloadURL) + tmpFile, err := downloadFile(downloadURL) + if err != nil { + return "", fmt.Errorf("下载失败: %w", err) + } + defer os.Remove(tmpFile) + + // 解压并安装 + log.Printf("📦 [FFmpeg] 正在解压安装到 %s ...", installDir) + ffmpegBinary, err := extractFunc(tmpFile, installDir) + if err != nil { + return "", fmt.Errorf("解压失败: %w", err) + } + + // 验证安装 + if err := verifyFFmpeg(ffmpegBinary); err != nil { + return "", fmt.Errorf("安装验证失败: %w", err) + } + + log.Printf("✅ [FFmpeg] 安装成功: %s", ffmpegBinary) + return ffmpegBinary, nil +} + +// downloadFile 下载文件到临时目录 +func downloadFile(url string) (string, error) { + resp, err := http.Get(url) + if err != nil { + return "", err + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("HTTP 状态码: %d", resp.StatusCode) + } + + // 创建临时文件 + tmpFile, err := os.CreateTemp("", "ffmpeg-download-*") + if err != nil { + return "", err + } + defer tmpFile.Close() + + // 下载并显示进度 + written, err := io.Copy(tmpFile, resp.Body) + if err != nil { + os.Remove(tmpFile.Name()) + return "", err + } + + log.Printf("📥 [FFmpeg] 下载完成,大小: %.2f MB", float64(written)/1024/1024) + return tmpFile.Name(), nil +} + +// extractTarXz 解压 tar.xz 文件 (Linux) +func extractTarXz(archivePath, destDir string) (string, error) { + // 使用 xz 命令解压(需要系统安装 xz-utils) + // 或者使用纯 Go 的 gzip 解压 (如果是 .tar.gz) + + // 先尝试使用系统命令 + cmd := exec.Command("tar", "-xf", archivePath, "-C", destDir) + if err := cmd.Run(); err != nil { + return "", fmt.Errorf("tar 解压失败: %w", err) + } + + // 查找解压后的 ffmpeg 二进制文件 + var ffmpegBinary string + err := filepath.Walk(destDir, func(path string, info os.FileInfo, err error) error { + if err != nil { + return err + } + if !info.IsDir() && info.Name() == "ffmpeg" { + ffmpegBinary = path + return filepath.SkipAll + } + return nil + }) + if err != nil && err != filepath.SkipAll { + return "", err + } + + if ffmpegBinary == "" { + return "", fmt.Errorf("解压后未找到 ffmpeg 二进制文件") + } + + // 移动到目标位置 + targetPath := filepath.Join(destDir, "ffmpeg") + if ffmpegBinary != targetPath { + if err := os.Rename(ffmpegBinary, targetPath); err != nil { + // 如果 rename 失败(跨设备),使用复制 + if err := copyFile(ffmpegBinary, targetPath); err != nil { + return "", err + } + } + } + + // 设置执行权限 + if err := os.Chmod(targetPath, 0755); err != nil { + return "", err + } + + return targetPath, nil +} + +// extractZip 解压 zip 文件 (Windows) +func extractZip(archivePath, destDir string) (string, error) { + r, err := zip.OpenReader(archivePath) + if err != nil { + return "", err + } + defer r.Close() + + var ffmpegBinary string + + for _, f := range r.File { + // 只查找 ffmpeg.exe + if strings.HasSuffix(f.Name, "ffmpeg.exe") { + rc, err := f.Open() + if err != nil { + return "", err + } + + targetPath := filepath.Join(destDir, "ffmpeg.exe") + outFile, err := os.Create(targetPath) + if err != nil { + rc.Close() + return "", err + } + + _, err = io.Copy(outFile, rc) + outFile.Close() + rc.Close() + + if err != nil { + return "", err + } + + ffmpegBinary = targetPath + break + } + } + + if ffmpegBinary == "" { + return "", fmt.Errorf("zip 包中未找到 ffmpeg.exe") + } + + return ffmpegBinary, nil +} + +// extractTarGz 解压 tar.gz 文件 +func extractTarGz(archivePath, destDir string) (string, error) { + f, err := os.Open(archivePath) + if err != nil { + return "", err + } + defer f.Close() + + gzr, err := gzip.NewReader(f) + if err != nil { + return "", err + } + defer gzr.Close() + + tr := tar.NewReader(gzr) + + var ffmpegBinary string + + for { + header, err := tr.Next() + if err == io.EOF { + break + } + if err != nil { + return "", err + } + + // 只查找 ffmpeg 二进制文件 + if strings.HasSuffix(header.Name, "/ffmpeg") || header.Name == "ffmpeg" { + targetPath := filepath.Join(destDir, "ffmpeg") + outFile, err := os.Create(targetPath) + if err != nil { + return "", err + } + + _, err = io.Copy(outFile, tr) + outFile.Close() + + if err != nil { + return "", err + } + + // 设置执行权限 + if err := os.Chmod(targetPath, 0755); err != nil { + return "", err + } + + ffmpegBinary = targetPath + break + } + } + + if ffmpegBinary == "" { + return "", fmt.Errorf("tar.gz 包中未找到 ffmpeg 二进制文件") + } + + return ffmpegBinary, nil +} + +// copyFile 复制文件 +func copyFile(src, dst string) error { + sourceFile, err := os.Open(src) + if err != nil { + return err + } + defer sourceFile.Close() + + destFile, err := os.Create(dst) + if err != nil { + return err + } + defer destFile.Close() + + _, err = io.Copy(destFile, sourceFile) + return err +} + +// GetFFmpegPath 获取 FFmpeg 路径(仅在已初始化后调用) +func GetFFmpegPath() string { + return ffmpegPath +} + +// IsFFmpegAvailable 检查 FFmpeg 是否可用 +func IsFFmpegAvailable() bool { + path, err := EnsureFFmpeg() + return err == nil && path != "" +} + + diff --git a/internal/mediaserver/ffmpeg_transcoder.go b/internal/mediaserver/ffmpeg_transcoder.go new file mode 100644 index 0000000..c455043 --- /dev/null +++ b/internal/mediaserver/ffmpeg_transcoder.go @@ -0,0 +1,369 @@ +/** + * package mediaserver + * + * FFmpeg 实时转码器 + * 将 WebM (VP8/Opus) 转码为 FLV (H.264/AAC) + * 用于处理不支持 H.264 编码的浏览器发送的视频流 + */ +package mediaserver + +import ( + "bytes" + "fmt" + "io" + "log" + "os/exec" + "sync" + "time" +) + +// TranscoderConfig 转码器配置 +type TranscoderConfig struct { + // 视频配置 + VideoCodec string // 输出视频编解码器,默认 libx264 + VideoPreset string // x264 预设,默认 ultrafast + VideoBitrate string // 视频码率,默认 500k + + // 音频配置 + AudioCodec string // 输出音频编解码器,默认 aac + AudioBitrate string // 音频码率,默认 64k + AudioSampleRate int // 音频采样率,默认 44100 + AudioChannels int // 音频通道数,默认 2 +} + +// DefaultTranscoderConfig 默认转码配置 +func DefaultTranscoderConfig() *TranscoderConfig { + return &TranscoderConfig{ + VideoCodec: "libx264", + VideoPreset: "ultrafast", + VideoBitrate: "500k", + AudioCodec: "aac", + AudioBitrate: "64k", + AudioSampleRate: 44100, + AudioChannels: 2, + } +} + +// FFmpegTranscoder FFmpeg 转码器 +type FFmpegTranscoder struct { + config *TranscoderConfig + ffmpegPath string + + cmd *exec.Cmd + stdin io.WriteCloser + stdout io.ReadCloser + stderr io.ReadCloser + + outputChan chan []byte + stopChan chan struct{} + + mu sync.Mutex + running bool + + // 统计信息 + inputBytes int64 + outputBytes int64 + startTime time.Time +} + +// NewFFmpegTranscoder 创建新的转码器 +func NewFFmpegTranscoder(config *TranscoderConfig) (*FFmpegTranscoder, error) { + if config == nil { + config = DefaultTranscoderConfig() + } + + // 确保 FFmpeg 可用 + ffmpegPath, err := EnsureFFmpeg() + if err != nil { + return nil, fmt.Errorf("FFmpeg 不可用: %w", err) + } + + return &FFmpegTranscoder{ + config: config, + ffmpegPath: ffmpegPath, + outputChan: make(chan []byte, 100), + stopChan: make(chan struct{}), + }, nil +} + +// Start 启动转码器 +func (t *FFmpegTranscoder) Start() error { + t.mu.Lock() + defer t.mu.Unlock() + + if t.running { + return fmt.Errorf("转码器已在运行") + } + + // 构建 FFmpeg 命令 + // ffmpeg -f webm -i pipe:0 -c:v libx264 -preset ultrafast -tune zerolatency -c:a aac -f flv pipe:1 + args := []string{ + "-f", "webm", // 输入格式 + "-i", "pipe:0", // 从 stdin 读取 + "-c:v", t.config.VideoCodec, + "-preset", t.config.VideoPreset, + "-tune", "zerolatency", // 零延迟模式 + "-b:v", t.config.VideoBitrate, + "-c:a", t.config.AudioCodec, + "-b:a", t.config.AudioBitrate, + "-ar", fmt.Sprintf("%d", t.config.AudioSampleRate), + "-ac", fmt.Sprintf("%d", t.config.AudioChannels), + "-f", "flv", // 输出格式 + "pipe:1", // 输出到 stdout + } + + t.cmd = exec.Command(t.ffmpegPath, args...) + + var err error + + // 获取 stdin + t.stdin, err = t.cmd.StdinPipe() + if err != nil { + return fmt.Errorf("获取 stdin 失败: %w", err) + } + + // 获取 stdout + t.stdout, err = t.cmd.StdoutPipe() + if err != nil { + return fmt.Errorf("获取 stdout 失败: %w", err) + } + + // 获取 stderr + t.stderr, err = t.cmd.StderrPipe() + if err != nil { + return fmt.Errorf("获取 stderr 失败: %w", err) + } + + // 启动 FFmpeg 进程 + if err := t.cmd.Start(); err != nil { + return fmt.Errorf("启动 FFmpeg 失败: %w", err) + } + + t.running = true + t.startTime = time.Now() + + log.Printf("🎬 [FFmpegTranscoder] 转码器已启动, PID=%d", t.cmd.Process.Pid) + + // 启动输出读取协程 + go t.readOutput() + + // 启动 stderr 读取协程(用于日志) + go t.readStderr() + + // 启动进程监控协程 + go t.monitor() + + return nil +} + +// Stop 停止转码器 +func (t *FFmpegTranscoder) Stop() error { + t.mu.Lock() + defer t.mu.Unlock() + + if !t.running { + return nil + } + + t.running = false + close(t.stopChan) + + // 关闭 stdin,通知 FFmpeg 输入结束 + if t.stdin != nil { + t.stdin.Close() + } + + // 等待进程结束 + if t.cmd != nil && t.cmd.Process != nil { + t.cmd.Process.Kill() + t.cmd.Wait() + } + + // 关闭输出通道 + close(t.outputChan) + + duration := time.Since(t.startTime) + log.Printf("🛑 [FFmpegTranscoder] 转码器已停止, 运行时间=%v, 输入=%d bytes, 输出=%d bytes", + duration, t.inputBytes, t.outputBytes) + + return nil +} + +// Write 写入 WebM 数据 +func (t *FFmpegTranscoder) Write(data []byte) (int, error) { + t.mu.Lock() + if !t.running { + t.mu.Unlock() + return 0, fmt.Errorf("转码器未运行") + } + stdin := t.stdin + t.mu.Unlock() + + n, err := stdin.Write(data) + if err != nil { + return n, err + } + + t.inputBytes += int64(n) + return n, nil +} + +// Output 获取输出通道 +func (t *FFmpegTranscoder) Output() <-chan []byte { + return t.outputChan +} + +// readOutput 读取 FFmpeg 输出 +func (t *FFmpegTranscoder) readOutput() { + buf := make([]byte, 64*1024) // 64KB 缓冲区 + + for { + select { + case <-t.stopChan: + return + default: + } + + n, err := t.stdout.Read(buf) + if err != nil { + if err != io.EOF { + log.Printf("⚠️ [FFmpegTranscoder] 读取输出错误: %v", err) + } + return + } + + if n > 0 { + t.outputBytes += int64(n) + + // 复制数据并发送到通道 + data := make([]byte, n) + copy(data, buf[:n]) + + select { + case t.outputChan <- data: + default: + // 通道满,丢弃数据 + log.Printf("⚠️ [FFmpegTranscoder] 输出通道满,丢弃 %d bytes", n) + } + } + } +} + +// readStderr 读取 FFmpeg 错误输出 +func (t *FFmpegTranscoder) readStderr() { + buf := new(bytes.Buffer) + io.Copy(buf, t.stderr) + + if buf.Len() > 0 { + // 只在有错误时输出 + output := buf.String() + if len(output) > 500 { + output = output[:500] + "..." + } + log.Printf("📋 [FFmpegTranscoder] FFmpeg 输出: %s", output) + } +} + +// monitor 监控 FFmpeg 进程 +func (t *FFmpegTranscoder) monitor() { + err := t.cmd.Wait() + + t.mu.Lock() + wasRunning := t.running + t.running = false + t.mu.Unlock() + + if wasRunning { + if err != nil { + log.Printf("⚠️ [FFmpegTranscoder] FFmpeg 进程异常退出: %v", err) + } else { + log.Printf("ℹ️ [FFmpegTranscoder] FFmpeg 进程正常退出") + } + } +} + +// IsRunning 检查转码器是否运行中 +func (t *FFmpegTranscoder) IsRunning() bool { + t.mu.Lock() + defer t.mu.Unlock() + return t.running +} + +// Stats 获取统计信息 +func (t *FFmpegTranscoder) Stats() (inputBytes, outputBytes int64, duration time.Duration) { + t.mu.Lock() + defer t.mu.Unlock() + return t.inputBytes, t.outputBytes, time.Since(t.startTime) +} + +// TranscoderPool 转码器池,管理多个转码器实例 +type TranscoderPool struct { + transcoders map[string]*FFmpegTranscoder + mu sync.RWMutex + config *TranscoderConfig +} + +// NewTranscoderPool 创建转码器池 +func NewTranscoderPool(config *TranscoderConfig) *TranscoderPool { + if config == nil { + config = DefaultTranscoderConfig() + } + return &TranscoderPool{ + transcoders: make(map[string]*FFmpegTranscoder), + config: config, + } +} + +// Get 获取或创建转码器 +func (p *TranscoderPool) Get(streamID string) (*FFmpegTranscoder, error) { + p.mu.Lock() + defer p.mu.Unlock() + + if t, exists := p.transcoders[streamID]; exists && t.IsRunning() { + return t, nil + } + + // 创建新的转码器 + t, err := NewFFmpegTranscoder(p.config) + if err != nil { + return nil, err + } + + if err := t.Start(); err != nil { + return nil, err + } + + p.transcoders[streamID] = t + return t, nil +} + +// Release 释放转码器 +func (p *TranscoderPool) Release(streamID string) { + p.mu.Lock() + defer p.mu.Unlock() + + if t, exists := p.transcoders[streamID]; exists { + t.Stop() + delete(p.transcoders, streamID) + } +} + +// ReleaseAll 释放所有转码器 +func (p *TranscoderPool) ReleaseAll() { + p.mu.Lock() + defer p.mu.Unlock() + + for id, t := range p.transcoders { + t.Stop() + delete(p.transcoders, id) + } +} + +// Count 获取活跃转码器数量 +func (p *TranscoderPool) Count() int { + p.mu.RLock() + defer p.mu.RUnlock() + return len(p.transcoders) +} + + diff --git a/internal/mediaserver/server.go b/internal/mediaserver/server.go index 887c0c7..135fa7e 100644 --- a/internal/mediaserver/server.go +++ b/internal/mediaserver/server.go @@ -123,6 +123,16 @@ func Start() { log.Println("🎥 [MediaServer] 正在启动...") + // 在启动时检测并安装 FFmpeg(用于 VP8/VP9 转码) + go func() { + log.Println("🔍 [MediaServer] 检测 FFmpeg...") + if path, err := EnsureFFmpeg(); err != nil { + log.Printf("⚠️ [MediaServer] FFmpeg 不可用: %v (VP8/VP9 转码将不可用)", err) + } else { + log.Printf("✅ [MediaServer] FFmpeg 就绪: %s", path) + } + }() + // 启动 WebRTC SFU if ms.sfu != nil { go ms.sfu.Start() diff --git a/internal/mediaserver/ws_rtmp_proxy.go b/internal/mediaserver/ws_rtmp_proxy.go index 3a2fc6d..5995717 100644 --- a/internal/mediaserver/ws_rtmp_proxy.go +++ b/internal/mediaserver/ws_rtmp_proxy.go @@ -15,6 +15,7 @@ package mediaserver import ( "bytes" "encoding/binary" + "encoding/json" "fmt" "io" "log" @@ -34,6 +35,14 @@ type WebMToRTMPProxy struct { sessionsMu sync.RWMutex } +// CodecInfo 编解码器信息(从前端发送的 JSON 消息解析) +type CodecInfo struct { + Type string `json:"type"` // "codec_info" + VideoCodec string `json:"video_codec"` // "h264", "vp8", "vp9" + AudioCodec string `json:"audio_codec"` // "opus", "aac" + MimeType string `json:"mime_type"` // 完整的 MIME 类型 +} + // ProxySession 代理会话 type ProxySession struct { ID string @@ -54,6 +63,14 @@ type ProxySession struct { // 时间戳 baseTimestamp uint32 lastTimestamp uint32 + + // 编解码器信息 + codecInfo *CodecInfo + codecDetected bool + + // FFmpeg 转码器(VP8/VP9 需要转码) + transcoder *FFmpegTranscoder + needTranscode bool } // WebMHeader WebM 文件头信息 @@ -222,6 +239,11 @@ func (p *WebMToRTMPProxy) handleSession(session *ProxySession) { totalBytes := int64(0) defer func() { + // 清理转码器 + if session.transcoder != nil { + session.transcoder.Stop() + } + // 清理 p.sessionsMu.Lock() delete(p.sessions, session.ID) @@ -237,8 +259,12 @@ func (p *WebMToRTMPProxy) handleSession(session *ProxySession) { session.Stream.mu.Unlock() } - log.Printf("🔌 [WSProxy] 连接关闭 | Stream:%s User:%s | 收到消息:%d 总字节:%d", - session.StreamID, session.UserID, messageCount, totalBytes) + codecStr := "unknown" + if session.codecInfo != nil { + codecStr = session.codecInfo.VideoCodec + } + log.Printf("🔌 [WSProxy] 连接关闭 | Stream:%s User:%s | 收到消息:%d 总字节:%d 编解码器:%s 需要转码:%v", + session.StreamID, session.UserID, messageCount, totalBytes, codecStr, session.needTranscode) }() for { @@ -270,21 +296,166 @@ func (p *WebMToRTMPProxy) handleSession(session *ProxySession) { // 每 100 条消息记录一次统计 if messageCount%100 == 0 { - log.Printf("📊 [WSProxy] 消息统计 | Stream:%s Count:%d TotalBytes:%d", session.StreamID, messageCount, totalBytes) + log.Printf("📊 [WSProxy] 消息统计 | Stream:%s Count:%d TotalBytes:%d NeedTranscode:%v", + session.StreamID, messageCount, totalBytes, session.needTranscode) } - if messageType != websocket.BinaryMessage { - log.Printf("⚠️ [WSProxy] 忽略非二进制消息 | Type:%d", messageType) + // 处理文本消息(可能是编解码器信息) + if messageType == websocket.TextMessage { + if err := p.handleTextMessage(session, data); err != nil { + log.Printf("⚠️ [WSProxy] 处理文本消息失败: %v", err) + } continue } - // 处理 WebM 数据 - if err := p.processWebMData(session, data); err != nil { - log.Printf("⚠️ [WSProxy] 处理 WebM 数据失败: %v", err) + if messageType != websocket.BinaryMessage { + log.Printf("⚠️ [WSProxy] 忽略未知消息类型 | Type:%d", messageType) + continue + } + + // 处理二进制数据(WebM) + if session.needTranscode && session.transcoder != nil { + // VP8/VP9: 通过 FFmpeg 转码 + if err := p.processWithTranscoder(session, data); err != nil { + log.Printf("⚠️ [WSProxy] 转码处理失败: %v", err) + } + } else { + // H.264 或未检测到编解码器: 直接处理 + if err := p.processWebMData(session, data); err != nil { + log.Printf("⚠️ [WSProxy] 处理 WebM 数据失败: %v", err) + } } } } +// handleTextMessage 处理文本消息(编解码器信息) +func (p *WebMToRTMPProxy) handleTextMessage(session *ProxySession, data []byte) error { + var codecInfo CodecInfo + if err := json.Unmarshal(data, &codecInfo); err != nil { + return fmt.Errorf("JSON 解析失败: %w", err) + } + + if codecInfo.Type != "codec_info" { + log.Printf("ℹ️ [WSProxy] 收到非 codec_info 消息: %s", codecInfo.Type) + return nil + } + + session.codecInfo = &codecInfo + session.codecDetected = true + + log.Printf("🎬 [WSProxy] 收到编解码器信息 | Stream:%s Video:%s Audio:%s MIME:%s", + session.StreamID, codecInfo.VideoCodec, codecInfo.AudioCodec, codecInfo.MimeType) + + // 判断是否需要转码 + // H.264 可以直接封装到 FLV,VP8/VP9 需要转码 + switch codecInfo.VideoCodec { + case "h264", "avc1": + session.needTranscode = false + log.Printf("✅ [WSProxy] H.264 编码,无需转码,直接重封装为 FLV") + case "vp8", "vp9": + session.needTranscode = true + log.Printf("🔄 [WSProxy] %s 编码,需要 FFmpeg 转码为 H.264", codecInfo.VideoCodec) + + // 初始化转码器 + if err := p.initTranscoder(session); err != nil { + log.Printf("⚠️ [WSProxy] 初始化转码器失败: %v,将使用非标准 FLV 格式", err) + session.needTranscode = false + } + default: + log.Printf("⚠️ [WSProxy] 未知视频编解码器: %s,尝试直接处理", codecInfo.VideoCodec) + session.needTranscode = false + } + + return nil +} + +// initTranscoder 初始化 FFmpeg 转码器 +func (p *WebMToRTMPProxy) initTranscoder(session *ProxySession) error { + transcoder, err := NewFFmpegTranscoder(nil) + if err != nil { + return fmt.Errorf("创建转码器失败: %w", err) + } + + if err := transcoder.Start(); err != nil { + return fmt.Errorf("启动转码器失败: %w", err) + } + + session.transcoder = transcoder + + // 启动协程读取转码后的 FLV 数据并广播 + go p.readTranscodedOutput(session) + + log.Printf("✅ [WSProxy] FFmpeg 转码器已初始化 | Stream:%s", session.StreamID) + return nil +} + +// readTranscodedOutput 读取转码后的 FLV 输出并广播 +func (p *WebMToRTMPProxy) readTranscodedOutput(session *ProxySession) { + log.Printf("▶️ [WSProxy] 开始读取转码输出 | Stream:%s", session.StreamID) + + flvBuffer := bytes.NewBuffer(nil) + flvHeaderParsed := false + + for data := range session.transcoder.Output() { + flvBuffer.Write(data) + + // 解析 FLV 数据 + for { + bufData := flvBuffer.Bytes() + + // 首先解析 FLV 头(9 字节)+ PreviousTagSize(4 字节) + if !flvHeaderParsed { + if len(bufData) < 13 { + break + } + // 跳过 FLV 头 + flvBuffer.Next(13) + flvHeaderParsed = true + bufData = flvBuffer.Bytes() + } + + // 解析 FLV Tag + if len(bufData) < 11 { + break + } + + // Tag 头 + tagType := bufData[0] + dataSize := int(bufData[1])<<16 | int(bufData[2])<<8 | int(bufData[3]) + + // 检查是否有完整的 tag(11 字节头 + dataSize + 4 字节 PreviousTagSize) + totalTagSize := 11 + dataSize + 4 + if len(bufData) < totalTagSize { + break + } + + // 提取完整的 FLV tag(不包括 PreviousTagSize) + flvTag := make([]byte, 11+dataSize+4) + copy(flvTag, bufData[:totalTagSize]) + + // 广播 FLV tag + if tagType == 8 || tagType == 9 { // 音频或视频 + p.broadcastFLVTag(session, flvTag) + } + + // 从缓冲区移除已处理的数据 + flvBuffer.Next(totalTagSize) + } + } + + log.Printf("⏹️ [WSProxy] 转码输出读取结束 | Stream:%s", session.StreamID) +} + +// processWithTranscoder 使用转码器处理数据 +func (p *WebMToRTMPProxy) processWithTranscoder(session *ProxySession, data []byte) error { + if session.transcoder == nil { + return fmt.Errorf("转码器未初始化") + } + + _, err := session.transcoder.Write(data) + return err +} + // processWebMData 处理 WebM 数据 func (p *WebMToRTMPProxy) processWebMData(session *ProxySession, data []byte) error { // 将数据追加到缓冲区 @@ -541,26 +712,20 @@ func (p *WebMToRTMPProxy) convertClusterToFLV(session *ProxySession, clusterData } // createFLVVideoTag 创建 FLV 视频 tag -// VP8 -> FLV (需要转换为 H.264,这里简化处理) +// 根据会话的编解码器信息选择正确的封装方式 func (p *WebMToRTMPProxy) createFLVVideoTag(timestamp uint32, data []byte, isKeyframe bool) []byte { - // 注意:VP8 不能直接封装到 FLV,需要转码为 H.264 - // 这里使用一个简化的方案:将 VP8 数据作为自定义格式封装 - // 实际生产环境需要使用 FFmpeg 或硬件编码器进行转码 - // FLV Video Tag Header: // FrameType (4 bits): 1=keyframe, 2=inter frame - // CodecID (4 bits): 7=AVC (H.264) - - // 由于 VP8 无法直接放入 FLV,这里使用一个变通方案 - // 将 VP8 数据标记为私有编码格式 + // CodecID (4 bits): 7=AVC (H.264), 12=VP8 (非标准) frameType := byte(2) // inter frame if isKeyframe { frameType = 1 // keyframe } - // 使用 CodecID=12 (VP8 - 非标准,仅用于内部传输) - // 或者可以考虑在服务端进行实时转码 + // 默认使用 VP8(非标准),如果是 H.264 则使用标准格式 + // 注意:VP8 不能被标准 FLV 播放器播放,需要转码 + // 这里保留 VP8 封装作为降级方案 codecID := byte(12) // 自定义:VP8 header := (frameType << 4) | codecID @@ -573,6 +738,43 @@ func (p *WebMToRTMPProxy) createFLVVideoTag(timestamp uint32, data []byte, isKey return p.createFLVTag(9, timestamp, videoData) // 9 = video } +// createFLVVideoTagH264 创建 H.264 FLV 视频 tag +// H.264 可以直接封装到标准 FLV +func (p *WebMToRTMPProxy) createFLVVideoTagH264(timestamp uint32, data []byte, isKeyframe bool) []byte { + // FLV Video Tag Header: + // FrameType (4 bits): 1=keyframe, 2=inter frame + // CodecID (4 bits): 7=AVC (H.264) + + frameType := byte(2) // inter frame + if isKeyframe { + frameType = 1 // keyframe + } + + codecID := byte(7) // AVC (H.264) + header := (frameType << 4) | codecID + + // AVC 数据需要额外的封装 + // AVCPacketType: 0=AVC sequence header, 1=AVC NALU + // CompositionTime: 3 bytes (通常为 0) + avcPacketType := byte(1) // AVC NALU + if isKeyframe { + // 关键帧可能需要先发送 sequence header + // 这里简化处理,假设数据已经是正确格式 + } + + // 构建 AVC 数据 + // 1 byte header + 1 byte AVCPacketType + 3 bytes CompositionTime + data + videoData := make([]byte, 5+len(data)) + videoData[0] = header + videoData[1] = avcPacketType + videoData[2] = 0 // CompositionTime + videoData[3] = 0 + videoData[4] = 0 + copy(videoData[5:], data) + + return p.createFLVTag(9, timestamp, videoData) // 9 = video +} + // createFLVAudioTag 创建 FLV 音频 tag // Opus -> FLV (需要转换为 AAC,这里简化处理) func (p *WebMToRTMPProxy) createFLVAudioTag(timestamp uint32, data []byte) []byte {