@@ -7,85 +7,52 @@ import (
77 "net/http"
88 "net/url"
99 "sync/atomic"
10- "time"
11-
12- "github.com/smallnest/ringbuffer"
1310)
1411
15- const defaultReadBufSize = 32 * 1024 / 2
16-
17- // StreamServer serves transcode output over HTTP using a blocking ring
18- // buffer. Write blocks when full, applying backpressure to ffmpeg so
19- // no data is ever silently overwritten.
12+ // StreamServer serves an io.Reader over HTTP. The consumer's read
13+ // pace drives the producer through OS pipe backpressure.
2014type StreamServer struct {
21- ring * ringbuffer.RingBuffer
22- cancel context.CancelFunc
15+ reader io.Reader
2316 active atomic.Bool
2417 contentType string
2518 extension string
2619 headers map [string ]string
2720 listener net.Listener
2821 server * http.Server
29- errCh chan error
22+ done chan struct {}
3023}
3124
32- // StreamServerConfig holds configuration for creating a StreamServer .
25+ // StreamServerConfig holds the parameters for NewStreamServer .
3326type StreamServerConfig struct {
34- LocalIP string
35- ContentType string
36- Extension string
37- Headers map [string ]string
38- BufferCapacity int
27+ LocalIP string
28+ ContentType string
29+ Extension string
30+ Headers map [string ]string
3931}
4032
41- // NewStreamServer creates, starts, and returns a StreamServer. It begins
42- // ingesting from stream and serving HTTP immediately.
43- func NewStreamServer (ctx context.Context , cfg StreamServerConfig , stream io.Reader ) (* StreamServer , error ) {
33+ func NewStreamServer (cfg StreamServerConfig , reader io.Reader ) (* StreamServer , error ) {
4434 ln , err := net .Listen ("tcp" , cfg .LocalIP + ":0" )
4535 if err != nil {
4636 return nil , err
4737 }
48-
49- ctx , cancel := context .WithCancel (ctx )
50-
5138 s := & StreamServer {
52- ring : ringbuffer .New (cfg .BufferCapacity ).SetBlocking (true ),
53- cancel : cancel ,
39+ reader : reader ,
5440 contentType : cfg .ContentType ,
5541 extension : cfg .Extension ,
5642 headers : cfg .Headers ,
5743 listener : ln ,
58- errCh : make (chan error , 1 ),
44+ done : make (chan struct {} ),
5945 }
60-
61- s .ring .WithCancel (ctx )
62-
6346 mux := http .NewServeMux ()
6447 mux .HandleFunc ("/stream" + cfg .Extension , s .handleStream )
65-
6648 s .server = & http.Server {Handler : mux }
67-
68- go func () {
69- s .ring .ReadFrom (stream )
70- s .ring .CloseWriter ()
71- }()
72-
7349 go func () {
74- if err := s .server .Serve (s .listener ); err != nil && err != http .ErrServerClosed {
75- s .errCh <- err
76- }
77- close (s .errCh )
50+ s .server .Serve (ln )
51+ close (s .done )
7852 }()
79-
8053 return s , nil
8154}
8255
83- // Stop cancels ingestion and closes the server.
84- func (s * StreamServer ) Stop () {
85- s .cancel ()
86- s .server .Close ()
87- }
88-
8956// URL returns the full URL the server is listening on.
9057func (s * StreamServer ) URL () * url.URL {
9158 return & url.URL {
@@ -95,72 +62,59 @@ func (s *StreamServer) URL() *url.URL {
9562 }
9663}
9764
65+ // Close shuts down the server.
66+ func (s * StreamServer ) Close () error { return s .server .Close () }
67+
9868// Wait blocks until the server exits or the context is cancelled.
9969func (s * StreamServer ) Wait (ctx context.Context ) error {
10070 select {
101- case err := <- s .errCh :
102- return err
103- case <- ctx .Done ():
71+ case <- s .done :
10472 return nil
105- }
106- }
107-
108- // WaitForData blocks until the buffer has at least minBytes of data or
109- // the context is cancelled.
110- func (s * StreamServer ) WaitForData (ctx context.Context , minBytes int ) error {
111- const tick = 50 * time .Millisecond
112- ticker := time .NewTicker (tick )
113- defer ticker .Stop ()
114-
115- for {
116- if s .ring .Length () >= minBytes {
117- return nil
118- }
119- select {
120- case <- ctx .Done ():
121- return ctx .Err ()
122- case <- ticker .C :
123- }
73+ case <- ctx .Done ():
74+ return ctx .Err ()
12475 }
12576}
12677
12778func (s * StreamServer ) handleStream (w http.ResponseWriter , r * http.Request ) {
128- if ! s .active .CompareAndSwap (false , true ) {
129- http .Error (w , "stream already has an active reader" , http .StatusServiceUnavailable )
130- return
131- }
132- defer s .active .Store (false )
133-
13479 w .Header ().Set ("Content-Type" , s .contentType )
13580 for k , v := range s .headers {
13681 w .Header ().Set (k , v )
13782 }
138-
13983 if r .Method == http .MethodHead {
14084 w .WriteHeader (http .StatusOK )
14185 return
14286 }
14387
144- w .WriteHeader (http .StatusOK )
88+ if ! s .active .CompareAndSwap (false , true ) {
89+ http .Error (w , "stream already active" , http .StatusServiceUnavailable )
90+ return
91+ }
92+ defer s .active .Store (false )
14593
94+ w .WriteHeader (http .StatusOK )
14695 if f , ok := w .(http.Flusher ); ok {
14796 f .Flush ()
14897 }
14998
150- buf := make ([]byte , defaultReadBufSize )
151-
99+ buf := make ([]byte , 32 * 1024 )
100+ readerDone := false
152101 for {
153- n , err := s .ring .Read (buf )
102+ n , err := s .reader .Read (buf )
154103 if n > 0 {
155- if _ , writeErr := w .Write (buf [:n ]); writeErr != nil {
156- return
104+ if _ , we := w .Write (buf [:n ]); we != nil {
105+ break // TV disconnected — keep server alive
157106 }
158107 if f , ok := w .(http.Flusher ); ok {
159108 f .Flush ()
160109 }
161110 }
162111 if err != nil {
163- return
112+ readerDone = true
113+ break
164114 }
165115 }
116+ if readerDone {
117+ // Shut down so Wait() unblocks. Goroutine avoids handler deadlock.
118+ go s .server .Close ()
119+ }
166120}
0 commit comments