Finds next queue for the querier. To support fair scheduling between users, client is expected to pass last user index returned by this function as argument. Is there was no previous last user index, use -1.
(lastUserIndex int, querierID string)
| 221 | // to pass last user index returned by this function as argument. Is there was no previous |
| 222 | // last user index, use -1. |
| 223 | func (q *queues) getNextQueueForQuerier(lastUserIndex int, querierID string) (userRequestQueue, string, int) { |
| 224 | uid := lastUserIndex |
| 225 | |
| 226 | q.queuesMx.RLock() |
| 227 | defer q.queuesMx.RUnlock() |
| 228 | |
| 229 | for iters := 0; iters < len(q.users); iters++ { |
| 230 | uid = uid + 1 |
| 231 | |
| 232 | // Don't use "mod len(q.users)", as that could skip users at the beginning of the list |
| 233 | // for example when q.users has shrunk since last call. |
| 234 | if uid >= len(q.users) { |
| 235 | uid = 0 |
| 236 | } |
| 237 | |
| 238 | u := q.users[uid] |
| 239 | if u == "" { |
| 240 | continue |
| 241 | } |
| 242 | |
| 243 | uq := q.userQueues[u] |
| 244 | |
| 245 | if uq.queriers != nil { |
| 246 | if _, ok := uq.queriers[querierID]; !ok { |
| 247 | // This querier is not handling the user. |
| 248 | continue |
| 249 | } |
| 250 | } |
| 251 | |
| 252 | return uq.queue, u, uid |
| 253 | } |
| 254 | return nil, "", uid |
| 255 | } |
| 256 | |
| 257 | func (q *queues) addQuerierConnection(querierID string) { |
| 258 | info := q.queriers[querierID] |
no outgoing calls