Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion pkg/buffer/pipe.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ type PipeBuffer struct {
offR int
offW int
rw sync.Mutex
ioWg sync.WaitGroup
block Block

readSignal chan struct{}
Expand Down Expand Up @@ -73,6 +74,8 @@ func (br *PipeBuffer) Read(p []byte) (int, error) {

off := br.offR
block := br.block
br.ioWg.Add(1)
defer br.ioWg.Done()
br.rw.Unlock()

n, err := block.ReadAt(p[:min(len(p), canRead)], int64(off))
Expand Down Expand Up @@ -109,6 +112,8 @@ func (br *PipeBuffer) Write(p []byte) (int, error) {

off := br.offW
block := br.block
br.ioWg.Add(1)
defer br.ioWg.Done()
br.rw.Unlock()

n, err := block.WriteAt(p[:min(canWrite, len(p))], int64(off))
Expand All @@ -131,6 +136,7 @@ func (br *PipeBuffer) Write(p []byte) (int, error) {
}

func (br *PipeBuffer) Reset(limit int) error {
br.ioWg.Wait()
br.rw.Lock()
defer br.rw.Unlock()
if br.block == nil {
Expand All @@ -147,11 +153,12 @@ func (br *PipeBuffer) Reset(limit int) error {

func (br *PipeBuffer) Close() error {
br.rw.Lock()
defer br.rw.Unlock()
if br.block != nil {
br.block = nil
br.readPending = false
close(br.readSignal)
}
br.rw.Unlock()
br.ioWg.Wait()
return nil
}
131 changes: 131 additions & 0 deletions pkg/buffer/pipe_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,131 @@
package buffer

import (
"context"
"errors"
"io"
"testing"
"time"
)

type blockingBlock struct {
data []byte
blockOn string
started chan struct{}
release chan struct{}
}

func newBlockingBlock(blockOn string) *blockingBlock {
return &blockingBlock{
data: make([]byte, 1),
blockOn: blockOn,
started: make(chan struct{}),
release: make(chan struct{}),
}
}

func (b *blockingBlock) Size() int64 {
return int64(len(b.data))
}

func (b *blockingBlock) ReadAt(p []byte, off int64) (int, error) {
if b.blockOn == "read" {
close(b.started)
<-b.release
}
n := copy(p, b.data[off:])
if n < len(p) {
return n, io.EOF
}
return n, nil
}

func (b *blockingBlock) WriteAt(p []byte, off int64) (int, error) {
if b.blockOn == "write" {
close(b.started)
<-b.release
}
n := copy(b.data[off:], p)
if n < len(p) {
return n, io.ErrShortWrite
}
return n, nil
}

func TestPipeBufferCloseWaitsForActiveIO(t *testing.T) {
for _, operation := range []string{"read", "write"} {
t.Run(operation, func(t *testing.T) {
block := newBlockingBlock(operation)
buf := NewPipeBuffer(context.Background(), block)
if operation == "read" {
if _, err := buf.Write([]byte{1}); err != nil {
t.Fatalf("prepare read: %v", err)
}
}

ioDone := make(chan error, 1)
go func() {
var err error
if operation == "read" {
_, err = buf.Read(make([]byte, 1))
} else {
_, err = buf.Write([]byte{1})
}
ioDone <- err
}()

select {
case <-block.started:
case <-time.After(time.Second):
t.Fatal("I/O did not start")
}

closeDone := make(chan error, 1)
go func() {
closeDone <- buf.Close()
}()

deadline := time.Now().Add(time.Second)
for {
buf.rw.Lock()
closed := buf.block == nil
buf.rw.Unlock()
if closed {
break
}
if time.Now().After(deadline) {
t.Fatal("buffer did not enter the closed state")
}
time.Sleep(time.Millisecond)
}

select {
case err := <-closeDone:
t.Fatalf("Close returned before active %s completed: %v", operation, err)
default:
}

close(block.release)
select {
case err := <-ioDone:
if err != nil {
t.Fatalf("active %s failed: %v", operation, err)
}
case <-time.After(time.Second):
t.Fatalf("active %s did not complete", operation)
}
select {
case err := <-closeDone:
if err != nil {
t.Fatalf("Close failed: %v", err)
}
case <-time.After(time.Second):
t.Fatal("Close did not wait for active I/O")
}

if _, err := buf.Write([]byte{1}); !errors.Is(err, io.ErrClosedPipe) {
t.Fatalf("write after Close error = %v, want %v", err, io.ErrClosedPipe)
}
})
}
}