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
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,12 @@
"1k_requests_inline_noninterleaved": 29100,
"1k_requests_interleaved": 36150,
"1k_requests_noninterleaved": 35100,
"client_server_h1_request_response": 275050,
"client_server_h1_request_response_inline": 263050,
"client_server_request_response": 253050,
"client_server_request_response_inline": 244050,
"client_server_request_response_many": 1198050,
"client_server_request_response_many_inline": 889050,
"client_server_h1_request_response": 279050,
"client_server_h1_request_response_inline": 267050,
"client_server_request_response": 257050,
"client_server_request_response_inline": 248050,
"client_server_request_response_many": 1202050,
"client_server_request_response_many_inline": 893050,
"create_client_stream_channel": 35050,
"create_client_stream_channel_inline": 35050,
"create_client_stream_channel_inline_no_promise_based_API": 35050,
Expand All @@ -18,6 +18,6 @@
"get_100000_headers_canonical_form_trimming_whitespace_from_long_string": 300050,
"get_100000_headers_canonical_form_trimming_whitespace_from_short_string": 200050,
"hpack_decoding": 5050,
"stream_teardown_100_concurrent": 253250,
"stream_teardown_100_concurrent_inline": 252350
"stream_teardown_100_concurrent": 253450,
"stream_teardown_100_concurrent_inline": 252550
}
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,12 @@
"1k_requests_inline_noninterleaved": 29100,
"1k_requests_interleaved": 36150,
"1k_requests_noninterleaved": 35100,
"client_server_h1_request_response": 275050,
"client_server_h1_request_response_inline": 263050,
"client_server_request_response": 253050,
"client_server_request_response_inline": 244050,
"client_server_request_response_many": 1198050,
"client_server_request_response_many_inline": 889050,
"client_server_h1_request_response": 279050,
"client_server_h1_request_response_inline": 267050,
"client_server_request_response": 257050,
"client_server_request_response_inline": 248050,
"client_server_request_response_many": 1202050,
"client_server_request_response_many_inline": 893050,
"create_client_stream_channel": 35050,
"create_client_stream_channel_inline": 35050,
"create_client_stream_channel_inline_no_promise_based_API": 35050,
Expand All @@ -18,6 +18,6 @@
"get_100000_headers_canonical_form_trimming_whitespace_from_long_string": 300050,
"get_100000_headers_canonical_form_trimming_whitespace_from_short_string": 200050,
"hpack_decoding": 5050,
"stream_teardown_100_concurrent": 253250,
"stream_teardown_100_concurrent_inline": 252350
"stream_teardown_100_concurrent": 253450,
"stream_teardown_100_concurrent_inline": 252550
}
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,12 @@
"1k_requests_inline_noninterleaved": 29100,
"1k_requests_interleaved": 36150,
"1k_requests_noninterleaved": 35100,
"client_server_h1_request_response": 275050,
"client_server_h1_request_response_inline": 263050,
"client_server_request_response": 253050,
"client_server_request_response_inline": 244050,
"client_server_request_response_many": 1198050,
"client_server_request_response_many_inline": 889050,
"client_server_h1_request_response": 279050,
"client_server_h1_request_response_inline": 267050,
"client_server_request_response": 257050,
"client_server_request_response_inline": 248050,
"client_server_request_response_many": 1202050,
"client_server_request_response_many_inline": 893050,
"create_client_stream_channel": 35050,
"create_client_stream_channel_inline": 35050,
"create_client_stream_channel_inline_no_promise_based_API": 35050,
Expand All @@ -18,6 +18,6 @@
"get_100000_headers_canonical_form_trimming_whitespace_from_long_string": 300050,
"get_100000_headers_canonical_form_trimming_whitespace_from_short_string": 200050,
"hpack_decoding": 5050,
"stream_teardown_100_concurrent": 253250,
"stream_teardown_100_concurrent_inline": 252350
"stream_teardown_100_concurrent": 253450,
"stream_teardown_100_concurrent_inline": 252550
}
47 changes: 39 additions & 8 deletions Sources/NIOHTTP2/Frame Buffers/OutboundFlowControlBuffer.swift
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,12 @@ internal struct OutboundFlowControlBuffer {
/// The streams with pending data to output.
private var flushableStreams: Set<HTTP2StreamID> = Set()

/// Round-robin order over flushableStreams. We rotate through this to give every stream a fair turn without starvation.
/// Streams that stop being flushable are left as stale entries and skipped when they surface, giving O(1) lazy deletion.
private var sendQueue: CircularBuffer<HTTP2StreamID> = CircularBuffer()
/// Stops us enqueuing a duplicate.
private var queuedStreams: Set<HTTP2StreamID> = Set()

/// The current size of the connection flow control window. May be negative.
internal var connectionWindowSize: Int

Expand All @@ -70,6 +76,8 @@ internal struct OutboundFlowControlBuffer {
// Avoid some resizes.
self.writableStreams.reserveCapacity(16)
self.flushableStreams.reserveCapacity(16)
self.sendQueue.reserveCapacity(16)
self.queuedStreams.reserveCapacity(16)
}

internal mutating func processOutboundFrame(
Expand Down Expand Up @@ -132,14 +140,25 @@ internal struct OutboundFlowControlBuffer {
}
if let actuallyWritable = actuallyWritable, actuallyWritable {
self.flushableStreams.insert(streamID)
if self.queuedStreams.insert(streamID).inserted {
self.sendQueue.append(streamID)
}
}
}

self.writableStreams.removeAll(keepingCapacity: true)
}

private func nextStreamToSend() -> HTTP2StreamID? {
self.flushableStreams.first
private mutating func nextStreamToSend() -> HTTP2StreamID? {
// Round-robin via the FIFO queue. Pop from the front, skipping streams that are no longer flushable.
// Each stale entry is discarded at most once, so this is O(1) amortised.
while let streamID = self.sendQueue.popFirst() {
self.queuedStreams.remove(streamID)
if self.flushableStreams.contains(streamID) {
return streamID
}
}
return nil
}

internal mutating func updateWindowOfStream(_ streamID: HTTP2StreamID, newSize: Int32) {
Expand All @@ -153,6 +172,9 @@ internal struct OutboundFlowControlBuffer {
case .changed(newValue: true):
// Became writable, and specifically became _flushable_.
self.flushableStreams.insert(streamID)
if self.queuedStreams.insert(streamID).inserted {
self.sendQueue.append(streamID)
}
case .changed(newValue: false):
// Became unwritable.
self.flushableStreams.remove(streamID)
Expand Down Expand Up @@ -196,31 +218,37 @@ internal struct OutboundFlowControlBuffer {

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

let nextWrite = self.streamDataBuffers.modify(streamID: nextStreamID) {
(state: inout StreamFlowControlState) -> DataBuffer.BufferElement in
(state: inout StreamFlowControlState) -> (DataBuffer.BufferElement, isFlushable: Bool) in
let (nextWrite, writabilityState) = state.nextWrite(
maxSize: min(self.connectionWindowSize, self.maxFrameSize)
)

switch writabilityState {
case .changed(newValue: false):
self.flushableStreams.remove(nextStreamID)
return (nextWrite, isFlushable: false)
case .changed(newValue: true), .unchanged:
()
return (nextWrite, isFlushable: true)
}

return nextWrite
}
guard let (payload, promise) = nextWrite else {
guard let ((payload, promise), isFlushable) = nextWrite else {
// The stream was not present. This is weird, it shouldn't ever happen, but we tolerate it, and recurse.
self.flushableStreams.remove(nextStreamID)
return self.nextFlushedWritableFrame()
}

// If the stream is still flushable, put it back at the end of the rotation.
if isFlushable {
let inserted = self.queuedStreams.insert(nextStreamID).inserted
assert(inserted, "\(nextStreamID) unexpectedly still present in queuedStreams")
self.sendQueue.append(nextStreamID)
}

let frame = HTTP2Frame(streamID: nextStreamID, payload: payload)
return (frame, promise)
}
Expand All @@ -234,6 +262,9 @@ internal struct OutboundFlowControlBuffer {
case .changed(newValue: true):
// Became flushable
self.flushableStreams.insert($0.streamID)
if self.queuedStreams.insert($0.streamID).inserted {
self.sendQueue.append($0.streamID)
}
case .changed(newValue: false):
// Became unflushable.
self.flushableStreams.remove($0.streamID)
Expand Down
16 changes: 16 additions & 0 deletions Tests/NIOHTTP2Tests/OutboundFlowControlBufferTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -393,6 +393,22 @@ class OutboundFlowControlBufferTests: XCTestCase {
XCTAssertNil(self.buffer.nextFlushedWritableFrame())
}

func testFlushableStreamsAreNotStarved() {
// Three streams with three frames each.
let streamIDs: [HTTP2StreamID] = [1, 3, 5]
self.buffer.maxFrameSize = 5

for streamID in streamIDs {
self.buffer.streamCreated(streamID, initialWindowSize: 15)
let frame = self.createDataFrame(streamID, byteBufferSize: 15)
XCTAssertNoThrow(try self.buffer.processOutboundFrame(frame, promise: nil).assertNothing())
self.buffer.flushReceived()
}

// Every stream must be served once per round so that none can be starved.
XCTAssertEqual(self.receivedFrames().map { $0.streamID }, [1, 3, 5, 1, 3, 5, 1, 3, 5])
}

func testRejectsPrioritySelfDependency() {
XCTAssertThrowsError(
try self.buffer.priorityUpdate(
Expand Down
Loading