@@ -24,6 +24,7 @@ const MaxFileSize = 50 * 1024 * 1024 // 50 MB
2424
2525type Proxy struct {
2626 sinks []* Sink
27+ sinkMutex sync.Mutex
2728 errors chan error
2829 transport * http.Transport
2930 ctx context.Context
@@ -32,7 +33,7 @@ type Proxy struct {
3233 metrics * ProxyMetrics
3334}
3435
35- func NewProxy (ctx context.Context , addr string , sinks [] * Sink ) (* Proxy , error ) {
36+ func NewProxy (ctx context.Context , conf Config ) (* Proxy , error ) {
3637 tr := & http.Transport {
3738 DialContext : (& net.Dialer {
3839 Timeout : 30 * time .Second ,
@@ -44,27 +45,25 @@ func NewProxy(ctx context.Context, addr string, sinks []*Sink) (*Proxy, error) {
4445 TLSHandshakeTimeout : 10 * time .Second ,
4546 ExpectContinueTimeout : 1 * time .Second ,
4647 }
48+ reg := prometheus .NewRegistry ()
49+ reg .MustRegister (
50+ collectors .NewProcessCollector (collectors.ProcessCollectorOpts {}),
51+ )
4752 p := & Proxy {
48- sinks : sinks ,
4953 transport : tr ,
5054 errors : make (chan error , 1 ),
5155 ctx : ctx ,
56+ metrics : NewProxyMetrics (reg ),
5257 }
5358
5459 mux := http .NewServeMux ()
55- srv := http.Server {Addr : addr , Handler : mux }
56-
57- reg := prometheus .NewRegistry ()
58- reg .MustRegister (
59- collectors .NewProcessCollector (collectors.ProcessCollectorOpts {}),
60- )
61- p .metrics = NewProxyMetrics (reg )
60+ srv := http.Server {Addr : conf .ListenAddress , Handler : mux }
6261
6362 // set routes
6463 mux .HandleFunc ("/" , p .HandleUpload )
6564 mux .Handle ("/metrics" , promhttp .HandlerFor (reg , promhttp.HandlerOpts {}))
6665
67- ln , err := net .Listen ("tcp" , addr )
66+ ln , err := net .Listen ("tcp" , conf . ListenAddress )
6867 if err != nil {
6968 return nil , err
7069 }
@@ -91,22 +90,100 @@ func NewProxy(ctx context.Context, addr string, sinks []*Sink) (*Proxy, error) {
9190 }
9291 }()
9392
94- // run sink uploaders
95- for _ , sink := range p .sinks {
96- log .Printf ("setup sink %+v\n " , sink )
93+ p .done .Add (1 )
94+ go p .runTimeout (ctx )
9795
98- // if the number of workers is >1 the server would have to deal with out of order playlists
99- sink .Start (ctx , p .transport , reg , 1 )
96+ // initial config update
97+ if err := p .UpdateConfig (ctx , conf ); err != nil {
98+ return nil , err
10099 }
101100
102101 return p , nil
103102}
104103
104+ func (p * Proxy ) UpdateConfig (ctx context.Context , conf Config ) error {
105+ p .sinkMutex .Lock ()
106+ defer p .sinkMutex .Unlock ()
107+ var newConfigs []SinkConfig
108+ var err error
109+ // mark removed sinks for graceful deletion
110+ for _ , s := range p .sinks {
111+ var found bool
112+ for _ , sinkConfig := range conf .Sinks {
113+ if s .Address () == sinkConfig .Address {
114+ found = true
115+ break
116+ }
117+ }
118+ if found {
119+ continue
120+ }
121+ s .StartGracePeriod ()
122+ }
123+ // update existing sinks
124+ for _ , sinkConfig := range conf .Sinks {
125+ var updated bool
126+ for _ , s := range p .sinks {
127+ if s .Address () != sinkConfig .Address {
128+ continue
129+ }
130+ // update existing sink
131+ if err2 := s .UpdateConfig (ctx , sinkConfig ); err2 != nil {
132+ log .Error ().Err (err2 ).Str ("sink" , s .url .Host ).Msg ("sink update failed" )
133+ firstErr (err2 , & err )
134+ }
135+ updated = true
136+ break
137+ }
138+ if updated {
139+ continue
140+ }
141+ newConfigs = append (newConfigs , sinkConfig )
142+ }
143+ // create new sinks
144+ for _ , sinkConfig := range newConfigs {
145+ sink , err2 := NewSink (ctx , sinkConfig , p .transport , p .metrics .reg )
146+ if err2 != nil {
147+ firstErr (err2 , & err )
148+ log .Error ().Err (err2 ).Str ("sink" , sinkConfig .Address ).Msg ("sink init failed" )
149+ continue
150+ }
151+ p .sinks = append (p .sinks , sink )
152+ log .Info ().Str ("sink" , sink .url .Host ).Str ("basePath" , sink .url .Path ).Str ("authType" , string (sink .conf .AuthType )).Msg ("added sink" )
153+ }
154+ return err
155+ }
156+
157+ func (p * Proxy ) runTimeout (ctx context.Context ) {
158+ defer p .done .Done ()
159+ ticker := time .NewTicker (time .Minute )
160+ defer ticker .Stop ()
161+ for {
162+ select {
163+ case <- ctx .Done ():
164+ return
165+ case <- ticker .C :
166+ p .sinkMutex .Lock ()
167+ sinks := p .sinks [:0 ]
168+ for _ , s := range p .sinks {
169+ if s .IsStale () {
170+ s .Stop ()
171+ log .Info ().Str ("sink" , s .url .Host ).Msg ("removed sink after timeout" )
172+ continue
173+ }
174+ sinks = append (sinks , s )
175+ }
176+ p .sinks = sinks
177+ p .sinkMutex .Unlock ()
178+ }
179+ }
180+ }
181+
105182// Wait for server to finish
106183func (p * Proxy ) Wait () {
107184 p .done .Wait ()
108185 for _ , sink := range p .sinks {
109- sink .wait ()
186+ sink .Stop ()
110187 }
111188 p .metrics .deregister ()
112189}
@@ -235,3 +312,9 @@ func (m *ProxyMetrics) deregister() {
235312 m .reg .Unregister (m .totalNumRead )
236313 m .reg .Unregister (m .totalReadDelay )
237314}
315+
316+ func firstErr (err error , errPtr * error ) {
317+ if * errPtr == nil && err != nil {
318+ * errPtr = err
319+ }
320+ }
0 commit comments