| 67 | } |
| 68 | |
| 69 | func NewRequestQueue(forgetDelay time.Duration, queueLength *prometheus.GaugeVec, discardedRequests *prometheus.CounterVec, limits Limits, registerer prometheus.Registerer) *RequestQueue { |
| 70 | q := &RequestQueue{ |
| 71 | queues: newUserQueues(forgetDelay, limits, queueLength), |
| 72 | connectedQuerierWorkers: atomic.NewInt32(0), |
| 73 | totalRequests: promauto.With(registerer).NewCounterVec(prometheus.CounterOpts{ |
| 74 | Name: "cortex_request_queue_requests_total", |
| 75 | Help: "Total number of query requests going to the request queue.", |
| 76 | }, []string{"user", "priority"}), |
| 77 | discardedRequests: discardedRequests, |
| 78 | } |
| 79 | |
| 80 | q.cond = sync.NewCond(&q.mtx) |
| 81 | q.Service = services.NewTimerService(forgetCheckPeriod, nil, q.forgetDisconnectedQueriers, q.stopping).WithName("request queue") |
| 82 | |
| 83 | return q |
| 84 | } |
| 85 | |
| 86 | // EnqueueRequest puts the request into the queue. MaxQueries is user-specific value that specifies how many queriers can |
| 87 | // this user use (zero or negative = all queriers). It is passed to each EnqueueRequest, because it can change |