From dcd4b388385232c69414ff276708edecfba551eb Mon Sep 17 00:00:00 2001 From: Alessandro Ros Date: Tue, 4 Aug 2026 22:48:57 +0200 Subject: [PATCH] srt: close sources immediately when path is closed (#6038) --- go.mod | 2 ++ go.sum | 4 ++-- internal/forward/rtmp/dest.go | 4 ++-- internal/forward/rtsp/dest.go | 2 +- internal/forward/srt/dest.go | 26 ++++++-------------------- internal/staticsources/srt/source.go | 2 +- 6 files changed, 14 insertions(+), 26 deletions(-) diff --git a/go.mod b/go.mod index 0c8f90ed..73fadcd5 100644 --- a/go.mod +++ b/go.mod @@ -107,3 +107,5 @@ require ( gopkg.in/warnings.v0 v0.1.2 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) + +replace github.com/datarhei/gosrt => github.com/aler9/gosrt v0.0.0-20260804162707-b6afcc8add8f diff --git a/go.sum b/go.sum index b2500dcb..22772e2d 100644 --- a/go.sum +++ b/go.sum @@ -23,6 +23,8 @@ github.com/alecthomas/kong v1.16.0 h1:g92/kUxBcdcTPOM79yE63viJgtcp5dNyrB3/O2cjYT github.com/alecthomas/kong v1.16.0/go.mod h1:wrlbXem1CWqUV5Vbmss5ISYhsVPkBb1Yo7YKJghju2I= github.com/alecthomas/repr v0.5.2 h1:SU73FTI9D1P5UNtvseffFSGmdNci/O6RsqzeXJtP0Qs= github.com/alecthomas/repr v0.5.2/go.mod h1:Fr0507jx4eOXV7AlPV6AVZLYrLIuIeSOWtW57eE/O/4= +github.com/aler9/gosrt v0.0.0-20260804162707-b6afcc8add8f h1:wGCt98qRB2SOKMK5tdWM10sjyQa6WtrfsBrx0pRFLLI= +github.com/aler9/gosrt v0.0.0-20260804162707-b6afcc8add8f/go.mod h1:Jc34G/QwMzPMJdHNcaxc9ddbx95h2FumO2h/rgV/eRA= github.com/anmitsu/go-shlex v0.0.0-20200514113438-38f4b401e2be h1:9AeTilPcZAjCFIImctFaOjnTIavg87rW78vTPkQqLI8= github.com/anmitsu/go-shlex v0.0.0-20200514113438-38f4b401e2be/go.mod h1:ySMOLuWl6zY27l47sB3qLNK6tF2fkHG55UZxx8oIVo4= github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5 h1:0CwZNZbxp69SHPdPJAN/hZIm0C4OItdklCFmMRWYpio= @@ -54,8 +56,6 @@ github.com/cloudwego/base64x v0.1.6/go.mod h1:OFcloc187FXDaYHvrNIjxSe8ncn0OOM8gE github.com/creack/pty v1.1.7/go.mod h1:lj5s0c3V2DBrqTV7llrYr5NG6My20zk30Fl46Y7DoTY= github.com/cyphar/filepath-securejoin v0.6.1 h1:5CeZ1jPXEiYt3+Z6zqprSAgSWiggmpVyciv8syjIpVE= github.com/cyphar/filepath-securejoin v0.6.1/go.mod h1:A8hd4EnAeyujCJRrICiOWqjS1AX0a9kM5XL+NwKoYSc= -github.com/datarhei/gosrt v0.11.0 h1:g3dGowSxrD1Oxr0Us6/w7x9bqzHHjBYu+EA3tNrhDeg= -github.com/datarhei/gosrt v0.11.0/go.mod h1:F5B25N3CFf68K4igNLQ1iARcKDbkv8riymjT8l5cbLg= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= diff --git a/internal/forward/rtmp/dest.go b/internal/forward/rtmp/dest.go index e3b9bfc2..f747507a 100644 --- a/internal/forward/rtmp/dest.go +++ b/internal/forward/rtmp/dest.go @@ -104,12 +104,12 @@ func (d *Dest) Run(ctx context.Context) error { err = conn.Initialize(ctx) if err != nil { - return fmt.Errorf("connect RTMP destination: %w", err) + return err } terminate := make(chan struct{}) - errChan := make(chan error, 1) + errChan := make(chan error) go func() { errChan <- d.runInner(conn, terminate) }() diff --git a/internal/forward/rtsp/dest.go b/internal/forward/rtsp/dest.go index da912a9b..7bf32f9d 100644 --- a/internal/forward/rtsp/dest.go +++ b/internal/forward/rtsp/dest.go @@ -77,7 +77,7 @@ func (d *Dest) Run(ctx context.Context) error { terminate := make(chan struct{}) - errChan := make(chan error, 1) + errChan := make(chan error) go func() { errChan <- d.runInner(client, desc, terminate) }() diff --git a/internal/forward/srt/dest.go b/internal/forward/srt/dest.go index 25efb0ff..e78e956a 100644 --- a/internal/forward/srt/dest.go +++ b/internal/forward/srt/dest.go @@ -69,43 +69,29 @@ func (d *Dest) Run(ctx context.Context) error { terminate := make(chan struct{}) - type runResult struct { - err error - } - - errChan := make(chan runResult, 1) + errChan := make(chan error) go func() { - errChan <- runResult{err: d.runInner(ctx, address, srtConf, terminate)} + errChan <- d.runInner(ctx, address, srtConf, terminate) }() select { - case res := <-errChan: - return res.err + case err = <-errChan: + return err case <-ctx.Done(): close(terminate) + <-errChan return fmt.Errorf("terminated") } } func (d *Dest) runInner(ctx context.Context, address string, srtConf srtlib.Config, terminate <-chan struct{}) error { - conn, err := srtlib.Dial("srt", address, srtConf) + conn, err := srtlib.DialWithContext(ctx, "srt", address, srtConf) if err != nil { - select { - case <-ctx.Done(): - return nil - default: - } return err } defer conn.Close() - select { - case <-ctx.Done(): - return nil - default: - } - d.mutex.Lock() d.outboundBytesFunc = func() uint64 { var stats srtlib.Statistics diff --git a/internal/staticsources/srt/source.go b/internal/staticsources/srt/source.go index 8fd46192..ad8cb117 100644 --- a/internal/staticsources/srt/source.go +++ b/internal/staticsources/srt/source.go @@ -47,7 +47,7 @@ func (s *Source) Run(params defs.StaticSourceRunParams) error { return err } - sconn, err := srt.Dial("srt", address, conf) + sconn, err := srt.DialWithContext(params.Context, "srt", address, conf) if err != nil { return err }