网站通话使用webrtc,新增ffmpeg

This commit is contained in:
2025-12-16 22:00:45 +08:00
parent e9a9c09d19
commit bb7891706b
5 changed files with 574 additions and 32 deletions

View File

@@ -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

View File

@@ -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)
}
}
}

View File

@@ -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...)

View File

@@ -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()

View File

@@ -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 检查来处理
}
}
}