2021-08-10 16:34:10 +00:00
|
|
|
package core
|
|
|
|
|
|
|
|
import (
|
|
|
|
"sync"
|
|
|
|
|
|
|
|
"github.com/aler9/gortsplib"
|
|
|
|
)
|
|
|
|
|
|
|
|
type streamNonRTSPReadersMap struct {
|
|
|
|
mutex sync.RWMutex
|
|
|
|
ma map[reader]struct{}
|
|
|
|
}
|
|
|
|
|
|
|
|
func newStreamNonRTSPReadersMap() *streamNonRTSPReadersMap {
|
|
|
|
return &streamNonRTSPReadersMap{
|
|
|
|
ma: make(map[reader]struct{}),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *streamNonRTSPReadersMap) close() {
|
|
|
|
m.mutex.Lock()
|
|
|
|
defer m.mutex.Unlock()
|
|
|
|
m.ma = nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *streamNonRTSPReadersMap) add(r reader) {
|
|
|
|
m.mutex.Lock()
|
|
|
|
defer m.mutex.Unlock()
|
|
|
|
m.ma[r] = struct{}{}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *streamNonRTSPReadersMap) remove(r reader) {
|
|
|
|
m.mutex.Lock()
|
|
|
|
defer m.mutex.Unlock()
|
|
|
|
delete(m.ma, r)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *streamNonRTSPReadersMap) forwardFrame(trackID int, streamType gortsplib.StreamType, payload []byte) {
|
|
|
|
m.mutex.RLock()
|
|
|
|
defer m.mutex.RUnlock()
|
|
|
|
|
|
|
|
for c := range m.ma {
|
2021-10-27 19:01:00 +00:00
|
|
|
c.onReaderFrame(trackID, streamType, payload)
|
2021-08-10 16:34:10 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
type stream struct {
|
|
|
|
nonRTSPReaders *streamNonRTSPReadersMap
|
|
|
|
rtspStream *gortsplib.ServerStream
|
|
|
|
}
|
|
|
|
|
|
|
|
func newStream(tracks gortsplib.Tracks) *stream {
|
|
|
|
s := &stream{
|
|
|
|
nonRTSPReaders: newStreamNonRTSPReadersMap(),
|
|
|
|
rtspStream: gortsplib.NewServerStream(tracks),
|
|
|
|
}
|
|
|
|
return s
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *stream) close() {
|
|
|
|
s.nonRTSPReaders.close()
|
|
|
|
s.rtspStream.Close()
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *stream) tracks() gortsplib.Tracks {
|
|
|
|
return s.rtspStream.Tracks()
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *stream) readerAdd(r reader) {
|
|
|
|
if _, ok := r.(pathRTSPSession); !ok {
|
|
|
|
s.nonRTSPReaders.add(r)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *stream) readerRemove(r reader) {
|
|
|
|
if _, ok := r.(pathRTSPSession); !ok {
|
|
|
|
s.nonRTSPReaders.remove(r)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *stream) onFrame(trackID int, streamType gortsplib.StreamType, payload []byte) {
|
|
|
|
// forward to RTSP readers
|
|
|
|
s.rtspStream.WriteFrame(trackID, streamType, payload)
|
|
|
|
|
|
|
|
// forward to non-RTSP readers
|
|
|
|
s.nonRTSPReaders.forwardFrame(trackID, streamType, payload)
|
|
|
|
}
|