From bb7891706b88bce6b91529600d34bf4a679c9603 Mon Sep 17 00:00:00 2001 From: liqi Date: Tue, 16 Dec 2025 22:00:45 +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/api/call_handler.go | 3 +- internal/mediaserver/ffmpeg_rtmp_receiver.go | 381 +++++++++++++++++++ internal/mediaserver/ffmpeg_transcoder.go | 47 ++- internal/mediaserver/rtmp.go | 74 +++- internal/mediaserver/rtmp_protocol.go | 101 ++++- 5 files changed, 574 insertions(+), 32 deletions(-) create mode 100644 internal/mediaserver/ffmpeg_rtmp_receiver.go diff --git a/internal/api/call_handler.go b/internal/api/call_handler.go index 4c82bf1..f69b7e0 100644 --- a/internal/api/call_handler.go +++ b/internal/api/call_handler.go @@ -222,7 +222,8 @@ func JoinCallRoomHandler(c *gin.Context) { return } - stream, err := rtmpServer.GenerateStreamURLs(req.RoomID, req.UserID) + // 使用 FFmpeg -listen 模式接收小程序推流 + stream, err := rtmpServer.GenerateStreamURLsWithPlatform(req.RoomID, req.UserID, req.Platform) if err != nil { utils.InternalError(c, "生成推流地址失败") return diff --git a/internal/mediaserver/ffmpeg_rtmp_receiver.go b/internal/mediaserver/ffmpeg_rtmp_receiver.go new file mode 100644 index 0000000..a8cddea --- /dev/null +++ b/internal/mediaserver/ffmpeg_rtmp_receiver.go @@ -0,0 +1,381 @@ +/** + * package mediaserver + * + * FFmpeg RTMP 接收器 + * 使用 FFmpeg -listen 模式接收小程序 RTMP 推流,转换为 FLV 输出 + * 解决自定义 RTMP 协议解析器无法正确处理微信小程序 live-pusher 的问题 + */ +package mediaserver + +import ( + "bytes" + "fmt" + "io" + "log" + "os/exec" + "sync" + "time" +) + +// FFmpegRTMPReceiver FFmpeg RTMP 接收器 +type FFmpegRTMPReceiver struct { + streamID string + port int + ffmpegPath string + rtmpServer *RTMPServer + + cmd *exec.Cmd + stdout io.ReadCloser + stderr io.ReadCloser + + outputChan chan []byte + stopChan chan struct{} + + mu sync.Mutex + running bool + + // 统计信息 + startTime time.Time + outputBytes int64 +} + +// FFmpegRTMPReceiverPool 接收器池,管理多个 FFmpeg RTMP 接收器 +type FFmpegRTMPReceiverPool struct { + receivers map[string]*FFmpegRTMPReceiver + mu sync.RWMutex + rtmpServer *RTMPServer + ffmpegPath string + basePort int + nextPort int + maxPort int + usedPorts map[int]bool +} + +// NewFFmpegRTMPReceiverPool 创建接收器池 +func NewFFmpegRTMPReceiverPool(rtmpServer *RTMPServer, basePort int) (*FFmpegRTMPReceiverPool, error) { + ffmpegPath, err := EnsureFFmpeg() + if err != nil { + return nil, fmt.Errorf("FFmpeg 不可用: %w", err) + } + + return &FFmpegRTMPReceiverPool{ + receivers: make(map[string]*FFmpegRTMPReceiver), + rtmpServer: rtmpServer, + ffmpegPath: ffmpegPath, + basePort: basePort, + nextPort: basePort, + maxPort: basePort + 100, // 最多支持 100 个并发流 + usedPorts: make(map[int]bool), + }, nil +} + +// allocatePortLocked 分配一个可用端口(调用者必须持有锁) +func (p *FFmpegRTMPReceiverPool) allocatePortLocked() (int, error) { + // 注意:此方法假设调用者已持有 p.mu 锁 + // 查找可用端口 + for port := p.basePort; port < p.maxPort; port++ { + if !p.usedPorts[port] { + p.usedPorts[port] = true + return port, nil + } + } + + return 0, fmt.Errorf("没有可用端口") +} + +// releasePortLocked 释放端口(调用者必须持有锁) +func (p *FFmpegRTMPReceiverPool) releasePortLocked(port int) { + // 注意:此方法假设调用者已持有 p.mu 锁 + delete(p.usedPorts, port) +} + +// GetOrCreate 获取或创建接收器 +func (p *FFmpegRTMPReceiverPool) GetOrCreate(streamID string) (*FFmpegRTMPReceiver, error) { + p.mu.Lock() + defer p.mu.Unlock() + + // 检查是否已存在 + if r, exists := p.receivers[streamID]; exists && r.IsRunning() { + return r, nil + } + + // 分配端口(使用不加锁版本,因为我们已经持有锁) + port, err := p.allocatePortLocked() + if err != nil { + return nil, err + } + + // 创建新的接收器 + r := &FFmpegRTMPReceiver{ + streamID: streamID, + port: port, + ffmpegPath: p.ffmpegPath, + rtmpServer: p.rtmpServer, + outputChan: make(chan []byte, 100), + stopChan: make(chan struct{}), + } + + p.receivers[streamID] = r + return r, nil +} + +// Release 释放接收器 +func (p *FFmpegRTMPReceiverPool) Release(streamID string) { + p.mu.Lock() + defer p.mu.Unlock() + + if r, exists := p.receivers[streamID]; exists { + p.releasePortLocked(r.port) + r.Stop() + delete(p.receivers, streamID) + } +} + +// ReleaseAll 释放所有接收器 +func (p *FFmpegRTMPReceiverPool) ReleaseAll() { + p.mu.Lock() + defer p.mu.Unlock() + + for id, r := range p.receivers { + p.releasePortLocked(r.port) + r.Stop() + delete(p.receivers, id) + } +} + +// GetRTMPPort 获取接收器的 RTMP 端口 +func (r *FFmpegRTMPReceiver) GetRTMPPort() int { + return r.port +} + +// GetRTMPURL 获取完整的 RTMP 推流地址 +func (r *FFmpegRTMPReceiver) GetRTMPURL() string { + return fmt.Sprintf("rtmp://127.0.0.1:%d/live/%s", r.port, r.streamID) +} + +// Start 启动接收器 +func (r *FFmpegRTMPReceiver) Start() error { + r.mu.Lock() + defer r.mu.Unlock() + + if r.running { + return fmt.Errorf("接收器已在运行") + } + + // 构建 FFmpeg 命令 + // ffmpeg -listen 1 -i rtmp://0.0.0.0:{port}/live/{stream_id} -c copy -f flv pipe:1 + args := []string{ + "-listen", "1", // 作为服务器监听 + "-timeout", "30000000", // 超时时间 30 秒(微秒) + "-i", fmt.Sprintf("rtmp://0.0.0.0:%d/live/%s", r.port, r.streamID), // 监听地址 + "-c", "copy", // 直接复制,不重新编码 + "-fflags", "nobuffer", // 禁用输入缓冲 + "-flags", "low_delay", // 低延迟模式 + "-f", "flv", // 输出格式 + "pipe:1", // 输出到 stdout + } + + r.cmd = exec.Command(r.ffmpegPath, args...) + + var err error + + // 获取 stdout + r.stdout, err = r.cmd.StdoutPipe() + if err != nil { + return fmt.Errorf("获取 stdout 失败: %w", err) + } + + // 获取 stderr + r.stderr, err = r.cmd.StderrPipe() + if err != nil { + return fmt.Errorf("获取 stderr 失败: %w", err) + } + + // 启动 FFmpeg 进程 + if err := r.cmd.Start(); err != nil { + return fmt.Errorf("启动 FFmpeg 失败: %w", err) + } + + r.running = true + r.startTime = time.Now() + + log.Printf("🎬 [FFmpegRTMPReceiver] 启动成功 | Stream:%s Port:%d PID:%d", + r.streamID, r.port, r.cmd.Process.Pid) + + // 启动输出读取协程 + go r.readOutput() + + // 启动 stderr 读取协程 + go r.readStderr() + + // 启动进程监控协程 + go r.monitor() + + return nil +} + +// Stop 停止接收器 +func (r *FFmpegRTMPReceiver) Stop() error { + r.mu.Lock() + defer r.mu.Unlock() + + if !r.running { + return nil + } + + r.running = false + close(r.stopChan) + + // 关闭进程 + if r.cmd != nil && r.cmd.Process != nil { + r.cmd.Process.Kill() + r.cmd.Wait() + } + + // 关闭输出通道 + close(r.outputChan) + + duration := time.Since(r.startTime) + log.Printf("🛑 [FFmpegRTMPReceiver] 已停止 | Stream:%s Port:%d Duration:%v Output:%d bytes", + r.streamID, r.port, duration, r.outputBytes) + + return nil +} + +// IsRunning 检查是否运行中 +func (r *FFmpegRTMPReceiver) IsRunning() bool { + r.mu.Lock() + defer r.mu.Unlock() + return r.running +} + +// Output 获取输出通道 +func (r *FFmpegRTMPReceiver) Output() <-chan []byte { + return r.outputChan +} + +// readOutput 读取 FFmpeg FLV 输出 +func (r *FFmpegRTMPReceiver) readOutput() { + log.Printf("▶️ [FFmpegRTMPReceiver] 开始读取 FLV 输出 | Stream:%s", r.streamID) + + buf := make([]byte, 64*1024) // 64KB 缓冲区 + flvHeaderSent := false + + for { + select { + case <-r.stopChan: + return + default: + } + + n, err := r.stdout.Read(buf) + if err != nil { + if err != io.EOF { + log.Printf("⚠️ [FFmpegRTMPReceiver] 读取输出错误: %v | Stream:%s", err, r.streamID) + } + return + } + + if n > 0 { + r.outputBytes += int64(n) + + // 复制数据 + data := make([]byte, n) + copy(data, buf[:n]) + + // 跳过 FLV 头(前 13 字节)- 只跳过一次 + if !flvHeaderSent && len(data) >= 13 { + // 验证 FLV 头 + if data[0] == 'F' && data[1] == 'L' && data[2] == 'V' { + log.Printf("✅ [FFmpegRTMPReceiver] 收到 FLV 头 | Stream:%s", r.streamID) + flvHeaderSent = true + // 广播 FLV 数据(包括头) + r.broadcastFLVData(data) + } else { + // 不是 FLV 头,直接广播 + r.broadcastFLVData(data) + } + } else { + // 广播 FLV 数据 + r.broadcastFLVData(data) + } + + // 发送到输出通道 + select { + case r.outputChan <- data: + default: + // 通道满,丢弃 + log.Printf("⚠️ [FFmpegRTMPReceiver] 输出通道满,丢弃 %d bytes | Stream:%s", n, r.streamID) + } + } + } +} + +// broadcastFLVData 广播 FLV 数据给所有订阅者 +func (r *FFmpegRTMPReceiver) broadcastFLVData(data []byte) { + if r.rtmpServer == nil { + return + } + + // 解析 FLV tag 并广播 + r.parseFLVAndBroadcast(data) +} + +// parseFLVAndBroadcast 解析 FLV 数据并广播 +func (r *FFmpegRTMPReceiver) parseFLVAndBroadcast(data []byte) { + // 获取流 + stream := r.rtmpServer.GetStream(r.streamID) + if stream == nil { + return + } + + // 简单广播原始数据给订阅者 + r.rtmpServer.subMu.RLock() + subs := r.rtmpServer.subscribers[r.streamID] + r.rtmpServer.subMu.RUnlock() + + if subs == nil { + return + } + + for _, sub := range subs { + select { + case sub.DataChan <- data: + default: + // 订阅者通道满,跳过 + } + } +} + +// readStderr 读取 FFmpeg 错误输出 +func (r *FFmpegRTMPReceiver) readStderr() { + buf := new(bytes.Buffer) + io.Copy(buf, r.stderr) + + if buf.Len() > 0 { + output := buf.String() + if len(output) > 500 { + output = output[:500] + "..." + } + log.Printf("📋 [FFmpegRTMPReceiver] FFmpeg 输出: %s | Stream:%s", output, r.streamID) + } +} + +// monitor 监控 FFmpeg 进程 +func (r *FFmpegRTMPReceiver) monitor() { + err := r.cmd.Wait() + + r.mu.Lock() + wasRunning := r.running + r.running = false + r.mu.Unlock() + + if wasRunning { + if err != nil { + log.Printf("⚠️ [FFmpegRTMPReceiver] FFmpeg 进程异常退出: %v | Stream:%s", err, r.streamID) + } else { + log.Printf("ℹ️ [FFmpegRTMPReceiver] FFmpeg 进程正常退出 | Stream:%s", r.streamID) + } + } +} + diff --git a/internal/mediaserver/ffmpeg_transcoder.go b/internal/mediaserver/ffmpeg_transcoder.go index 3ff29ea..6fdbdb6 100644 --- a/internal/mediaserver/ffmpeg_transcoder.go +++ b/internal/mediaserver/ffmpeg_transcoder.go @@ -103,37 +103,56 @@ func (t *FFmpegTranscoder) Start() error { if t.config.CopyMode { // H.264 copy 模式:只重封装,不重新编码(低延迟) - // ffmpeg -f webm -i pipe:0 -c:v copy -c:a aac -f flv pipe:1 + // 添加低延迟参数:禁用缓冲、减少探测时间 args = []string{ - "-f", "webm", // 输入格式 - "-i", "pipe:0", // 从 stdin 读取 - "-c:v", "copy", // 视频直接复制,不重新编码 + // 低延迟输入参数 + "-fflags", "nobuffer", // 禁用输入缓冲 + "-flags", "low_delay", // 低延迟模式 + "-probesize", "32", // 减少探测大小(字节) + "-analyzeduration", "0", // 禁用分析时长 + "-f", "webm", // 输入格式 + "-i", "pipe:0", // 从 stdin 读取 + // 输出参数 + "-c:v", "copy", // 视频直接复制,不重新编码 "-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 + // 低延迟输出参数 + "-fflags", "+genpts", // 生成时间戳 + "-f", "flv", // 输出格式 + "pipe:1", // 输出到 stdout } - log.Printf("🎬 [FFmpegTranscoder] 使用 copy 模式(H.264 重封装)") + log.Printf("🎬 [FFmpegTranscoder] 使用 copy 模式(H.264 重封装,低延迟)") } else { // VP8/VP9 转码模式:重新编码为 H.264 - // 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 读取 + // 低延迟输入参数 + "-fflags", "nobuffer", // 禁用输入缓冲 + "-flags", "low_delay", // 低延迟模式 + "-probesize", "32", // 减少探测大小(字节) + "-analyzeduration", "0", // 禁用分析时长 + "-f", "webm", // 输入格式 + "-i", "pipe:0", // 从 stdin 读取 + // 视频编码参数 "-c:v", t.config.VideoCodec, "-preset", t.config.VideoPreset, - "-tune", "zerolatency", // 零延迟模式 + "-tune", "zerolatency", // 零延迟调优 "-b:v", t.config.VideoBitrate, + "-g", "30", // GOP 大小(减少关键帧间隔) + "-keyint_min", "15", // 最小关键帧间隔 + // 音频编码参数 "-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 + // 低延迟输出参数 + "-fflags", "+genpts", // 生成时间戳 + "-f", "flv", // 输出格式 + "pipe:1", // 输出到 stdout } - log.Printf("🎬 [FFmpegTranscoder] 使用转码模式(VP8/VP9 → H.264)") + log.Printf("🎬 [FFmpegTranscoder] 使用转码模式(VP8/VP9 → H.264,低延迟)") } t.cmd = exec.Command(t.ffmpegPath, args...) diff --git a/internal/mediaserver/rtmp.go b/internal/mediaserver/rtmp.go index 3d780fd..9e743ec 100644 --- a/internal/mediaserver/rtmp.go +++ b/internal/mediaserver/rtmp.go @@ -42,6 +42,9 @@ type RTMPServer struct { // 订阅者管理 subscribers map[string]map[string]*Subscriber // streamID -> subscriberID -> subscriber subMu sync.RWMutex + + // FFmpeg RTMP 接收器池(用于接收小程序推流) + ffmpegReceiverPool *FFmpegRTMPReceiverPool } // RTMPStream RTMP 流信息 @@ -94,7 +97,17 @@ func (r *RTMPServer) Start() { r.running = true r.mu.Unlock() - // 启动 RTMP 服务器 + // 初始化 FFmpeg RTMP 接收器池(用于接收小程序推流) + // 使用 19350 作为基础端口,支持 100 个并发流 + pool, err := NewFFmpegRTMPReceiverPool(r, 19350) + if err != nil { + log.Printf("⚠️ [RTMP] 初始化 FFmpeg 接收器池失败: %v,小程序推流可能无法正常工作", err) + } else { + r.ffmpegReceiverPool = pool + log.Printf("✅ [RTMP] FFmpeg 接收器池已初始化 | 端口范围: 19350-19450") + } + + // 启动 RTMP 服务器(作为备用,处理不使用 FFmpeg 的推流) go r.startRTMPServer() // 启动 HTTP-FLV 服务器 @@ -264,6 +277,15 @@ func (h *ConnectionHandler) handleSetChunkSize(msg *RTMPMessage) error { } newSize := binary.BigEndian.Uint32(msg.Data) + + // 关键验证:chunkSize 必须在合理范围内 (1 到 16MB) + // 参考 RTMP 规范:最大值通常不超过 16MB,常见值为 128-65536 + // 如果值异常,可能是协议解析错误,忽略此消息 + if newSize == 0 || newSize > 16*1024*1024 { + log.Printf("⚠️ [RTMP] SetChunkSize 值异常,忽略: %d (remote: %s)", newSize, h.remoteAddr) + return nil + } + log.Printf("📝 [RTMP] SetChunkSize: %d -> %d (remote: %s)", h.reader.GetChunkSize(), newSize, h.remoteAddr) // 关键:立即更新 chunk size @@ -1002,6 +1024,12 @@ func (r *RTMPServer) Stop() { // GenerateStreamURLs 为用户生成推拉流地址 func (r *RTMPServer) GenerateStreamURLs(roomID, userID string) (*RTMPStream, error) { + return r.GenerateStreamURLsWithPlatform(roomID, userID, "web") +} + +// GenerateStreamURLsWithPlatform 为用户生成推拉流地址(指定平台) +// platform: "web", "h5", "app", "miniprogram", "wxapp" +func (r *RTMPServer) GenerateStreamURLsWithPlatform(roomID, userID, platform string) (*RTMPStream, error) { r.mu.RLock() if !r.running { r.mu.RUnlock() @@ -1018,12 +1046,45 @@ func (r *RTMPServer) GenerateStreamURLs(roomID, userID string) (*RTMPStream, err token := generateStreamToken(streamID) + var pushURL, pullURL string + + // 判断是否是小程序 + isMiniprogram := platform == "miniprogram" || platform == "wxapp" + + if isMiniprogram && r.ffmpegReceiverPool != nil { + // 小程序使用 FFmpeg -listen 模式接收推流 + receiver, err := r.ffmpegReceiverPool.GetOrCreate(streamID) + if err != nil { + log.Printf("⚠️ [RTMP] 创建 FFmpeg 接收器失败: %v,使用默认 RTMP 端口", err) + pushURL = fmt.Sprintf("rtmp://%s:%d/live/%s?token=%s", publicIP, r.config.RTMPPort, streamID, token) + pullURL = fmt.Sprintf("rtmp://%s:%d/live/%s", publicIP, r.config.RTMPPort, streamID) + } else { + // 启动 FFmpeg 接收器 + if err := receiver.Start(); err != nil { + log.Printf("⚠️ [RTMP] 启动 FFmpeg 接收器失败: %v,使用默认 RTMP 端口", err) + r.ffmpegReceiverPool.Release(streamID) + pushURL = fmt.Sprintf("rtmp://%s:%d/live/%s?token=%s", publicIP, r.config.RTMPPort, streamID, token) + pullURL = fmt.Sprintf("rtmp://%s:%d/live/%s", publicIP, r.config.RTMPPort, streamID) + } else { + // 使用 FFmpeg 接收器的端口 + ffmpegPort := receiver.GetRTMPPort() + pushURL = fmt.Sprintf("rtmp://%s:%d/live/%s?token=%s", publicIP, ffmpegPort, streamID, token) + pullURL = fmt.Sprintf("rtmp://%s:%d/live/%s", publicIP, ffmpegPort, streamID) + log.Printf("✅ [RTMP] 小程序流使用 FFmpeg 接收器 | Stream:%s Port:%d", streamID, ffmpegPort) + } + } + } else { + // Web/H5/App 使用默认 RTMP 端口 + pushURL = fmt.Sprintf("rtmp://%s:%d/live/%s?token=%s", publicIP, r.config.RTMPPort, streamID, token) + pullURL = fmt.Sprintf("rtmp://%s:%d/live/%s", publicIP, r.config.RTMPPort, streamID) + } + stream := &RTMPStream{ ID: streamID, RoomID: roomID, UserID: userID, - PushURL: fmt.Sprintf("rtmp://%s:%d/live/%s?token=%s", publicIP, r.config.RTMPPort, streamID, token), - PullURL: fmt.Sprintf("rtmp://%s:%d/live/%s", publicIP, r.config.RTMPPort, streamID), + PushURL: pushURL, + PullURL: pullURL, // 使用 HTTPS 域名,通过 Nginx 反向代理到 HTTP-FLV 服务 // 避免浏览器混合内容安全策略阻止 HTTP 请求 FLVURL: fmt.Sprintf("https://g-ws.nailaoyun.cn/live/%s.flv", streamID), @@ -1037,7 +1098,7 @@ func (r *RTMPServer) GenerateStreamURLs(roomID, userID string) (*RTMPStream, err r.streams[streamID] = stream r.streamsMu.Unlock() - log.Printf("🎬 [RTMP] 创建流 | Room:%s User:%s Stream:%s", roomID, userID, streamID) + log.Printf("🎬 [RTMP] 创建流 | Room:%s User:%s Stream:%s Platform:%s", roomID, userID, streamID, platform) return stream, nil } @@ -1095,6 +1156,11 @@ func (r *RTMPServer) RemoveStream(streamID string) { log.Printf("🎬 [RTMP] 移除流 | Stream:%s", streamID) } + // 清理 FFmpeg 接收器 + if r.ffmpegReceiverPool != nil { + r.ffmpegReceiverPool.Release(streamID) + } + r.subMu.Lock() delete(r.subscribers, streamID) r.subMu.Unlock() diff --git a/internal/mediaserver/rtmp_protocol.go b/internal/mediaserver/rtmp_protocol.go index c9ffc60..1f02004 100644 --- a/internal/mediaserver/rtmp_protocol.go +++ b/internal/mediaserver/rtmp_protocol.go @@ -189,14 +189,53 @@ func (r *ChunkReader) ReadMessage() (*RTMPMessage, error) { toRead = remaining } - // toRead=0 说明 MessageLength=0,可能是 fmt=3 没有有效的 prevHeader - // 记录详细日志以便调试,然后跳过此 chunk 继续尝试读取下一个 + // toRead=0 说明 MessageLength=0,可能是 fmt=2/3 没有有效的 prevHeader + // 尝试通过搜索下一个 fmt=0 chunk 来重新同步 if toRead == 0 { - log.Printf("⚠️ [RTMP Protocol] toRead=0 on csid=%d fmt=%d msgLen=%d bytesRead=%d, 跳过此 chunk", - header.ChunkStreamID, header.Format, state.header.MessageLength, state.bytesRead) - // 清除此 csid 的无效状态,等待有效的 fmt=0 重新开始 + log.Printf("⚠️ [RTMP Protocol] csid=%d fmt=%d msgLen=0, 尝试重新同步...", + header.ChunkStreamID, header.Format) + + // 清除此 csid 的状态 delete(r.messageBuffer, header.ChunkStreamID) - continue // 继续尝试读取下一个 chunk,而不是返回错误 + + // 尝试寻找下一个有效的 chunk 起始位置 + // 读取并丢弃字节,直到找到可能的 fmt=0 chunk header + maxScanBytes := 4096 // 最多扫描 4KB + scanned := 0 + foundSync := false + + for scanned < maxScanBytes { + b, err := r.reader.ReadByte() + if err != nil { + return nil, fmt.Errorf("重新同步时读取失败: %w", err) + } + scanned++ + + // 检查是否可能是 fmt=0 的 chunk header + // fmt=0 的 basic header 第一个字节: 高 2 位为 00 + potentialFmt := (b >> 6) & 0x03 + potentialCsid := uint32(b & 0x3F) + + if potentialFmt == 0 && potentialCsid >= 2 && potentialCsid < 64 { + // 可能找到了 fmt=0 chunk,尝试验证 + // 先放回这个字节 + if err := r.reader.UnreadByte(); err != nil { + return nil, fmt.Errorf("UnreadByte 失败: %w", err) + } + + log.Printf("ℹ️ [RTMP Protocol] 可能找到同步点: 跳过了 %d 字节, 潜在 csid=%d", + scanned-1, potentialCsid) + foundSync = true + break + } + } + + if !foundSync { + log.Printf("❌ [RTMP Protocol] 扫描 %d 字节后仍未找到同步点", scanned) + return nil, fmt.Errorf("RTMP 协议错误: 无法重新同步") + } + + continue // 重新尝试读取 chunk header } // 读取 chunk 数据 @@ -216,6 +255,13 @@ func (r *ChunkReader) ReadMessage() (*RTMPMessage, error) { StreamID: state.header.MessageSID, Data: state.data, } + + // 调试日志:记录消息完成信息 + if msg.TypeID == RTMP_MSG_AMF0_DATA || msg.TypeID == RTMP_MSG_AUDIO || msg.TypeID == RTMP_MSG_VIDEO { + log.Printf("✅ [RTMP Debug] 消息完成: csid=%d type=%d len=%d", + header.ChunkStreamID, msg.TypeID, len(msg.Data)) + } + // 清除消息缓冲 delete(r.messageBuffer, header.ChunkStreamID) return msg, nil @@ -233,6 +279,12 @@ func (r *ChunkReader) readChunkHeader() (*ChunkHeader, error) { format := (firstByte >> 6) & 0x03 csid := uint32(firstByte & 0x3F) + + // 调试日志:记录每个 chunk header 的原始字节 + // 注意:csid > 64 或异常 csid 值可能表示流同步丢失 + if csid > 20 || (format != 0 && format != 3) { + log.Printf("🔍 [RTMP Debug] firstByte=0x%02X fmt=%d csid=%d", firstByte, format, csid) + } // 扩展 chunk stream ID if csid == 0 { @@ -282,6 +334,10 @@ func (r *ChunkReader) readChunkHeader() (*ChunkHeader, error) { header.MessageSID = binary.LittleEndian.Uint32(data[7:11]) // 标记此 csid 已收到有效头部 r.knownCSIDs[csid] = true + + // 调试日志:记录 fmt=0 的完整信息 + log.Printf("🔍 [RTMP Debug] fmt=0 csid=%d ts=%d msgLen=%d typeID=%d streamID=%d", + csid, header.Timestamp, header.MessageLength, header.MessageTypeID, header.MessageSID) case CHUNK_FMT_1: // 7 bytes: timestamp delta(3) + length(3) + typeID(1) @@ -304,17 +360,30 @@ func (r *ChunkReader) readChunkHeader() (*ChunkHeader, error) { } delta := uint32(data[0])<<16 | uint32(data[1])<<8 | uint32(data[2]) header.Timestamp = prevHeader.Timestamp + delta - // 放宽验证:如果没有有效的 prevHeader,记录警告但不立即失败 - // 依赖后续 ReadMessage 中的 toRead==0 检查来捕获真正的无效数据 + + // 关键修复:如果没有有效的 prevHeader,尝试根据 csid 推断消息类型 + // 微信小程序 live-pusher 可能首次发送音视频时就使用 fmt=2 if prevHeader.MessageLength == 0 && !r.knownCSIDs[csid] { - log.Printf("⚠️ [RTMP Protocol] fmt=2 on csid %d without valid prevHeader, may cause issues", csid) + // 尝试推断消息类型(基于常见的 RTMP 实现) + // csid >= 4 通常用于音视频数据 + if csid >= 4 { + // 设置默认消息长度为 chunkSize(会在后续 chunk 中累积) + // 这是一个启发式处理,让数据能够被读取 + r.mu.RLock() + header.MessageLength = r.chunkSize + r.mu.RUnlock() + // 猜测消息类型:偶数 csid 可能是音频,奇数是视频(这是一个启发式) + // 实际上我们需要从数据本身来判断 + log.Printf("⚠️ [RTMP Protocol] fmt=2 csid=%d 无 prevHeader,设置临时 msgLen=%d", + csid, header.MessageLength) + } else { + log.Printf("⚠️ [RTMP Protocol] fmt=2 on csid %d without valid prevHeader", csid) + } } case CHUNK_FMT_3: // 0 bytes: 使用上一个头部的所有字段(已经复制了 prevHeader 的值) - // 放宽验证:如果没有有效的 prevHeader,记录警告但不立即失败 - // 某些 RTMP 客户端(如小程序 live-pusher)可能在首次使用某 csid 时就用 fmt=3 - // 检查是否在 messageBuffer 中有未完成的消息(消息续传场景) + // 关键修复:检查消息续传或尝试推断 if prevHeader.MessageLength == 0 && !r.knownCSIDs[csid] { // 检查是否是消息续传 if state, exists := r.messageBuffer[csid]; exists && state.bytesRead < state.header.MessageLength { @@ -324,9 +393,15 @@ func (r *ChunkReader) readChunkHeader() (*ChunkHeader, error) { header.MessageSID = state.header.MessageSID header.Timestamp = state.header.Timestamp log.Printf("ℹ️ [RTMP Protocol] fmt=3 续传 csid=%d msgLen=%d", csid, header.MessageLength) + } else if csid >= 4 { + // 对于可能的音视频 csid,设置默认消息长度 + r.mu.RLock() + header.MessageLength = r.chunkSize + r.mu.RUnlock() + log.Printf("⚠️ [RTMP Protocol] fmt=3 csid=%d 无 prevHeader,设置临时 msgLen=%d", + csid, header.MessageLength) } else { log.Printf("⚠️ [RTMP Protocol] fmt=3 on csid %d without valid prevHeader (首次使用), msgLen=0", csid) - // 不返回错误,让 ReadMessage 中的 toRead==0 检查来处理 } } }