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
44 changes: 44 additions & 0 deletions Sources/SwiftNetwork/QUIC/FlowControl.swift
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,31 @@ struct FlowControlState: ~Copyable {
return true
}

// Accounts for inbound bytes that the peer sent, and that therefore consumed
// receive window, but that will never be delivered to the application because
// the stream carrying them was closed.
//
// `inboundMaxData` is anchored on `totalInboundBytesDelivered`, so unless
// discarded bytes are counted as consumed, the credit they used is never
// returned to the peer: the usable receive window shrinks by that amount for
// the remaining life of the connection, and enough discarded bytes stall it
// outright.
//
// The caller must have already added these bytes to
// `totalInOrderInboundBytesRead`; this only advances the delivered counter to
// match, which is what moves the MAX_DATA anchor.
fileprivate mutating func creditDiscardedInboundBytes(_ bytes: UInt64) {
guard bytes > 0 else { return }

let (newDelivered, deliveredOverflow) = totalInboundBytesDelivered.addingReportingOverflow(bytes)
guard !deliveredOverflow else { return }

// Delivered can never exceed the in-order total: the difference between
// them is what remains buffered awaiting the application.
totalInboundBytesDelivered = min(newDelivered, totalInOrderInboundBytesRead)
inboundBytesDeliveredSinceLastUpdate += bytes
}

// Outbound values (sending):

// Maximum number of bytes allowed to be sent to the peer, as
Expand Down Expand Up @@ -342,9 +367,28 @@ extension QUICConnection {
log.datapath(
"Zombie adjusted in-order inbound bytes changed from \(oldTotalInbound) to \(newValue))"
)
// The stream is already gone, so these bytes can never be delivered.
// Count them as consumed to release the credit they used.
flowControlState.creditDiscardedInboundBytes(delta)
}
}

// Releases the receive-window credit used by inbound bytes that arrived on a
// stream but that the application will never read, because the stream was
// closed with those bytes still buffered.
//
// The caller must have already accounted for `bytes` in the connection's
// in-order inbound total.
func creditDiscardedInboundBytes(_ bytes: UInt64) {
guard bytes > 0 else { return }
flowControlState.creditDiscardedInboundBytes(bytes)
log.datapath(
"Credited \(bytes) discarded inbound bytes; connection MAX_DATA anchor is now "
+ "\(self.flowControlState.totalInOrderInboundBytesRead)"
)
sendInboundFlowControlCredit()
}

func updateLastReceivedOffsetForZombie(lastOffsetDelta: UInt64) {
let connectionMaxData = flowControlState.inboundMaxData
let connectionCurrentLargestData = flowControlState.largestInboundByteOffsetReceived
Expand Down
48 changes: 47 additions & 1 deletion Sources/SwiftNetwork/QUIC/QUICConnection.swift
Original file line number Diff line number Diff line change
Expand Up @@ -387,6 +387,18 @@ public final class QUICConnection: ManyToManyApplicationStreamProtocol,

private var pendOutboundData = false // Don't immediately process application sends

// Set while `recovery` is exclusively borrowed for ACK processing.
//
// Acknowledging a packet can close a stream (a fully-ACKed RESET_STREAM or
// FIN), and closing a stream wants to flush frames. The no-argument
// `sendFrames()` passes `&recovery` inout, so doing that from inside the ACK
// walk would be a second overlapping modification of `recovery` and traps
// under exclusivity enforcement. While this is set, `sendFrames()` records
// the request instead of performing it, and the ACK path flushes once the
// borrow ends.
private var isProcessingAcks = false
private var deferredSendFramesRequested = false

// false == IPv6, true == IPv4
private(set) var initialAddressIsIPv4 = false

Expand Down Expand Up @@ -3046,7 +3058,15 @@ public final class QUICConnection: ManyToManyApplicationStreamProtocol,
// Adds recovery and applicationPendingItems to avoid extra begin/end acccess checking overhead
@discardableResult
func sendFrames(ignoreCongestionWindow: Bool = false, delayedACK: Bool = false) -> Bool {
sendFrames(
// ACK processing holds `recovery` exclusively, and acknowledging a packet
// can close a stream, which in turn wants to flush frames. Passing
// `&recovery` again here would overlap that borrow and trap, so record
// the request and let the ACK path flush once its borrow has ended.
if isProcessingAcks {
deferredSendFramesRequested = true
return false

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

So there must have been an overlapping access with Recovery here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confirmed that this was a defensive check

}
return sendFrames(
ignoreCongestionWindow: ignoreCongestionWindow,
delayedACK: delayedACK,
sentPackets: &sentPackets,
Expand Down Expand Up @@ -4131,6 +4151,12 @@ public final class QUICConnection: ManyToManyApplicationStreamProtocol,
localMaxStreamData: stream.flowControlState.inboundMaxData
)

// The application has given up on reading, so any bytes still
// buffered are discarded here. Return the receive-window credit they
// consumed; the zombie's final-size handling covers only the bytes
// still in flight beyond what we have already received.
stream.discardUnreadInboundBytes(connection: self)

// Do we delete this 'stream' somehow, now that it's a zombie?
// Once it's got no more references it will automatically taken care of
// with ARC. It may have references in pendingStartStreams (removed above)
Expand Down Expand Up @@ -4667,6 +4693,14 @@ public final class QUICConnection: ManyToManyApplicationStreamProtocol,
log.error("Error sending frames on stream close: \(error)")
}
}

// Anything the application did not read above is now unreachable: the
// flow is about to be torn down. Those bytes consumed connection receive
// window when they arrived, so hand that credit back before the buffers
// are dropped, otherwise the window shrinks for the rest of the
// connection's life.
stream.discardUnreadInboundBytes(connection: self)

stream.closed = true
deliverDisconnectedEvent(flow: flowID, error: error)
knownFlows.removeValue(forKey: streamID)
Expand Down Expand Up @@ -5222,6 +5256,11 @@ extension QUICConnection {
return false
}

// Acknowledging a packet can close a stream, and closing a stream wants
// to flush frames. Suppress those nested flushes for the duration of the
// `recovery` borrow below, then perform one flush afterwards if any were
// requested.
isProcessingAcks = true
recovery.receivedAck(
ack: frame,
ackedPath: path,
Expand All @@ -5232,6 +5271,13 @@ extension QUICConnection {
var sentPackets = NetworkUniqueDeque<SentPacketRecord>()
path.pmtudState.tryToSend(on: path, sentPackets: &sentPackets)
recovery.recordSentPackets(&sentPackets, connection: self)

// The borrow of `recovery` has ended, so it is safe to flush again.
isProcessingAcks = false
if deferredSendFramesRequested {
deferredSendFramesRequested = false
sendFrames()
}
return true
}

Expand Down
47 changes: 47 additions & 0 deletions Sources/SwiftNetwork/QUIC/QUICStream.swift
Original file line number Diff line number Diff line change
Expand Up @@ -875,6 +875,53 @@ public final class QUICStreamInstance: MultiplexedStreamFlow<QUICConnection>,
self.sendInboundFlowControlCreditIfNeeded(connection: connection)
}

// Discards every inbound byte still buffered for this stream and returns the
// flow control credit those bytes consumed.
//
// Called when the stream is closed with data the application never read,
// either still sitting in the reassembly queue or already dequeued into the
// upper receive queue awaiting a read. Those bytes counted against the
// connection's receive window when they arrived; without this the credit is
// never given back and the usable window shrinks permanently.
func discardUnreadInboundBytes(connection: QUICConnection) {
// Bytes handed to the upper layer but not yet read by the application.
// Dequeuing already advanced the reassembly queue's `currentOffset` past
// these, but flow control only counts them once the application reads,
// so they are still missing from the in-order total.
let pendingDelivery = UInt64(upperReceiveQueue.unclaimedLength)
// Contiguous bytes reassembled but not yet dequeued.
let pendingDequeue = UInt64(max(reassemblyQueue.availableToDequeue, 0))

// Anything the reassembly queue holds beyond the contiguous run is not
// yet part of the in-order total, so it has no credit to return here;
// the RESET_STREAM and zombie final-size paths cover those gaps.
let discardedBytes = pendingDelivery + pendingDequeue
guard discardedBytes > 0 else { return }

log.datapath(
"Discarding \(discardedBytes) unread inbound bytes on close "
+ "(\(pendingDelivery) awaiting read, \(pendingDequeue) awaiting dequeue)"
)

// Release the frames themselves before crediting, so the buffers are
// freed even if the connection is already tearing down.
upperReceiveQueue.finalizeAllFramesAsFailed()

// Advance the in-order total to cover everything the queue holds
// contiguously, which is what the application could have read. This adds
// the same delta to the connection-wide total.
let newInOrderTotal = UInt64(reassemblyQueue.currentOffset) + pendingDequeue
reassemblyQueue.dequeueAll()
updateFlowControlWithTotalInOrderInboundBytesRead(
newInOrderTotal,
connection: connection
)

// Both components are now part of the in-order total, so credit them as
// consumed to move the MAX_DATA anchor past them.
connection.creditDiscardedInboundBytes(discardedBytes)
}

@_optimize(speed)
func dequeueReassembledData(connection: QUICConnection) -> FrameArray? {
let totalLength = reassemblyQueue.availableToDequeue
Expand Down
Loading
Loading