Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
97 changes: 82 additions & 15 deletions config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@ package config

import (
"encoding/json"
"fmt"
"io/ioutil"
"os"
)

var defaultConfig Config
Expand All @@ -18,6 +20,7 @@ type Config struct {
LalSvrConfigPath string `json:"lal_config_path"` // lal配置文件路径,兼容旧版配置
LogicConfig LogicConfig `json:"logic_config"` // 扩展流组配置
LalRawContent []byte `json:"-"` // lal 原始配置内容
ConfFilePath string `json:"-"` // 配置文件路径,用于持久化
}

type SrtConfig struct {
Expand Down Expand Up @@ -85,22 +88,54 @@ type GB28181MediaConfig struct {
MultiPortMaxIncrement uint16 `json:"multi_port_max_increment"` // 多端口范围 ListenPort+1至ListenPort+MultiPortMax
}

// ZlmCompatHookConfig ZLM 兼容 hook URL 配置
// 为什么独立结构体:隔离 ZLM 适配层,lalmax 原有字段保持不变
type ZlmCompatHookConfig struct {
ZlmOnStreamChanged string `json:"zlm_on_stream_changed"`
ZlmOnServerKeepalive string `json:"zlm_on_server_keepalive"`
ZlmOnStreamNoneReader string `json:"zlm_on_stream_none_reader"`
ZlmOnRtpServerTimeout string `json:"zlm_on_rtp_server_timeout"`
ZlmOnRecordMp4 string `json:"zlm_on_record_mp4"`
ZlmOnPublish string `json:"zlm_on_publish"`
ZlmOnPlay string `json:"zlm_on_play"`
ZlmOnStreamNotFound string `json:"zlm_on_stream_not_found"`
ZlmOnServerStarted string `json:"zlm_on_server_started"`
}

// HasZlmHooks 任一 ZLM 兼容 hook 字段有值即返回 true
// 为什么:ZLM 回调与 lalmax 原有回调二选一,此方法为判断条件
func (c ZlmCompatHookConfig) HasZlmHooks() bool {
return c.ZlmOnStreamChanged != "" ||
c.ZlmOnServerKeepalive != "" ||
c.ZlmOnStreamNoneReader != "" ||
c.ZlmOnRtpServerTimeout != "" ||
c.ZlmOnRecordMp4 != "" ||
c.ZlmOnPublish != "" ||
c.ZlmOnPlay != "" ||
c.ZlmOnStreamNotFound != ""
}

type HttpNotifyConfig struct {
Enable bool `json:"enable"`
UpdateIntervalSec int `json:"update_interval_sec"`
OnServerStart string `json:"on_server_start"`
OnUpdate string `json:"on_update"`
OnGroupStart string `json:"on_group_start"`
OnGroupStop string `json:"on_group_stop"`
OnStreamActive string `json:"on_stream_active"`
OnPubStart string `json:"on_pub_start"`
OnPubStop string `json:"on_pub_stop"`
OnSubStart string `json:"on_sub_start"`
OnSubStop string `json:"on_sub_stop"`
OnRelayPullStart string `json:"on_relay_pull_start"`
OnRelayPullStop string `json:"on_relay_pull_stop"`
OnRtmpConnect string `json:"on_rtmp_connect"`
OnHlsMakeTs string `json:"on_hls_make_ts"`
Enable bool `json:"enable"`
UpdateIntervalSec int `json:"update_interval_sec"`
KeepaliveIntervalSec int `json:"keepalive_interval_sec"`
HookTimeoutSec int `json:"hook_timeout_sec"`
OnServerStart string `json:"on_server_start"`
OnUpdate string `json:"on_update"`
OnGroupStart string `json:"on_group_start"`
OnGroupStop string `json:"on_group_stop"`
OnStreamActive string `json:"on_stream_active"`
OnPubStart string `json:"on_pub_start"`
OnPubStop string `json:"on_pub_stop"`
OnSubStart string `json:"on_sub_start"`
OnSubStop string `json:"on_sub_stop"`
OnRelayPullStart string `json:"on_relay_pull_start"`
OnRelayPullStop string `json:"on_relay_pull_stop"`
OnRtmpConnect string `json:"on_rtmp_connect"`
OnHlsMakeTs string `json:"on_hls_make_ts"`

// --- ZLM 兼容 hook 配置 ---
ZlmCompatHookConfig
}

type LogicConfig struct {
Expand Down Expand Up @@ -189,3 +224,35 @@ func unmarshalConfig(data []byte, cfg *Config) error {
func GetConfig() *Config {
return &defaultConfig
}

// SaveToFile 将当前配置持久化到配置文件
// 为什么:setServerConfig 动态修改后需落盘,重启后配置仍生效
func (c *Config) SaveToFile() error {
if c.ConfFilePath == "" {
return nil
}

data, err := os.ReadFile(c.ConfFilePath)
if err != nil {
return fmt.Errorf("read config file: %w", err)
}

var file map[string]json.RawMessage
if err := json.Unmarshal(data, &file); err != nil {
return fmt.Errorf("parse config file: %w", err)
}

lalmax, err := json.MarshalIndent(c, "", " ")
if err != nil {
return fmt.Errorf("marshal lalmax config: %w", err)
}
file["lalmax"] = lalmax

out, err := json.MarshalIndent(file, "", " ")
if err != nil {
return fmt.Errorf("marshal config file: %w", err)
}
out = append(out, '\n')

return os.WriteFile(c.ConfFilePath, out, 0o644)
}
10 changes: 10 additions & 0 deletions gb28181/rtppub/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,16 @@ func (m *Manager) CheckSsrc(ssrc uint32) (*mediaserver.MediaInfo, bool) {
func (m *Manager) NotifyClose(streamName string) {
}

// UpdatePortRange 动态更新端口范围,由 setServerConfig 接口调用
// 为什么:owl 通过 setServerConfig 下发 rtp_proxy.port_range,需运行时生效
func (m *Manager) UpdatePortRange(portMin, portMax int) {
m.mu.Lock()
defer m.mu.Unlock()
m.portMin = portMin
m.portMax = portMax
nazalog.Infof("rtp pub port range updated. min=%d, max=%d", portMin, portMax)
}

func (m *Manager) OnRtpPacket(streamName string, mediaKey string) {
m.mu.Lock()
defer m.mu.Unlock()
Expand Down
16 changes: 16 additions & 0 deletions logic/group_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package logic

import (
"sync"
"time"

"github.com/q191201771/lalmax/fmp4/hls"
"github.com/q191201771/naza/pkg/nazalog"
Expand Down Expand Up @@ -211,6 +212,21 @@ func (m *ComplexGroupManager) GetGroupByStreamName(streamName string) (bool, *Gr
return m.GetGroup(StreamKeyFromStreamName(streamName))
}

// WaitGroup 等待流就绪,轮询 interval 间隔,总超时 timeout
// 为什么:GB28181 设备推流有延迟,播放端先于推流端到达,需短暂等待
func (m *ComplexGroupManager) WaitGroup(key StreamKey, interval, timeout time.Duration) (bool, *Group) {
deadline := time.Now().Add(timeout)
for {
if ok, g := m.GetGroup(key); ok {
return true, g
}
if time.Now().After(deadline) {
return false, nil
}
time.Sleep(interval)
}
}

// streamName 单独查找只在匹配唯一 appName 时成功,避免跨 app 串流。
func (m *ComplexGroupManager) getGroupByOnlyStreamNameLocked(streamName string) (bool, *Group) {
var found *Group
Expand Down
1 change: 1 addition & 0 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ func main() {
}

maxConf := config.GetConfig()
maxConf.ConfFilePath = confFilename

svr, err := server.NewLalMaxServer(maxConf)
if err != nil {
Expand Down
1 change: 0 additions & 1 deletion rtc/jessibucasession.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"math"
"sync"
"sync/atomic"

"github.com/gofrs/uuid"
"github.com/pion/webrtc/v3"
"github.com/q191201771/lal/pkg/base"
Expand Down
76 changes: 72 additions & 4 deletions rtc/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,10 @@ import (
"fmt"
"net"
"net/http"
"time"

config "github.com/q191201771/lalmax/config"
maxlogic "github.com/q191201771/lalmax/logic"

"github.com/gin-gonic/gin"
"github.com/pion/ice/v2"
Expand All @@ -14,11 +16,37 @@ import (
"github.com/q191201771/naza/pkg/nazalog"
)

// StreamNotFoundFn 流不存在时的回调,触发 on_stream_not_found 通知上层拉流
type StreamNotFoundFn func(app, stream, schema string)

type RtcServer struct {
config config.RtcConfig
lalServer logic.ILalServer
udpMux ice.UDPMux
tcpMux ice.TCPMux
config config.RtcConfig
lalServer logic.ILalServer
udpMux ice.UDPMux
tcpMux ice.TCPMux
streamNotFoundFn StreamNotFoundFn
}

// SetStreamNotFoundFn 注入流不存在回调
func (s *RtcServer) SetStreamNotFoundFn(fn StreamNotFoundFn) {
s.streamNotFoundFn = fn
}

// waitStreamReady 触发 on_stream_not_found 后轮询等待流就绪
// 为什么:WebRTC 播放请求先于 GB28181 设备推流到达,需通知上层拉流后等待
func (s *RtcServer) waitStreamReady(appName, streamid, schema string) bool {
key := maxlogic.NewStreamKey(appName, streamid)
if ok, _ := maxlogic.GetGroupManagerInstance().GetGroup(key); ok {
return true
}

if s.streamNotFoundFn != nil {
nazalog.Infof("stream not found, triggering on_stream_not_found. app=%s, stream=%s", appName, streamid)
s.streamNotFoundFn(appName, streamid, schema)
}

ok, _ := maxlogic.GetGroupManagerInstance().WaitGroup(key, 500*time.Millisecond, 5*time.Second)
return ok
}

func NewRtcServer(config config.RtcConfig, lal logic.ILalServer) (*RtcServer, error) {
Expand Down Expand Up @@ -136,6 +164,12 @@ func (s *RtcServer) HandleJessibuca(c *gin.Context) {
return
}

if !s.waitStreamReady(appName, streamid, "rtsp") {
nazalog.Errorf("stream not ready after waiting. app=%s, stream=%s", appName, streamid)
c.Status(http.StatusNotFound)
return
}

pc, err := newPeerConnection(s.config.ICEHostNATToIPs, s.udpMux, s.tcpMux)
if err != nil {
c.Status(http.StatusInternalServerError)
Expand Down Expand Up @@ -183,6 +217,12 @@ func (s *RtcServer) HandleWHEP(c *gin.Context) {
return
}

if !s.waitStreamReady(appName, streamid, "rtsp") {
nazalog.Errorf("stream not ready after waiting. app=%s, stream=%s", appName, streamid)
c.Status(http.StatusNotFound)
return
}

pc, err := newPeerConnection(s.config.ICEHostNATToIPs, s.udpMux, s.tcpMux)
if err != nil {
c.Status(http.StatusInternalServerError)
Expand All @@ -209,3 +249,31 @@ func (s *RtcServer) HandleWHEP(c *gin.Context) {

c.Data(http.StatusCreated, "application/sdp", []byte(sdp))
}

// HandleZlmWebrtcPlay ZLM 兼容 WebRTC 播放,返回 SDP answer
// 为什么独立方法:ZLM 信令格式为 JSON {"code":0,"sdp":"..."},与 WHEP 纯 SDP 不同
func (s *RtcServer) HandleZlmWebrtcPlay(app, stream, offer string) (string, error) {
if !s.waitStreamReady(app, stream, "rtsp") {
return "", fmt.Errorf("stream not found: %s/%s", app, stream)
}

pc, err := newPeerConnection(s.config.ICEHostNATToIPs, s.udpMux, s.tcpMux)
if err != nil {
return "", fmt.Errorf("create peer connection: %w", err)
}

session := NewWhepSession(app, stream, s.config.WriteChanSize, pc, s.lalServer)
if session == nil {
pc.Close()
return "", fmt.Errorf("create session failed: %s/%s", app, stream)
}

sdp := session.GetAnswerSDP(offer)
if sdp == "" {
session.Close()
return "", fmt.Errorf("generate answer sdp failed")
}

go session.Run()
return sdp, nil
}
35 changes: 35 additions & 0 deletions server/hook_builtin_http_plugin.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,9 @@ func (p *hookBuiltinHTTPPlugin) OnHookEvent(event HookEvent) error {
if p.hub.cfg.OnServerStart != "" {
p.hub.asyncPostEvent(p.hub.cfg.OnServerStart, event)
}
if p.hub.cfg.ZlmOnServerStarted != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnServerStarted, event)
}
case HookEventUpdate:
if p.hub.cfg.OnUpdate != "" {
p.hub.asyncPostEvent(p.hub.cfg.OnUpdate, event)
Expand Down Expand Up @@ -72,6 +75,38 @@ func (p *hookBuiltinHTTPPlugin) OnHookEvent(event HookEvent) error {
if p.hub.cfg.OnHlsMakeTs != "" {
p.hub.asyncPostEvent(p.hub.cfg.OnHlsMakeTs, event)
}
case HookEventStreamChanged:
if p.hub.cfg.ZlmOnStreamChanged != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnStreamChanged, event)
}
case HookEventServerKeepalive:
if p.hub.cfg.ZlmOnServerKeepalive != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnServerKeepalive, event)
}
case HookEventStreamNoneReader:
if p.hub.cfg.ZlmOnStreamNoneReader != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnStreamNoneReader, event)
}
case HookEventRtpServerTimeout:
if p.hub.cfg.ZlmOnRtpServerTimeout != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnRtpServerTimeout, event)
}
case HookEventRecordMp4:
if p.hub.cfg.ZlmOnRecordMp4 != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnRecordMp4, event)
}
case HookEventPublish:
if p.hub.cfg.ZlmOnPublish != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnPublish, event)
}
case HookEventPlay:
if p.hub.cfg.ZlmOnPlay != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnPlay, event)
}
case HookEventStreamNotFound:
if p.hub.cfg.ZlmOnStreamNotFound != "" {
p.hub.asyncPostEvent(p.hub.cfg.ZlmOnStreamNotFound, event)
}
}

return nil
Expand Down
Loading
Loading