Code refactoring for #1762

This commit is contained in:
Alex X
2025-10-01 16:57:39 +03:00
parent c196b82a72
commit 22cc8ac2c4
5 changed files with 79 additions and 188 deletions
+26 -54
View File
@@ -5,6 +5,7 @@ import (
"github.com/AlexxIT/go2rtc/internal/api" "github.com/AlexxIT/go2rtc/internal/api"
"github.com/AlexxIT/go2rtc/internal/app" "github.com/AlexxIT/go2rtc/internal/app"
"github.com/AlexxIT/go2rtc/pkg/core"
"github.com/AlexxIT/go2rtc/pkg/probe" "github.com/AlexxIT/go2rtc/pkg/probe"
) )
@@ -27,7 +28,7 @@ func apiStreams(w http.ResponseWriter, r *http.Request) {
return return
} }
cons := probe.NewProbe(query) cons := probe.Create("probe", query)
if len(cons.Medias) != 0 { if len(cons.Medias) != 0 {
cons.WithRequest(r) cons.WithRequest(r)
if err := stream.AddConsumer(cons); err != nil { if err := stream.AddConsumer(cons); err != nil {
@@ -126,73 +127,44 @@ func apiStreamsDOT(w http.ResponseWriter, r *http.Request) {
func apiPreload(w http.ResponseWriter, r *http.Request) { func apiPreload(w http.ResponseWriter, r *http.Request) {
query := r.URL.Query() query := r.URL.Query()
src := query.Get("src") src := query.Get("src")
query.Del("src")
if src == "" { // check if stream exists
http.Error(w, "no source", http.StatusBadRequest) stream := Get(src)
if stream == nil {
http.Error(w, "", http.StatusNotFound)
return return
} }
switch r.Method { switch r.Method {
case "PUT": case "PUT":
// check if stream exists // it's safe to delete from map while iterating
stream := Get(src) for k := range query {
if stream == nil { switch k {
http.Error(w, "stream not found", http.StatusNotFound) case core.KindVideo, core.KindAudio, "microphone":
default:
delete(query, k)
}
}
rawQuery := query.Encode()
if err := AddPreload(stream, rawQuery); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return return
} }
// check if consumer exists
if cons, ok := preloads[src]; ok {
stream.RemoveConsumer(cons)
delete(preloads, src)
}
// parse query parameters
var rawQuery string
if query.Has("video") {
if videoQuery := query.Get("video"); videoQuery != "" {
rawQuery += "video=" + videoQuery + "#"
} else {
rawQuery += "video#"
}
}
if query.Has("audio") {
if audioQuery := query.Get("audio"); audioQuery != "" {
rawQuery += "audio=" + audioQuery + "#"
} else {
rawQuery += "audio#"
}
}
if query.Has("microphone") {
if micQuery := query.Get("microphone"); micQuery != "" {
rawQuery += "microphone=" + micQuery + "#"
} else {
rawQuery += "microphone#"
}
}
if err := app.PatchConfig([]string{"preload", src}, rawQuery); err != nil { if err := app.PatchConfig([]string{"preload", src}, rawQuery); err != nil {
log.Error().Err(err).Str("src", src).Msg("Failed to patch config for PUT") http.Error(w, err.Error(), http.StatusInternalServerError)
http.Error(w, err.Error(), http.StatusBadRequest) }
case "DELETE":
if err := DelPreload(stream); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return return
} }
Preload(src, rawQuery)
case "DELETE":
if cons, ok := preloads[src]; ok {
if stream := Get(src); stream != nil {
stream.RemoveConsumer(cons)
} else {
cons.Stop()
}
delete(preloads, src)
}
if err := app.PatchConfig([]string{"preload", src}, nil); err != nil { if err := app.PatchConfig([]string{"preload", src}, nil); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest) http.Error(w, err.Error(), http.StatusInternalServerError)
} }
default: default:
+40 -16
View File
@@ -1,34 +1,58 @@
package streams package streams
import ( import (
"errors"
"net/url" "net/url"
"sync"
"github.com/AlexxIT/go2rtc/pkg/preload" "github.com/AlexxIT/go2rtc/pkg/probe"
) )
var preloads = map[string]*preload.Preload{} var preloads = map[*Stream]*probe.Probe{}
var preloadsMu sync.Mutex
func (s *Stream) Preload(name string, query url.Values) error { func Preload(stream *Stream, rawQuery string) {
cons := preload.NewPreload(name, query) if err := AddPreload(stream, rawQuery); err != nil {
preloads[name] = cons log.Error().Err(err).Caller().Send()
}
}
if err := s.AddConsumer(cons); err != nil { func AddPreload(stream *Stream, rawQuery string) error {
if rawQuery == "" {
rawQuery = "video&audio"
}
query, err := url.ParseQuery(rawQuery)
if err != nil {
return err return err
} }
preloadsMu.Lock()
defer preloadsMu.Unlock()
if cons := preloads[stream]; cons != nil {
stream.RemoveConsumer(cons)
}
cons := probe.Create("preload", query)
if err = stream.AddConsumer(cons); err != nil {
return err
}
preloads[stream] = cons
return nil return nil
} }
func Preload(src string, rawQuery string) { func DelPreload(stream *Stream) error {
// skip if exists preloadsMu.Lock()
if _, ok := preloads[src]; ok { defer preloadsMu.Unlock()
return
if cons := preloads[stream]; cons != nil {
stream.RemoveConsumer(cons)
delete(preloads, stream)
return nil
} }
if stream := Get(src); stream != nil { return errors.New("streams: preload not found")
query := ParseQuery(rawQuery)
if err := stream.Preload(src, query); err != nil {
log.Error().Err(err).Caller().Send()
}
}
} }
+3 -5
View File
@@ -36,17 +36,15 @@ func Init() {
} }
time.AfterFunc(time.Second, func() { time.AfterFunc(time.Second, func() {
if cfg.Publish != nil { // range for nil map is OK
for name, dst := range cfg.Publish { for name, dst := range cfg.Publish {
if stream := Get(name); stream != nil { if stream := Get(name); stream != nil {
Publish(stream, dst) Publish(stream, dst)
} }
} }
}
if cfg.Preload != nil {
for name, rawQuery := range cfg.Preload { for name, rawQuery := range cfg.Preload {
Preload(name, rawQuery) if stream := Get(name); stream != nil {
Preload(stream, rawQuery)
} }
} }
}) })
-82
View File
@@ -1,82 +0,0 @@
package preload
import (
"net/url"
"strings"
"github.com/AlexxIT/go2rtc/pkg/core"
"github.com/pion/rtp"
)
type Preload struct {
core.Connection
closed core.Waiter
}
func NewPreload(name string, query url.Values) *Preload {
medias := core.ParseQuery(query)
for _, value := range query["microphone"] {
media := &core.Media{Kind: core.KindAudio, Direction: core.DirectionRecvonly}
for _, name := range strings.Split(value, ",") {
name = strings.ToUpper(name)
switch name {
case "", "COPY":
name = core.CodecAny
}
media.Codecs = append(media.Codecs, &core.Codec{Name: name})
}
medias = append(medias, media)
}
if len(medias) == 0 {
medias = []*core.Media{
{
Kind: core.KindVideo,
Direction: core.DirectionSendonly,
Codecs: []*core.Codec{{Name: core.CodecAny}},
},
{
Kind: core.KindAudio,
Direction: core.DirectionSendonly,
Codecs: []*core.Codec{{Name: core.CodecAny}},
},
}
}
return &Preload{
Connection: core.Connection{
ID: core.NewID(),
FormatName: "preload",
Medias: medias,
},
}
}
func (p *Preload) AddTrack(media *core.Media, codec *core.Codec, track *core.Receiver) error {
sender := core.NewSender(media, track.Codec)
sender.Handler = func(pkt *rtp.Packet) {
p.Send += pkt.MarshalSize()
}
sender.HandleRTP(track)
p.Senders = append(p.Senders, sender)
return nil
}
func (p *Preload) Start() error {
p.closed.Wait()
return nil
}
func (p *Preload) Stop() error {
for _, receiver := range p.Receivers {
receiver.Close()
}
for _, sender := range p.Senders {
sender.Close()
}
p.closed.Done(nil)
return nil
}
@@ -11,7 +11,7 @@ type Probe struct {
core.Connection core.Connection
} }
func NewProbe(query url.Values) *Probe { func Create(name string, query url.Values) *Probe {
medias := core.ParseQuery(query) medias := core.ParseQuery(query)
for _, value := range query["microphone"] { for _, value := range query["microphone"] {
@@ -32,39 +32,18 @@ func NewProbe(query url.Values) *Probe {
return &Probe{ return &Probe{
Connection: core.Connection{ Connection: core.Connection{
ID: core.NewID(), ID: core.NewID(),
FormatName: "probe", FormatName: name,
Medias: medias, Medias: medias,
}, },
} }
} }
func (p *Probe) GetMedias() []*core.Media {
return p.Medias
}
func (p *Probe) AddTrack(media *core.Media, codec *core.Codec, track *core.Receiver) error { func (p *Probe) AddTrack(media *core.Media, codec *core.Codec, track *core.Receiver) error {
sender := core.NewSender(media, track.Codec) sender := core.NewSender(media, track.Codec)
sender.Bind(track) sender.Handler = func(pkt *core.Packet) {
p.Send += len(pkt.Payload)
}
sender.HandleRTP(track)
p.Senders = append(p.Senders, sender) p.Senders = append(p.Senders, sender)
return nil return nil
} }
func (p *Probe) GetTrack(media *core.Media, codec *core.Codec) (*core.Receiver, error) {
receiver := core.NewReceiver(media, codec)
p.Receivers = append(p.Receivers, receiver)
return receiver, nil
}
func (p *Probe) Start() error {
return nil
}
func (p *Probe) Stop() error {
for _, receiver := range p.Receivers {
receiver.Close()
}
for _, sender := range p.Senders {
sender.Close()
}
return nil
}