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
34 changes: 28 additions & 6 deletions gb28181/mediaserver/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,9 +51,11 @@ type Conn struct {
buffer *bytes.Buffer
key string

mediaServer *GB28181MediaServer
one sync.Once
oneSaveConn sync.Once
mediaServer *GB28181MediaServer
preferMediaKeyLookup bool
readTimeout time.Duration
one sync.Once
oneSaveConn sync.Once
}

func NewConn(conn net.Conn, observer IGbObserver, lal logic.ILalServer) *Conn {
Expand All @@ -76,12 +78,18 @@ func (c *Conn) SetMediaServer(mediaServer *GB28181MediaServer) {
func (c *Conn) SetKey(key string) {
c.key = key
}
func (c *Conn) SetPreferMediaKeyLookup(prefer bool) {
c.preferMediaKeyLookup = prefer
}
func (c *Conn) SetReadTimeout(timeout time.Duration) {
c.readTimeout = timeout
}
func (c *Conn) Serve() (err error) {
defer func() {
nazalog.Info("conn close, err:", err)
c.Close()

if c.observer != nil {
if c.observer != nil && c.streamName != "" {
c.observer.NotifyClose(c.streamName)
}
if c.psDumpFile != nil {
Expand All @@ -95,7 +103,9 @@ func (c *Conn) Serve() (err error) {
nazalog.Info("gb28181 conn, remoteaddr:", c.conn.RemoteAddr().String(), " localaddr:", c.conn.LocalAddr().String())

for {
c.conn.SetReadDeadline(time.Now().Add(10 * time.Second))
if c.readTimeout > 0 {
c.conn.SetReadDeadline(time.Now().Add(c.readTimeout))
}
pkt := &rtp.Packet{}
if c.conn.RemoteAddr().Network() == "udp" {
buf := make([]byte, 1472*4)
Expand Down Expand Up @@ -132,7 +142,13 @@ func (c *Conn) Serve() (err error) {
if !c.check && c.observer != nil {
var mediaInfo *MediaInfo
var ok bool
if pkt.SSRC != 0 {
if c.preferMediaKeyLookup {
mediaInfo, ok = c.observer.GetMediaInfoByKey(c.key)
if !ok {
nazalog.Error("get mediaInfo :", c.key)
return fmt.Errorf("get mediaInfo:%s", c.key)
}
} else if pkt.SSRC != 0 {
mediaInfo, ok = c.observer.CheckSsrc(pkt.SSRC)
if !ok {
nazalog.Error("invalid ssrc:", pkt.SSRC)
Expand All @@ -159,6 +175,9 @@ func (c *Conn) Serve() (err error) {
}
}
nazalog.Info("gb28181 ssrc check success, streamName:", c.streamName)
if c.observer != nil {
c.observer.OnRtpPacket(c.streamName, c.key)
}

session, err := c.lalServer.AddCustomizePubSession(mediaInfo.StreamName)
if err != nil {
Expand All @@ -174,6 +193,9 @@ func (c *Conn) Serve() (err error) {
c.lalSession = session
}
c.rtpPts = uint64(pkt.Header.Timestamp)
if c.observer != nil && c.streamName != "" {
c.observer.OnRtpPacket(c.streamName, c.key)
}
if c.demuxer != nil {
if c.psDumpFile != nil {
c.psDumpFile.WriteWithType(pkt.Payload, base.DumpTypePsRtpData)
Expand Down
56 changes: 44 additions & 12 deletions gb28181/mediaserver/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,20 @@ import (
"errors"
"net"
"sync"
"sync/atomic"
"time"

"github.com/q191201771/lal/pkg/logic"
"github.com/q191201771/naza/pkg/nazalog"
)

const defaultReadTimeout = 10 * time.Second

type IGbObserver interface {
CheckSsrc(ssrc uint32) (*MediaInfo, bool)
GetMediaInfoByKey(key string) (*MediaInfo, bool)
NotifyClose(streamName string)
OnRtpPacket(streamName string, mediaKey string)
}

type GB28181MediaServer struct {
Expand All @@ -22,52 +26,79 @@ type GB28181MediaServer struct {

listener net.Listener

disposeOnce sync.Once
observer IGbObserver
mediaKey string
disposeOnce sync.Once
disposed atomic.Bool
observer IGbObserver
mediaKey string
preferMediaKeyLookup bool
readTimeout time.Duration

conns sync.Map //增加链接对象,目前只适用于多端口
}

func NewGB28181MediaServer(listenPort int, mediaKey string, observer IGbObserver, lal logic.ILalServer) *GB28181MediaServer {
return &GB28181MediaServer{
listenPort: listenPort,
lalServer: lal,
observer: observer,
mediaKey: mediaKey,
listenPort: listenPort,
lalServer: lal,
observer: observer,
mediaKey: mediaKey,
readTimeout: defaultReadTimeout,
}
}

func (s *GB28181MediaServer) WithPreferMediaKeyLookup(prefer bool) *GB28181MediaServer {
s.preferMediaKeyLookup = prefer
return s
}

func (s *GB28181MediaServer) WithReadTimeout(timeout time.Duration) *GB28181MediaServer {
s.readTimeout = timeout
return s
}

func (s *GB28181MediaServer) GetListenerPort() uint16 {
return uint16(s.listenPort)
}
func (s *GB28181MediaServer) Start(listener net.Listener) (err error) {
s.listener = listener
if s.listener != nil {
go func() {
if listener != nil {
go func(listener net.Listener) {
for {
if s.listener == nil {
if s.disposed.Load() {
return
}
conn, err := s.listener.Accept()
conn, err := listener.Accept()
if err != nil {
var ne net.Error
if ok := errors.As(err, &ne); ok && ne.Timeout() {
nazalog.Error("Accept failed: timeout error, retrying...")
time.Sleep(time.Second / 20)
continue
} else {
break
}
}
if conn == nil {
continue
}
if s.disposed.Load() {
conn.Close()
return
}

c := NewConn(conn, s.observer, s.lalServer)
c.SetKey(s.mediaKey)
c.SetMediaServer(s)
c.SetPreferMediaKeyLookup(s.preferMediaKeyLookup)
c.SetReadTimeout(s.readTimeout)
s.conns.Store(c, c)
go func() {
c.Serve()
s.conns.Delete(c)
s.conns.Delete(c.streamName)
}()
}
}()
}(listener)
}
return
}
Expand All @@ -79,6 +110,7 @@ func (s *GB28181MediaServer) CloseConn(streamName string) {
}
func (s *GB28181MediaServer) Dispose() {
s.disposeOnce.Do(func() {
s.disposed.Store(true)
s.conns.Range(func(_, value any) bool {
conn := value.(*Conn)
conn.Close()
Expand Down
Loading
Loading