blob: aa1a3d39c80ef4a453e43a1f411c740f4260c3fa (
plain)
| 1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
 | // Package bgreader provides a io.Reader that can optionally buffer reads in the background.
package bgreader
import (
	"io"
	"sync"
	"github.com/jackc/pgx/v5/internal/iobufpool"
)
const (
	bgReaderStatusStopped = iota
	bgReaderStatusRunning
	bgReaderStatusStopping
)
// BGReader is an io.Reader that can optionally buffer reads in the background. It is safe for concurrent use.
type BGReader struct {
	r io.Reader
	cond           *sync.Cond
	bgReaderStatus int32
	readResults    []readResult
}
type readResult struct {
	buf *[]byte
	err error
}
// Start starts the backgrounder reader. If the background reader is already running this is a no-op. The background
// reader will stop automatically when the underlying reader returns an error.
func (r *BGReader) Start() {
	r.cond.L.Lock()
	defer r.cond.L.Unlock()
	switch r.bgReaderStatus {
	case bgReaderStatusStopped:
		r.bgReaderStatus = bgReaderStatusRunning
		go r.bgRead()
	case bgReaderStatusRunning:
		// no-op
	case bgReaderStatusStopping:
		r.bgReaderStatus = bgReaderStatusRunning
	}
}
// Stop tells the background reader to stop after the in progress Read returns. It is safe to call Stop when the
// background reader is not running.
func (r *BGReader) Stop() {
	r.cond.L.Lock()
	defer r.cond.L.Unlock()
	switch r.bgReaderStatus {
	case bgReaderStatusStopped:
		// no-op
	case bgReaderStatusRunning:
		r.bgReaderStatus = bgReaderStatusStopping
	case bgReaderStatusStopping:
		// no-op
	}
}
func (r *BGReader) bgRead() {
	keepReading := true
	for keepReading {
		buf := iobufpool.Get(8192)
		n, err := r.r.Read(*buf)
		*buf = (*buf)[:n]
		r.cond.L.Lock()
		r.readResults = append(r.readResults, readResult{buf: buf, err: err})
		if r.bgReaderStatus == bgReaderStatusStopping || err != nil {
			r.bgReaderStatus = bgReaderStatusStopped
			keepReading = false
		}
		r.cond.L.Unlock()
		r.cond.Broadcast()
	}
}
// Read implements the io.Reader interface.
func (r *BGReader) Read(p []byte) (int, error) {
	r.cond.L.Lock()
	defer r.cond.L.Unlock()
	if len(r.readResults) > 0 {
		return r.readFromReadResults(p)
	}
	// There are no unread background read results and the background reader is stopped.
	if r.bgReaderStatus == bgReaderStatusStopped {
		return r.r.Read(p)
	}
	// Wait for results from the background reader
	for len(r.readResults) == 0 {
		r.cond.Wait()
	}
	return r.readFromReadResults(p)
}
// readBackgroundResults reads a result previously read by the background reader. r.cond.L must be held.
func (r *BGReader) readFromReadResults(p []byte) (int, error) {
	buf := r.readResults[0].buf
	var err error
	n := copy(p, *buf)
	if n == len(*buf) {
		err = r.readResults[0].err
		iobufpool.Put(buf)
		if len(r.readResults) == 1 {
			r.readResults = nil
		} else {
			r.readResults = r.readResults[1:]
		}
	} else {
		*buf = (*buf)[n:]
		r.readResults[0].buf = buf
	}
	return n, err
}
func New(r io.Reader) *BGReader {
	return &BGReader{
		r: r,
		cond: &sync.Cond{
			L: &sync.Mutex{},
		},
	}
}
 |