diff --git a/queue/ring.go b/queue/ring.go index ed1de3c9..2c8aee95 100644 --- a/queue/ring.go +++ b/queue/ring.go @@ -93,6 +93,14 @@ L: return false, ErrDisposed } + if offer { + deq := atomic.LoadUint64(&rb.dequeue) + pos = atomic.LoadUint64(&rb.queue) + if pos-deq >= rb.Cap() { + return false, nil + } + } + n = &rb.nodes[pos&rb.mask] seq := atomic.LoadUint64(&n.position) switch dif := seq - pos; { @@ -100,16 +108,13 @@ L: if atomic.CompareAndSwapUint64(&rb.queue, pos, pos+1) { break L } + pos = atomic.LoadUint64(&rb.queue) case dif < 0: panic(`Ring buffer in a compromised state during a put operation.`) default: pos = atomic.LoadUint64(&rb.queue) } - if offer { - return false, nil - } - runtime.Gosched() // free up the cpu before the next iteration } diff --git a/queue/ring_test.go b/queue/ring_test.go index 21b020c0..794c70ad 100644 --- a/queue/ring_test.go +++ b/queue/ring_test.go @@ -150,6 +150,37 @@ func TestOffer(t *testing.T) { assert.Equal(t, "bar", item) } +func TestRingQueueOffer_parallel(t *testing.T) { + size := 256 + parallelGoroutines := 8 + + rb := NewRingBuffer(uint64(size * parallelGoroutines)) + + wg := new(sync.WaitGroup) + wg.Add(parallelGoroutines) + + for i := 0; i < parallelGoroutines; i++ { + go func(id int) { + defer wg.Done() + + for el := 1; el <= size; el++ { + ok, err := rb.Offer(el) + if err != nil { + t.Errorf("error in goroutine-%d: %v", id, err) + return + } + + if !ok { + t.Errorf("queue full before expected on adding %d element, len: %d, cap: %d", el, rb.Len(), rb.Cap()) + } + } + }(i) + } + + wg.Wait() +} + + func TestRingGetEmpty(t *testing.T) { rb := NewRingBuffer(3)