protect memorybuffer race

This commit is contained in:
Daisuke Maki 2016-06-23 04:32:09 -04:00
parent f4ea32e0e3
commit 5d265d60e9

View file

@ -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")
}