()
| 1027 | const maxBatchSize = 50 |
| 1028 | |
| 1029 | func (q *querier) peerUpdateWorker() { |
| 1030 | defer q.wg.Done() |
| 1031 | defer q.logger.Debug(q.ctx, "peerUpdate worker exited") |
| 1032 | eb := backoff.NewExponentialBackOff() |
| 1033 | eb.MaxElapsedTime = 0 // retry indefinitely |
| 1034 | eb.MaxInterval = dbMaxBackoff |
| 1035 | bkoff := backoff.WithContext(eb, q.ctx) |
| 1036 | for { |
| 1037 | allKeys, err := q.peerUpdateQ.acquireBatch(maxBatchSize) |
| 1038 | if err != nil { |
| 1039 | return |
| 1040 | } |
| 1041 | peers := make([]uuid.UUID, 0, len(allKeys)) |
| 1042 | peers = append(peers, allKeys...) |
| 1043 | err = backoff.Retry(func() error { |
| 1044 | return q.peerUpdate(peers) |
| 1045 | }, bkoff) |
| 1046 | if err != nil { |
| 1047 | bkoff.Reset() |
| 1048 | } |
| 1049 | q.peerUpdateQ.done(allKeys...) |
| 1050 | } |
| 1051 | } |
| 1052 | |
| 1053 | func (q *querier) mappingWorker() { |
| 1054 | defer q.wg.Done() |
no test coverage detected