Skip to content

Commit 2811703

Browse files
authored
Fix outbound stream starvation (#559)
Motivation: `nextStreamToSend()` returned `self.flushableStreams.first`, and since `Set.first` always yields the same stable element, under load the buffer kept handing the whole connection window to one stream while the rest starved to near-zero. Modifiations: Rotate through the flushable streams instead of always picking the same one. `flushableStreams` stays the source of truth for membership, and alongside it we keep a queue of those streams (`sendQueue`) in the order to serve them. We pop the front to pick the next stream, and append it to the back once it's been served, so every flushable stream gets a turn per sweep and none can be starved. Both the pop and the append are O(1). Removing a stream from the middle of the queue would be O(n), so we avoid this. A stream that stops being flushable is left as a stale entry and simply skipped when it reaches the front (lazy deletion), with a `queuedStreams` set to stop us ever enqueuing a duplicate.
1 parent 45bdf67 commit 2811703

5 files changed

Lines changed: 79 additions & 32 deletions

File tree

IntegrationTests/tests_01_allocation_counters/Thresholds/6.1.json

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
"1k_requests_inline_noninterleaved": 29100,
44
"1k_requests_interleaved": 36150,
55
"1k_requests_noninterleaved": 35100,
6-
"client_server_h1_request_response": 275050,
7-
"client_server_h1_request_response_inline": 263050,
8-
"client_server_request_response": 253050,
9-
"client_server_request_response_inline": 244050,
10-
"client_server_request_response_many": 1198050,
11-
"client_server_request_response_many_inline": 889050,
6+
"client_server_h1_request_response": 279050,
7+
"client_server_h1_request_response_inline": 267050,
8+
"client_server_request_response": 257050,
9+
"client_server_request_response_inline": 248050,
10+
"client_server_request_response_many": 1202050,
11+
"client_server_request_response_many_inline": 893050,
1212
"create_client_stream_channel": 35050,
1313
"create_client_stream_channel_inline": 35050,
1414
"create_client_stream_channel_inline_no_promise_based_API": 35050,
@@ -18,6 +18,6 @@
1818
"get_100000_headers_canonical_form_trimming_whitespace_from_long_string": 300050,
1919
"get_100000_headers_canonical_form_trimming_whitespace_from_short_string": 200050,
2020
"hpack_decoding": 5050,
21-
"stream_teardown_100_concurrent": 253250,
22-
"stream_teardown_100_concurrent_inline": 252350
21+
"stream_teardown_100_concurrent": 253450,
22+
"stream_teardown_100_concurrent_inline": 252550
2323
}

IntegrationTests/tests_01_allocation_counters/Thresholds/6.2.json

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
"1k_requests_inline_noninterleaved": 29100,
44
"1k_requests_interleaved": 36150,
55
"1k_requests_noninterleaved": 35100,
6-
"client_server_h1_request_response": 275050,
7-
"client_server_h1_request_response_inline": 263050,
8-
"client_server_request_response": 253050,
9-
"client_server_request_response_inline": 244050,
10-
"client_server_request_response_many": 1198050,
11-
"client_server_request_response_many_inline": 889050,
6+
"client_server_h1_request_response": 279050,
7+
"client_server_h1_request_response_inline": 267050,
8+
"client_server_request_response": 257050,
9+
"client_server_request_response_inline": 248050,
10+
"client_server_request_response_many": 1202050,
11+
"client_server_request_response_many_inline": 893050,
1212
"create_client_stream_channel": 35050,
1313
"create_client_stream_channel_inline": 35050,
1414
"create_client_stream_channel_inline_no_promise_based_API": 35050,
@@ -18,6 +18,6 @@
1818
"get_100000_headers_canonical_form_trimming_whitespace_from_long_string": 300050,
1919
"get_100000_headers_canonical_form_trimming_whitespace_from_short_string": 200050,
2020
"hpack_decoding": 5050,
21-
"stream_teardown_100_concurrent": 253250,
22-
"stream_teardown_100_concurrent_inline": 252350
21+
"stream_teardown_100_concurrent": 253450,
22+
"stream_teardown_100_concurrent_inline": 252550
2323
}

IntegrationTests/tests_01_allocation_counters/Thresholds/6.3.json

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
"1k_requests_inline_noninterleaved": 29100,
44
"1k_requests_interleaved": 36150,
55
"1k_requests_noninterleaved": 35100,
6-
"client_server_h1_request_response": 275050,
7-
"client_server_h1_request_response_inline": 263050,
8-
"client_server_request_response": 253050,
9-
"client_server_request_response_inline": 244050,
10-
"client_server_request_response_many": 1198050,
11-
"client_server_request_response_many_inline": 889050,
6+
"client_server_h1_request_response": 279050,
7+
"client_server_h1_request_response_inline": 267050,
8+
"client_server_request_response": 257050,
9+
"client_server_request_response_inline": 248050,
10+
"client_server_request_response_many": 1202050,
11+
"client_server_request_response_many_inline": 893050,
1212
"create_client_stream_channel": 35050,
1313
"create_client_stream_channel_inline": 35050,
1414
"create_client_stream_channel_inline_no_promise_based_API": 35050,
@@ -18,6 +18,6 @@
1818
"get_100000_headers_canonical_form_trimming_whitespace_from_long_string": 300050,
1919
"get_100000_headers_canonical_form_trimming_whitespace_from_short_string": 200050,
2020
"hpack_decoding": 5050,
21-
"stream_teardown_100_concurrent": 253250,
22-
"stream_teardown_100_concurrent_inline": 252350
21+
"stream_teardown_100_concurrent": 253450,
22+
"stream_teardown_100_concurrent_inline": 252550
2323
}

Sources/NIOHTTP2/Frame Buffers/OutboundFlowControlBuffer.swift

Lines changed: 39 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,12 @@ internal struct OutboundFlowControlBuffer {
5555
/// The streams with pending data to output.
5656
private var flushableStreams: Set<HTTP2StreamID> = Set()
5757

58+
/// Round-robin order over flushableStreams. We rotate through this to give every stream a fair turn without starvation.
59+
/// Streams that stop being flushable are left as stale entries and skipped when they surface, giving O(1) lazy deletion.
60+
private var sendQueue: CircularBuffer<HTTP2StreamID> = CircularBuffer()
61+
/// Stops us enqueuing a duplicate.
62+
private var queuedStreams: Set<HTTP2StreamID> = Set()
63+
5864
/// The current size of the connection flow control window. May be negative.
5965
internal var connectionWindowSize: Int
6066

@@ -70,6 +76,8 @@ internal struct OutboundFlowControlBuffer {
7076
// Avoid some resizes.
7177
self.writableStreams.reserveCapacity(16)
7278
self.flushableStreams.reserveCapacity(16)
79+
self.sendQueue.reserveCapacity(16)
80+
self.queuedStreams.reserveCapacity(16)
7381
}
7482

7583
internal mutating func processOutboundFrame(
@@ -132,14 +140,25 @@ internal struct OutboundFlowControlBuffer {
132140
}
133141
if let actuallyWritable = actuallyWritable, actuallyWritable {
134142
self.flushableStreams.insert(streamID)
143+
if self.queuedStreams.insert(streamID).inserted {
144+
self.sendQueue.append(streamID)
145+
}
135146
}
136147
}
137148

138149
self.writableStreams.removeAll(keepingCapacity: true)
139150
}
140151

141-
private func nextStreamToSend() -> HTTP2StreamID? {
142-
self.flushableStreams.first
152+
private mutating func nextStreamToSend() -> HTTP2StreamID? {
153+
// Round-robin via the FIFO queue. Pop from the front, skipping streams that are no longer flushable.
154+
// Each stale entry is discarded at most once, so this is O(1) amortised.
155+
while let streamID = self.sendQueue.popFirst() {
156+
self.queuedStreams.remove(streamID)
157+
if self.flushableStreams.contains(streamID) {
158+
return streamID
159+
}
160+
}
161+
return nil
143162
}
144163

145164
internal mutating func updateWindowOfStream(_ streamID: HTTP2StreamID, newSize: Int32) {
@@ -153,6 +172,9 @@ internal struct OutboundFlowControlBuffer {
153172
case .changed(newValue: true):
154173
// Became writable, and specifically became _flushable_.
155174
self.flushableStreams.insert(streamID)
175+
if self.queuedStreams.insert(streamID).inserted {
176+
self.sendQueue.append(streamID)
177+
}
156178
case .changed(newValue: false):
157179
// Became unwritable.
158180
self.flushableStreams.remove(streamID)
@@ -196,31 +218,37 @@ internal struct OutboundFlowControlBuffer {
196218

197219
internal mutating func nextFlushedWritableFrame() -> (HTTP2Frame, EventLoopPromise<Void>?)? {
198220
// If the channel isn't writable, we don't want to send anything.
199-
guard let nextStreamID = self.nextStreamToSend(), self.connectionWindowSize > 0 else {
221+
guard self.connectionWindowSize > 0, let nextStreamID = self.nextStreamToSend() else {
200222
return nil
201223
}
202224

203225
let nextWrite = self.streamDataBuffers.modify(streamID: nextStreamID) {
204-
(state: inout StreamFlowControlState) -> DataBuffer.BufferElement in
226+
(state: inout StreamFlowControlState) -> (DataBuffer.BufferElement, isFlushable: Bool) in
205227
let (nextWrite, writabilityState) = state.nextWrite(
206228
maxSize: min(self.connectionWindowSize, self.maxFrameSize)
207229
)
208230

209231
switch writabilityState {
210232
case .changed(newValue: false):
211233
self.flushableStreams.remove(nextStreamID)
234+
return (nextWrite, isFlushable: false)
212235
case .changed(newValue: true), .unchanged:
213-
()
236+
return (nextWrite, isFlushable: true)
214237
}
215-
216-
return nextWrite
217238
}
218-
guard let (payload, promise) = nextWrite else {
239+
guard let ((payload, promise), isFlushable) = nextWrite else {
219240
// The stream was not present. This is weird, it shouldn't ever happen, but we tolerate it, and recurse.
220241
self.flushableStreams.remove(nextStreamID)
221242
return self.nextFlushedWritableFrame()
222243
}
223244

245+
// If the stream is still flushable, put it back at the end of the rotation.
246+
if isFlushable {
247+
let inserted = self.queuedStreams.insert(nextStreamID).inserted
248+
assert(inserted, "\(nextStreamID) unexpectedly still present in queuedStreams")
249+
self.sendQueue.append(nextStreamID)
250+
}
251+
224252
let frame = HTTP2Frame(streamID: nextStreamID, payload: payload)
225253
return (frame, promise)
226254
}
@@ -234,6 +262,9 @@ internal struct OutboundFlowControlBuffer {
234262
case .changed(newValue: true):
235263
// Became flushable
236264
self.flushableStreams.insert($0.streamID)
265+
if self.queuedStreams.insert($0.streamID).inserted {
266+
self.sendQueue.append($0.streamID)
267+
}
237268
case .changed(newValue: false):
238269
// Became unflushable.
239270
self.flushableStreams.remove($0.streamID)

Tests/NIOHTTP2Tests/OutboundFlowControlBufferTests.swift

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -393,6 +393,22 @@ class OutboundFlowControlBufferTests: XCTestCase {
393393
XCTAssertNil(self.buffer.nextFlushedWritableFrame())
394394
}
395395

396+
func testFlushableStreamsAreNotStarved() {
397+
// Three streams with three frames each.
398+
let streamIDs: [HTTP2StreamID] = [1, 3, 5]
399+
self.buffer.maxFrameSize = 5
400+
401+
for streamID in streamIDs {
402+
self.buffer.streamCreated(streamID, initialWindowSize: 15)
403+
let frame = self.createDataFrame(streamID, byteBufferSize: 15)
404+
XCTAssertNoThrow(try self.buffer.processOutboundFrame(frame, promise: nil).assertNothing())
405+
self.buffer.flushReceived()
406+
}
407+
408+
// Every stream must be served once per round so that none can be starved.
409+
XCTAssertEqual(self.receivedFrames().map { $0.streamID }, [1, 3, 5, 1, 3, 5, 1, 3, 5])
410+
}
411+
396412
func testRejectsPrioritySelfDependency() {
397413
XCTAssertThrowsError(
398414
try self.buffer.priorityUpdate(

0 commit comments

Comments
 (0)