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
79 changes: 65 additions & 14 deletions Sources/SwiftNetwork/Connection/Connection.swift
Original file line number Diff line number Diff line change
Expand Up @@ -2106,6 +2106,21 @@ extension NetworkChannel where ApplicationProtocol: StreamProtocol {
public let isComplete: Bool
}

public struct StreamSpanMessage: ~Escapable {
@_lifetime(borrow content)
init(content: RawSpan? = nil, offset: Int = 0, isComplete: Bool = false, lastChunkOfBatch: Bool = false) {
self.content = content
self.offset = offset
self.isComplete = isComplete
self.lastChunkOfBatch = lastChunkOfBatch
}

public let content: RawSpan?
public let offset: Int
public let isComplete: Bool
public let lastChunkOfBatch: Bool
}

public func send(_ message: StreamMessage, completion: (@Sendable (Result<Void, NetworkError>) -> Void)? = nil) {
let endpointFlow = self.endpointFlow
endpointFlow.async {
Expand Down Expand Up @@ -2141,15 +2156,48 @@ extension NetworkChannel where ApplicationProtocol: StreamProtocol {
atMost maxBytes: Int,
completion: @escaping @Sendable (Result<StreamMessage, NetworkError>) -> Void
) {
let readRequest = ReadRequest(minimumBytes: minBytes, maximumBytes: maxBytes, maximumFrames: Int.max) {
(content, isComplete, isFinal, error) in
if let error = error {
completion(.failure(error))
} else {
completion(.success(.message(content: content, isComplete: isComplete)))
let endpointFlow = self.endpointFlow
endpointFlow.async {
let readRequest = ReadRequest(minimumBytes: minBytes, maximumBytes: maxBytes) {
(content, isComplete, isFinal, error) in
if let error = error {
completion(.failure(error))
} else {
completion(.success(.message(content: content, isComplete: isComplete)))
}
}
self.endpointFlow.addReadRequestOnContext(readRequest)
}
}

public func receive(
atLeast minBytes: Int,
atMost maxBytes: Int,
maximumChunks: Int,
completion: @escaping @Sendable (Result<StreamSpanMessage, NetworkError>) -> Void
) {
let endpointFlow = self.endpointFlow
endpointFlow.async {
let readRequest = ReadRequest(minimumBytes: minBytes, maximumBytes: maxBytes, maximumFrames: maximumChunks)
{
(content, offset, isComplete, isFinal, lastChunkOfBatch, error) in
if let error = error {
completion(.failure(error))
} else {
completion(
.success(
.init(
content: content,
offset: offset,
isComplete: isComplete,
lastChunkOfBatch: lastChunkOfBatch
)
)
)
}
}
self.endpointFlow.addReadRequestOnContext(readRequest)
}
self.endpointFlow.addReadRequest(readRequest)
}
}

Expand Down Expand Up @@ -2191,14 +2239,17 @@ extension NetworkChannel where ApplicationProtocol: DatagramProtocol {
}

public func receive(completion: @escaping @Sendable (Result<DatagramMessage, NetworkError>) -> Void) {
let readRequest = ReadRequest(minimumBytes: 1, maximumBytes: Int.max, maximumFrames: 1) {
(content, isComplete, isFinal, error) in
if let error = error {
completion(.failure(error))
} else {
completion(.success(.message(content: content)))
let endpointFlow = self.endpointFlow
endpointFlow.async {
let readRequest = ReadRequest(maximumFrames: 1) {
(content, isComplete, isFinal, error) in
if let error = error {
completion(.failure(error))
} else {
completion(.success(.message(content: content)))
}
}
self.endpointFlow.addReadRequestOnContext(readRequest)
}
self.endpointFlow.addReadRequest(readRequest)
}
}
125 changes: 90 additions & 35 deletions Sources/SwiftNetwork/EndpointFlow/EndpointFlow.swift
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ final class EndpointFlow: CustomDebugStringConvertible {
let context: NetworkContext
let identifier: UInt64
var writeRequests = NetworkUniqueDeque<WriteRequest>()
var readRequests = [ReadRequest]()
var readRequests = NetworkUniqueDeque<ReadRequest>()
var stateUpdateHandler: ((State) -> Void)? = nil
var cancelRequested = false
var teardownComplete = false
Expand Down Expand Up @@ -233,37 +233,80 @@ final class EndpointFlow: CustomDebugStringConvertible {

switch self.flowProtocol {
case .stream(let flow):
while true {
if let readRequest = self.readRequests.first {
if let content = flow.read(
minimumBytes: readRequest.minimumBytes,
maximumBytes: readRequest.maximumBytes
) {
// TODO: Get the actual metadata
readRequest.complete(content: content, isComplete: false, isFinal: true)
// TODO: This is not efficient. Probably better to use an ArraySlice here
self.readRequests.removeFirst()
} else {
while !self.readRequests.isEmpty {
if self.readRequests[0].expectsSpan {
guard
var frames = flow.readFrames(
minimumBytes: self.readRequests[0].minimumBytes,
maximumBytes: self.readRequests[0].maximumBytes
)
else {
flow.waitForInboundDataAvailable(completion: self.inputAvailable)
break
}
let readRequest = self.readRequests.removeFirst()
var offset = 0
while var frame = frames.popFirst() {
let isLastFrame = frames.isEmpty
if let bytes = frame.bytes {

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.

Anything else we can do here at the moment besides going through this computed property?

readRequest.complete(
bytes: bytes,
offset: offset,
isComplete: frame.metadataComplete,
isFinal: true,
lastChunkOfBatch: isLastFrame
)
offset += bytes.byteCount
}
frame.finalize(success: true)
}
} else {
break
guard
let content = flow.read(
minimumBytes: self.readRequests[0].minimumBytes,
maximumBytes: self.readRequests[0].maximumBytes
)
else {
flow.waitForInboundDataAvailable(completion: self.inputAvailable)
break
}
// TODO: Get the actual metadata
let readRequest = self.readRequests.removeFirst()
readRequest.complete(content: content, isComplete: false, isFinal: true)
}
}
case .datagram(let flow):
while true {
if let readRequest = self.readRequests.first {
if let content = flow.read() {
readRequest.complete(content: content, isComplete: true, isFinal: false)
// TODO: This is not efficient. Probably better to use an ArraySlice here
self.readRequests.removeFirst()
} else {
while !self.readRequests.isEmpty {
if self.readRequests[0].expectsSpan {
guard var frames = flow.readFrames(maximumFrames: self.readRequests[0].maximumFrames) else {
flow.waitForInboundDataAvailable(completion: self.inputAvailable)
break
}

let readRequest = self.readRequests.removeFirst()
var offset = 0
while var frame = frames.popFirst() {
let isLastFrame = frames.isEmpty
if let bytes = frame.bytes {
readRequest.complete(
bytes: bytes,
offset: offset,
isComplete: frame.metadataComplete,
isFinal: false,
lastChunkOfBatch: isLastFrame
)
offset += bytes.byteCount
}
frame.finalize(success: true)
}
} else {
break
guard let content = flow.read() else {
flow.waitForInboundDataAvailable(completion: self.inputAvailable)
break
}

let readRequest = self.readRequests.removeFirst()
readRequest.complete(content: content, isComplete: true, isFinal: false)
}
}
case .none:
Expand All @@ -277,24 +320,25 @@ final class EndpointFlow: CustomDebugStringConvertible {

func addWriteRequestOnContext(_ writeRequest: consuming WriteRequest) {
var writeRequest: WriteRequest? = writeRequest
self.startIfNeeded()
startIfNeeded()
if let takenRequest = writeRequest.take() {
self.writeRequests.append(takenRequest)
writeRequests.append(takenRequest)
}
if self.state == .ready {
self.write()
if state == .ready {
write()
}
}

func addReadRequest(_ readRequest: ReadRequest) {
self.parameters.context.async {
self.startIfNeeded()
self.readRequests.append(readRequest)
// If state is ready and this is the first read request, then try to start reading.
// Otherwise, wait for inputAvailable to trigger a call to read()
if self.state == .ready && self.readRequests.count == 1 {
self.read()
}
func addReadRequestOnContext(_ readRequest: consuming ReadRequest) {
var readRequest: ReadRequest? = readRequest
startIfNeeded()
if let takenRequest = readRequest.take() {
readRequests.append(takenRequest)
}
// If state is ready and this is the first read request, then try to start reading.
// Otherwise, wait for inputAvailable to trigger a call to read()
if state == .ready && readRequests.count == 1 {
read()
}
}

Expand Down Expand Up @@ -396,7 +440,18 @@ final class EndpointFlow: CustomDebugStringConvertible {
}
while !self.readRequests.isEmpty {
let readRequest = self.readRequests.removeFirst()
readRequest.complete(content: nil, isComplete: false, isFinal: true, error: .posix(ECANCELED))
if readRequest.expectsSpan {
readRequest.complete(
bytes: nil,
offset: 0,
isComplete: false,
isFinal: true,
lastChunkOfBatch: true,
error: .posix(ECANCELED)
)
} else {
readRequest.complete(content: nil, isComplete: false, isFinal: true, error: .posix(ECANCELED))
}
}
}

Expand Down
27 changes: 27 additions & 0 deletions Sources/SwiftNetwork/EndpointFlow/EndpointFlowProtocols.swift
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,19 @@ final class DatagramEndpointFlowProtocol: EndpointFlowProtocol<InboundDatagramLi
}
}

func readFrames(maximumFrames: Int) -> FrameArray? {
fromExternal {
do throws(NetworkError) {
return try lower.invokeReceiveDatagrams(
reference,
maximumDatagramCount: maximumFrames
)
} catch {
return nil
}
}
}

func read() -> [UInt8]? {
fromExternal {
do throws(NetworkError) {
Expand Down Expand Up @@ -528,6 +541,20 @@ final class StreamEndpointFlowProtocol: EndpointFlowProtocol<InboundStreamLinkag
}
}

func readFrames(minimumBytes: Int, maximumBytes: Int) -> FrameArray? {
fromExternal {
do throws(NetworkError) {
return try lower.invokeReceiveStreamData(
reference,
minimumBytes: minimumBytes,
maximumBytes: maximumBytes
)
} catch {
return nil
}
}
}

func read(minimumBytes: Int, maximumBytes: Int) -> [UInt8]? {
fromExternal {
do throws(NetworkError) {
Expand Down
Loading