From e9a9c09d19b41709440c35d2740b26601d2c8bdd Mon Sep 17 00:00:00 2001 From: liqi Date: Tue, 16 Dec 2025 16:42:58 +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/rtmp_protocol.go | 56 ++++++++++++++++++++------- 1 file changed, 43 insertions(+), 13 deletions(-) diff --git a/internal/mediaserver/rtmp_protocol.go b/internal/mediaserver/rtmp_protocol.go index d8f628e..c9ffc60 100644 --- a/internal/mediaserver/rtmp_protocol.go +++ b/internal/mediaserver/rtmp_protocol.go @@ -82,6 +82,8 @@ type ChunkReader struct { prevHeaders map[uint32]*ChunkHeader // 缓存不完整的消息数据 messageBuffer map[uint32]*messageState + // 追踪已经收到过有效头部的 csid(用于 fmt 2/3 验证) + knownCSIDs map[uint32]bool } // messageState 追踪消息的读取状态 @@ -107,6 +109,7 @@ func NewChunkReader(conn net.Conn) *ChunkReader { chunkSize: DEFAULT_CHUNK_SIZE, prevHeaders: make(map[uint32]*ChunkHeader), messageBuffer: make(map[uint32]*messageState), + knownCSIDs: make(map[uint32]bool), } } @@ -160,12 +163,19 @@ func (r *ChunkReader) ReadMessage() (*RTMPMessage, error) { state, exists := r.messageBuffer[header.ChunkStreamID] if !exists || state.bytesRead >= state.header.MessageLength { // 新消息或上一条消息已完成,创建新状态 + // 注意:如果 header.MessageLength 为 0(fmt=3 无有效 prevHeader), + // 这里会创建一个容量为 0 的 slice,后续 toRead 检查会捕获这种情况 state = &messageState{ header: header, data: make([]byte, 0, header.MessageLength), bytesRead: 0, } r.messageBuffer[header.ChunkStreamID] = state + } else { + // 消息续传:更新时间戳(如果 header 有新的时间戳) + if header.Timestamp > 0 { + state.header.Timestamp = header.Timestamp + } } // 计算本次要读取的字节数 @@ -179,11 +189,14 @@ func (r *ChunkReader) ReadMessage() (*RTMPMessage, error) { toRead = remaining } - // toRead=0 说明 MessageLength=0,这是无效的 RTMP 消息 - // 直接返回错误而不是继续,避免失去同步 + // toRead=0 说明 MessageLength=0,可能是 fmt=3 没有有效的 prevHeader + // 记录详细日志以便调试,然后跳过此 chunk 继续尝试读取下一个 if toRead == 0 { - return nil, fmt.Errorf("invalid RTMP: toRead=0 on csid=%d msgLen=%d bytesRead=%d", - header.ChunkStreamID, state.header.MessageLength, state.bytesRead) + 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 重新开始 + delete(r.messageBuffer, header.ChunkStreamID) + continue // 继续尝试读取下一个 chunk,而不是返回错误 } // 读取 chunk 数据 @@ -267,6 +280,8 @@ func (r *ChunkReader) readChunkHeader() (*ChunkHeader, error) { header.MessageLength = uint32(data[3])<<16 | uint32(data[4])<<8 | uint32(data[5]) header.MessageTypeID = data[6] header.MessageSID = binary.LittleEndian.Uint32(data[7:11]) + // 标记此 csid 已收到有效头部 + r.knownCSIDs[csid] = true case CHUNK_FMT_1: // 7 bytes: timestamp delta(3) + length(3) + typeID(1) @@ -278,12 +293,10 @@ func (r *ChunkReader) readChunkHeader() (*ChunkHeader, error) { header.Timestamp = prevHeader.Timestamp + delta header.MessageLength = uint32(data[3])<<16 | uint32(data[4])<<8 | uint32(data[5]) header.MessageTypeID = data[6] + // 标记此 csid 已收到有效头部 + r.knownCSIDs[csid] = true case CHUNK_FMT_2: - // 验证:fmt 2 必须有有效的 prevHeader(MessageLength 从 prevHeader 继承) - if prevHeader.MessageLength == 0 { - return nil, fmt.Errorf("invalid RTMP: fmt 2 on csid %d without valid prevHeader (msgLen=0)", csid) - } // 3 bytes: timestamp delta(3) data := make([]byte, 3) if _, err := io.ReadFull(r.reader, data); err != nil { @@ -291,14 +304,31 @@ 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 检查来捕获真正的无效数据 + if prevHeader.MessageLength == 0 && !r.knownCSIDs[csid] { + log.Printf("⚠️ [RTMP Protocol] fmt=2 on csid %d without valid prevHeader, may cause issues", csid) + } case CHUNK_FMT_3: - // 验证:fmt 3 必须有有效的 prevHeader(所有字段从 prevHeader 继承) - if prevHeader.MessageLength == 0 { - return nil, fmt.Errorf("invalid RTMP: fmt 3 on csid %d without valid prevHeader (msgLen=0)", csid) + // 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 { + // 是消息续传,使用已保存的消息头信息 + header.MessageLength = state.header.MessageLength + header.MessageTypeID = state.header.MessageTypeID + header.MessageSID = state.header.MessageSID + header.Timestamp = state.header.Timestamp + log.Printf("ℹ️ [RTMP Protocol] fmt=3 续传 csid=%d msgLen=%d", csid, header.MessageLength) + } else { + log.Printf("⚠️ [RTMP Protocol] fmt=3 on csid %d without valid prevHeader (首次使用), msgLen=0", csid) + // 不返回错误,让 ReadMessage 中的 toRead==0 检查来处理 + } } - // 0 bytes: 使用上一个头部的所有字段 - // 已经复制了 prevHeader 的值 } // 检查是否有扩展时间戳