Files
mediamtx/internal/stream/reader.go
T
Alessandro RosandGitHub f87d9e659e hls: track sessions (#962) (#5683)
sessions are now tracked through cookies or query parameters.

This provides the ability to inspect sessions through logs, metrics and
API, allows more precise tracking of outbound bytes, decreases load on
external HTTP authentication URLs since they are now called once per
session and not once per request.
2026-04-25 21:10:34 +02:00

132 lines
2.8 KiB
Go

package stream
import (
"fmt"
"github.com/bluenviron/gortsplib/v5/pkg/description"
"github.com/bluenviron/gortsplib/v5/pkg/format"
"github.com/bluenviron/gortsplib/v5/pkg/ringbuffer"
"github.com/bluenviron/mediamtx/internal/counterdumper"
"github.com/bluenviron/mediamtx/internal/logger"
"github.com/bluenviron/mediamtx/internal/unit"
)
// OnDataFunc is the callback passed to OnData().
type OnDataFunc func(*unit.Unit) error
// Reader is a stream reader.
type Reader struct {
SkipOutboundBytes bool
Parent logger.Writer
onDatas map[*description.Media]map[format.Format]OnDataFunc
queueSize int
buffer *ringbuffer.RingBuffer
outboundFramesDiscarded *counterdumper.Dumper
// out
err chan error
}
// OnData registers a callback that is called when data from given format is available.
func (r *Reader) OnData(medi *description.Media, forma format.Format, cb OnDataFunc) {
if r.onDatas == nil {
r.onDatas = make(map[*description.Media]map[format.Format]OnDataFunc)
}
if r.onDatas[medi] == nil {
r.onDatas[medi] = make(map[format.Format]OnDataFunc)
}
r.onDatas[medi][forma] = cb
}
// Formats returns all formats for which the reader has registered a OnData callback.
func (r *Reader) Formats() []format.Format {
n := 0
for _, formats := range r.onDatas {
for range formats {
n++
}
}
if n == 0 {
return nil
}
out := make([]format.Format, n)
n = 0
for _, formats := range r.onDatas {
for forma := range formats {
out[n] = forma
n++
}
}
return out
}
// OutboundFramesDiscarded returns the number of frames discarded because the reader is too slow.
func (r *Reader) OutboundFramesDiscarded() uint64 {
return r.outboundFramesDiscarded.Get()
}
// error returns whenever there's an error.
// It can be called only after stream.AddReader().
func (r *Reader) Error() chan error {
return r.err
}
func (r *Reader) start() {
buffer, _ := ringbuffer.New(uint64(r.queueSize))
r.buffer = buffer
r.err = make(chan error)
r.outboundFramesDiscarded = &counterdumper.Dumper{
OnReport: func(val uint64) {
r.Parent.Log(logger.Warn, "reader is too slow, discarding %d %s",
val,
func() string {
if val == 1 {
return "frame"
}
return "frames"
}())
},
}
r.outboundFramesDiscarded.Start()
go r.run()
}
func (r *Reader) stop() {
r.buffer.Close()
r.outboundFramesDiscarded.Stop()
<-r.err
}
func (r *Reader) run() {
r.err <- r.runInner()
close(r.err)
}
func (r *Reader) runInner() error {
for {
cb, ok := r.buffer.Pull()
if !ok {
return fmt.Errorf("terminated")
}
err := cb.(func() error)()
if err != nil {
return err
}
}
}
func (r *Reader) push(cb func() error) {
ok := r.buffer.Push(cb)
if !ok {
r.outboundFramesDiscarded.Increase()
}
}