Repository navigation
Expand file tree
/
Copy pathflow_controller_connection.go
More file actions
178 lines (152 loc) · 5.3 KB
/
Copy pathflow_controller_connection.go
File metadata and controls
178 lines (152 loc) · 5.3 KB
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
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
package quic
import (
"fmt"
"sync"
"github.com/quic-go/quic-go/internal/monotime"
"github.com/quic-go/quic-go/internal/protocol"
"github.com/quic-go/quic-go/internal/qerr"
"github.com/quic-go/quic-go/internal/utils"
)
type connectionFlowController struct {
receiveFlowController
// Protects send-side state, which TryWriteAll can access from application goroutines.
sendMutex sync.Mutex
bytesSent protocol.ByteCount
sendWindow protocol.ByteCount
lastBlockedAt protocol.ByteCount
}
// newConnectionFlowController gets a new flow controller for the connection.
// It is created before we receive the peer's transport parameters, thus it starts with a sendWindow of 0.
func newConnectionFlowController(
receiveWindow protocol.ByteCount,
maxReceiveWindow protocol.ByteCount,
allowWindowIncrease func(size protocol.ByteCount) bool,
rttStats *utils.RTTStats,
logger utils.Logger,
) *connectionFlowController {
return &connectionFlowController{
receiveFlowController: receiveFlowController{
rttStats: rttStats,
receiveWindow: receiveWindow,
receiveWindowSize: receiveWindow,
maxReceiveWindowSize: maxReceiveWindow,
allowWindowIncrease: allowWindowIncrease,
logger: logger,
},
}
}
// IncrementHighestReceived adds an increment to the highestReceived value
func (c *connectionFlowController) IncrementHighestReceived(increment protocol.ByteCount, now monotime.Time) error {
c.mutex.Lock()
defer c.mutex.Unlock()
// If this is the first frame received on this connection, start flow-control auto-tuning.
if c.highestReceived == 0 {
c.startNewAutoTuningEpoch(now)
}
c.highestReceived += increment
if c.checkFlowControlViolation() {
return &qerr.TransportError{
ErrorCode: qerr.FlowControlError,
ErrorMessage: fmt.Sprintf("received %d bytes for the connection, allowed %d bytes", c.highestReceived, c.receiveWindow),
}
}
return nil
}
func (c *connectionFlowController) AddBytesRead(n protocol.ByteCount) (hasWindowUpdate bool) {
c.mutex.Lock()
defer c.mutex.Unlock()
c.addBytesRead(n)
return c.hasWindowUpdate()
}
// TryAddBytesSent adds n bytes if sufficient connection-level send credit is available.
func (c *connectionFlowController) TryAddBytesSent(n protocol.ByteCount) bool {
c.sendMutex.Lock()
defer c.sendMutex.Unlock()
if c.bytesSent > c.sendWindow || n > c.sendWindow-c.bytesSent {
return false
}
c.bytesSent += n
return true
}
// AddBytesSentWithLimiter adds the limiter-approved portion of the available connection-level send credit.
func (c *connectionFlowController) AddBytesSentWithLimiter(
n protocol.ByteCount,
limiter func(int) int,
) (protocol.ByteCount, bool) {
c.sendMutex.Lock()
defer c.sendMutex.Unlock()
if c.bytesSent >= c.sendWindow {
return 0, false
}
n = min(n, c.sendWindow-c.bytesSent)
added := min(
max(protocol.ByteCount(limiter(int(n))), 0),
n,
)
c.bytesSent += added
return added, added < n
}
// UpdateSendWindow is called after receiving a MAX_DATA frame.
func (c *connectionFlowController) UpdateSendWindow(offset protocol.ByteCount) (updated bool) {
c.sendMutex.Lock()
defer c.sendMutex.Unlock()
if offset > c.sendWindow {
c.sendWindow = offset
return true
}
return false
}
func (c *connectionFlowController) SendWindowSize() protocol.ByteCount {
c.sendMutex.Lock()
defer c.sendMutex.Unlock()
return c.sendWindow - c.bytesSent
}
// IsNewlyBlocked says if it is newly blocked by connection flow control.
// For every offset, it only returns true once.
// If it is blocked, the offset is returned.
func (c *connectionFlowController) IsNewlyBlocked() (bool, protocol.ByteCount) {
c.sendMutex.Lock()
defer c.sendMutex.Unlock()
if c.bytesSent < c.sendWindow || c.sendWindow == c.lastBlockedAt {
return false, 0
}
c.lastBlockedAt = c.sendWindow
return true, c.sendWindow
}
func (c *connectionFlowController) GetWindowUpdate(now monotime.Time) protocol.ByteCount {
c.mutex.Lock()
defer c.mutex.Unlock()
oldWindowSize := c.receiveWindowSize
offset := c.getWindowUpdate(now)
if c.logger.Debug() && oldWindowSize < c.receiveWindowSize {
c.logger.Debugf("Increasing receive flow control window for the connection to %d kB", c.receiveWindowSize/(1<<10))
}
return offset
}
// EnsureMinimumWindowSize sets a minimum window size
// it should make sure that the connection-level window is increased when a stream-level window grows
func (c *connectionFlowController) EnsureMinimumWindowSize(inc protocol.ByteCount, now monotime.Time) {
c.mutex.Lock()
defer c.mutex.Unlock()
if inc <= c.receiveWindowSize {
return
}
newSize := min(inc, c.maxReceiveWindowSize)
if delta := newSize - c.receiveWindowSize; delta > 0 && c.allowWindowIncrease(delta) {
c.receiveWindowSize = newSize
if c.logger.Debug() {
c.logger.Debugf("Increasing receive flow control window for the connection to %d, in response to stream flow control window increase", newSize)
}
}
c.startNewAutoTuningEpoch(now)
}
// Reset rests the flow controller. This happens when 0-RTT is rejected.
// All stream data is invalidated, it's as if we had never opened a stream and never sent any data.
// At that point, we only have sent stream data, but we didn't have the keys to open 1-RTT keys yet.
func (c *connectionFlowController) Reset() {
c.sendMutex.Lock()
defer c.sendMutex.Unlock()
c.bytesSent = 0
c.lastBlockedAt = 0
c.sendWindow = 0
}