reflector.go/peer/http3/worker.go

48 lines
932 B
Go
Raw Normal View History

2021-04-06 20:00:36 +02:00
package http3
import (
"net/http"
2021-04-06 20:28:29 +02:00
"sync"
2021-04-06 20:00:36 +02:00
"github.com/lbryio/reflector.go/internal/metrics"
"github.com/lbryio/lbry.go/v2/extras/stop"
)
type blobRequest struct {
request *http.Request
reply http.ResponseWriter
2021-04-06 20:28:29 +02:00
finished *sync.WaitGroup
2021-04-06 20:00:36 +02:00
}
var getReqCh = make(chan *blobRequest, 20000)
2021-04-06 20:00:36 +02:00
func InitWorkers(server *Server, workers int) {
2021-04-06 20:00:36 +02:00
stopper := stop.New(server.grp)
for i := 0; i < workers; i++ {
2021-05-21 00:12:30 +02:00
metrics.RoutinesQueue.WithLabelValues("http3", "worker").Inc()
2021-04-06 20:00:36 +02:00
go func(worker int) {
2021-05-21 00:12:30 +02:00
defer metrics.RoutinesQueue.WithLabelValues("http3", "worker").Dec()
for {
select {
case <-stopper.Ch():
case r := <-getReqCh:
metrics.Http3BlobReqQueue.Dec()
process(server, r)
}
2021-04-06 20:00:36 +02:00
}
}(i)
}
return
2021-04-06 20:00:36 +02:00
}
func enqueue(b *blobRequest) {
metrics.Http3BlobReqQueue.Inc()
getReqCh <- b
}
func process(server *Server, r *blobRequest) {
server.HandleGetBlob(r.reply, r.request)
r.finished.Done()
2021-04-06 20:00:36 +02:00
}