From 5d265d60e911a4d82c57d9792dabf5df0b783be3 Mon Sep 17 00:00:00 2001 From: Daisuke Maki Date: Thu, 23 Jun 2016 04:32:09 -0400 Subject: [PATCH] protect memorybuffer race --- buffer.go | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/buffer.go b/buffer.go index 4683d09..1b47c84 100644 --- a/buffer.go +++ b/buffer.go @@ -62,19 +62,27 @@ func NewMemoryBuffer() *MemoryBuffer { } func (mb *MemoryBuffer) Size() int { + mb.mutex.Lock() + defer mb.mutex.Unlock() return len(mb.lines) } func (mb *MemoryBuffer) Reset() { + mb.mutex.Lock() + defer mb.mutex.Unlock() mb.done = make(chan struct{}) mb.lines = []Line(nil) } func (mb *MemoryBuffer) Done() <-chan struct{} { + mb.mutex.Lock() + defer mb.mutex.Unlock() return mb.done } func (mb *MemoryBuffer) Accept(ctx context.Context, p pipeline.Producer) { + mb.mutex.Lock() + defer mb.mutex.Unlock() if pdebug.Enabled { g := pdebug.Marker("MemoryBuffer.Accept") defer g.End() @@ -93,7 +101,7 @@ func (mb *MemoryBuffer) Accept(ctx context.Context, p pipeline.Producer) { case error: if pipeline.IsEndMark(v.(error)) { if pdebug.Enabled { - pdebug.Printf("MemoryBuffer received end mark (read %d lines)", mb.Size()) + pdebug.Printf("MemoryBuffer received end mark (read %d lines)", len(mb.lines)) } return } @@ -108,7 +116,7 @@ func (mb *MemoryBuffer) LineAt(n int) (Line, error) { mb.mutex.Lock() defer mb.mutex.Unlock() - if s := mb.Size(); s <= 0 || n >= s { + if s := len(mb.lines); s <= 0 || n >= s { return nil, errors.New("empty buffer") }