Next returns the next chunk from the queue, or errDone if all chunks have been returned. It blocks until the chunk is available. Concurrent Next() calls may return the same chunk.
()
| 224 | // Next returns the next chunk from the queue, or errDone if all chunks have been returned. It |
| 225 | // blocks until the chunk is available. Concurrent Next() calls may return the same chunk. |
| 226 | func (q *chunkQueue) Next() (*chunk, error) { |
| 227 | q.Lock() |
| 228 | var chunk *chunk |
| 229 | index, err := q.nextUp() |
| 230 | if err == nil { |
| 231 | chunk, err = q.load(index) |
| 232 | if err == nil { |
| 233 | q.chunkReturned[index] = true |
| 234 | } |
| 235 | } |
| 236 | q.Unlock() |
| 237 | if chunk != nil || err != nil { |
| 238 | return chunk, err |
| 239 | } |
| 240 | |
| 241 | select { |
| 242 | case _, ok := <-q.WaitFor(index): |
| 243 | if !ok { |
| 244 | return nil, errDone // queue closed |
| 245 | } |
| 246 | case <-time.After(chunkTimeout): |
| 247 | return nil, errTimeout |
| 248 | } |
| 249 | |
| 250 | q.Lock() |
| 251 | defer q.Unlock() |
| 252 | chunk, err = q.load(index) |
| 253 | if err != nil { |
| 254 | return nil, err |
| 255 | } |
| 256 | q.chunkReturned[index] = true |
| 257 | return chunk, nil |
| 258 | } |
| 259 | |
| 260 | // nextUp returns the next chunk to be returned, or errDone if all chunks have been returned. The |
| 261 | // caller must hold the mutex lock. |