mirror of
https://github.com/swift-server/async-http-client.git
synced 2026-06-02 07:37:34 +00:00
Full support for bidirectional streaming (#879)
> ## Note: > This is a long LLM generated PR description. However it captures very well, what has been changed and has already been reduced for brevity. The PR is sadly quite complex but I think the description captures the changes quite well. This is foundational work needed to properly support HTTP trailers and scenarios where the server sends a complete response before the client finishes uploading (e.g., early rejection, 100-continue flows, or bidirectional streaming protocols). ## Changes ### State Machine Improvements - **Added `endForwarded` state** to `Transaction.StateMachine.RequestStreamState` - This new state distinguishes between "request data forwarded to the channel" and "request data written to the network" - Properly handles the race condition where response completes before the request write completes - **Renamed `succeedRequest` → `forwardResponseEnd`** in both `HTTPRequestStateMachine.Action` and `HTTP1ConnectionStateMachine.Action` - Better reflects the semantic meaning: we're forwarding the end of the response stream, not necessarily succeeding the entire request yet - More accurate naming for bidirectional streaming scenarios ### Protocol Changes - **Added `requestBodyStreamSent()` to `HTTPExecutableRequest` protocol** - Called by the channel handler when the request body stream has been fully written to the network - Allows proper coordination between request and response stream completion - Implemented in both `Transaction` and `RequestBag` ### Request State Machine Updates - **Updated `FinalSuccessfulRequestAction`** - Changed `.sendRequestEnd(EventLoopPromise<Void>?)` to simpler `.requestDone` - Added `.none` case for when response completes but request is still in-flight - Removed the need to pass promises around, simplifying the state machine - **`sendRequestEnd` action now includes `FinalSuccessfulRequestAction`** - Allows the state machine to signal what should happen after the request completes - Enables proper cleanup coordination (idle connection, close, or continue) ### Channel Handler Updates - **HTTP1ClientChannelHandler** - `sendRequestEnd` now properly handles scenarios where response has already completed - Added future callback to coordinate request completion with final actions - Properly manages connection state (idle vs close) based on both streams completing - **HTTP2ClientRequestHandler** - Updated to handle new `sendRequestEnd` signature - Properly ignores HTTP/1-specific final actions (like `.requestDone`) ### RequestBag State Machine - **Added `endReceived` state to `ResponseStreamState`** - Tracks when the response has completed while request is still ongoing - Enables proper sequencing: response end → request end → task completion - **Updated `FinishAction`** - Added `.forwardStreamFinishedAndSucceedTask` for the case where both streams complete simultaneously - Ensures delegate methods are called in the correct order ### Error Handling - **Improved failure handling in `Transaction.StateMachine`** - Now properly handles errors that occur after response completes but before request finishes - Added `cancelExecutor` action to the fail path - Executor is now passed to `failRequestStreamContinuation` for proper cleanup ## Technical Details ### The Problem Previously, when a server sent a complete response before the client finished uploading the request body, AHC would: 1. Receive the full response (head, body, end) 2. But NOT inform the user that the response was complete if the request was still streaming 3. Only succeed the request after both streams completed This made it impossible to implement proper bidirectional streaming or handle scenarios like: - Server rejecting a large upload early (e.g., 413 Payload Too Large) - 100-continue flows where the server responds before request completes - HTTP trailers sent by the server ### The Solution The new state machine properly tracks four completion states: 1. **Neither complete**: Normal request/response in flight 2. **Response complete, request ongoing**: New `endForwarded`/`endReceived` states 3. **Request complete, response ongoing**: Existing logic 4. **Both complete**: Request succeeds The key insight is the `endForwarded` state, which represents "we've given all request data to the channel, but it hasn't been written to the network yet". This allows us to: - Immediately forward response completion to the user - Wait for the write to complete before cleaning up resources - Properly sequence connection state transitions ## Future Work This PR lays the groundwork for: - Proper internal HTTP trailer support (both sending and receiving) --------- Co-authored-by: George Barnett <gbarnett@apple.com>
This commit is contained in:
co-authored by
George Barnett
parent
4b99975677
commit
e2ab0d176f
@@ -242,7 +242,46 @@ final class HTTP1ClientChannelHandler: ChannelDuplexHandler {
|
||||
case .sendBodyPart(let part, let writePromise):
|
||||
context.writeAndFlush(self.wrapOutboundOut(.body(part)), promise: writePromise)
|
||||
|
||||
case .sendRequestEnd(let writePromise):
|
||||
case .sendRequestEnd(let writePromise, let finalAction):
|
||||
|
||||
let writePromise = writePromise ?? context.eventLoop.makePromise(of: Void.self)
|
||||
// We need to defer succeeding the old request to avoid ordering issues
|
||||
|
||||
writePromise.futureResult.hop(to: context.eventLoop).assumeIsolated().whenComplete { result in
|
||||
guard let oldRequest = self.request else {
|
||||
// in the meantime an error might have happened, which is why this request is
|
||||
// not reference anymore.
|
||||
return
|
||||
}
|
||||
oldRequest.requestBodyStreamSent()
|
||||
switch result {
|
||||
case .success:
|
||||
// If our final action is not `none`, that means we've already received
|
||||
// the complete response. As a result, once we've uploaded all the body parts
|
||||
// we need to tell the pool that the connection is idle or, if we were asked to
|
||||
// close when we're done, send the close. Either way, we then succeed the request
|
||||
|
||||
switch finalAction {
|
||||
case .none:
|
||||
// we must not nil out the request here, as we are still uploading the request
|
||||
// and therefore still need the reference to it.
|
||||
break
|
||||
|
||||
case .informConnectionIsIdle:
|
||||
self.request = nil
|
||||
self.onConnectionIdle()
|
||||
|
||||
case .close:
|
||||
self.request = nil
|
||||
context.close(promise: nil)
|
||||
}
|
||||
|
||||
case .failure(let error):
|
||||
context.close(promise: nil)
|
||||
oldRequest.fail(error)
|
||||
}
|
||||
}
|
||||
|
||||
context.writeAndFlush(self.wrapOutboundOut(.end(nil)), promise: writePromise)
|
||||
|
||||
if let readTimeoutAction = self.idleReadTimeoutStateMachine?.requestEndSent() {
|
||||
@@ -300,7 +339,7 @@ final class HTTP1ClientChannelHandler: ChannelDuplexHandler {
|
||||
// that the request is neither failed nor finished yet
|
||||
self.request!.receiveResponseBodyParts(buffer)
|
||||
|
||||
case .succeedRequest(let finalAction, let buffer):
|
||||
case .forwardResponseEnd(let finalAction, let buffer):
|
||||
// We can force unwrap the request here, as we have just validated in the state machine,
|
||||
// that the request is neither failed nor finished yet
|
||||
|
||||
@@ -312,39 +351,20 @@ final class HTTP1ClientChannelHandler: ChannelDuplexHandler {
|
||||
// other way around.
|
||||
|
||||
let oldRequest = self.request!
|
||||
self.request = nil
|
||||
self.runTimeoutAction(.clearIdleReadTimeoutTimer, context: context)
|
||||
self.runTimeoutAction(.clearIdleWriteTimeoutTimer, context: context)
|
||||
|
||||
switch finalAction {
|
||||
case .close:
|
||||
self.request = nil
|
||||
context.close(promise: nil)
|
||||
oldRequest.receiveResponseEnd(buffer, trailers: nil)
|
||||
case .sendRequestEnd(let writePromise, let shouldClose):
|
||||
let writePromise = writePromise ?? context.eventLoop.makePromise(of: Void.self)
|
||||
// We need to defer succeeding the old request to avoid ordering issues
|
||||
writePromise.futureResult.hop(to: context.eventLoop).assumeIsolated().whenComplete { result in
|
||||
switch result {
|
||||
case .success:
|
||||
// If our final action was `sendRequestEnd`, that means we've already received
|
||||
// the complete response. As a result, once we've uploaded all the body parts
|
||||
// we need to tell the pool that the connection is idle or, if we were asked to
|
||||
// close when we're done, send the close. Either way, we then succeed the request
|
||||
if shouldClose {
|
||||
context.close(promise: nil)
|
||||
} else {
|
||||
self.onConnectionIdle()
|
||||
}
|
||||
|
||||
oldRequest.receiveResponseEnd(buffer, trailers: nil)
|
||||
case .failure(let error):
|
||||
context.close(promise: nil)
|
||||
oldRequest.fail(error)
|
||||
}
|
||||
}
|
||||
case .none:
|
||||
oldRequest.receiveResponseEnd(buffer, trailers: nil)
|
||||
|
||||
context.writeAndFlush(self.wrapOutboundOut(.end(nil)), promise: writePromise)
|
||||
case .informConnectionIsIdle:
|
||||
self.request = nil
|
||||
self.onConnectionIdle()
|
||||
oldRequest.receiveResponseEnd(buffer, trailers: nil)
|
||||
}
|
||||
|
||||
@@ -27,18 +27,12 @@ struct HTTP1ConnectionStateMachine {
|
||||
}
|
||||
|
||||
enum Action {
|
||||
/// A action to execute, when we consider a request "done".
|
||||
/// An additional action to execute, when either the response or request stream has finished.
|
||||
enum FinalSuccessfulStreamAction {
|
||||
/// Nothing todo
|
||||
case none
|
||||
/// Close the connection
|
||||
case close
|
||||
/// If the server has replied, with a status of 200...300 before all data was sent, a request is considered succeeded,
|
||||
/// as soon as we wrote the request end onto the wire.
|
||||
///
|
||||
/// The promise is an optional write promise.
|
||||
///
|
||||
/// `shouldClose` records whether we have attached a Connection: close header to this request, and so the connection should
|
||||
/// be terminated
|
||||
case sendRequestEnd(EventLoopPromise<Void>?, shouldClose: Bool)
|
||||
/// Inform an observer that the connection has become idle
|
||||
case informConnectionIsIdle
|
||||
}
|
||||
@@ -63,7 +57,7 @@ struct HTTP1ConnectionStateMachine {
|
||||
startIdleTimer: Bool
|
||||
)
|
||||
case sendBodyPart(IOData, EventLoopPromise<Void>?)
|
||||
case sendRequestEnd(EventLoopPromise<Void>?)
|
||||
case sendRequestEnd(EventLoopPromise<Void>?, FinalSuccessfulStreamAction)
|
||||
case failSendBodyPart(Error, EventLoopPromise<Void>?)
|
||||
case failSendStreamFinished(Error, EventLoopPromise<Void>?)
|
||||
|
||||
@@ -72,9 +66,9 @@ struct HTTP1ConnectionStateMachine {
|
||||
|
||||
case forwardResponseHead(HTTPResponseHead, pauseRequestBodyStream: Bool)
|
||||
case forwardResponseBodyParts(CircularBuffer<ByteBuffer>)
|
||||
case forwardResponseEnd(FinalSuccessfulStreamAction, CircularBuffer<ByteBuffer>)
|
||||
|
||||
case failRequest(Error, FinalFailedStreamAction)
|
||||
case succeedRequest(FinalSuccessfulStreamAction, CircularBuffer<ByteBuffer>)
|
||||
|
||||
case read
|
||||
case close
|
||||
@@ -433,15 +427,11 @@ extension HTTP1ConnectionStateMachine.State {
|
||||
return .resumeRequestBodyStream
|
||||
case .sendBodyPart(let part, let writePromise):
|
||||
return .sendBodyPart(part, writePromise)
|
||||
case .sendRequestEnd(let writePromise):
|
||||
return .sendRequestEnd(writePromise)
|
||||
case .forwardResponseHead(let head, let pauseRequestBodyStream):
|
||||
return .forwardResponseHead(head, pauseRequestBodyStream: pauseRequestBodyStream)
|
||||
case .forwardResponseBodyParts(let parts):
|
||||
return .forwardResponseBodyParts(parts)
|
||||
case .succeedRequest(let finalAction, let finalParts):
|
||||
case .sendRequestEnd(let writePromise, let finalAction):
|
||||
guard case .inRequest(_, close: let close) = self else {
|
||||
fatalError("Invalid state: \(self)")
|
||||
assertionFailure("Invalid state: \(self)")
|
||||
self = .closing
|
||||
return .failRequest(HTTPClientError.internalStateFailure(), .close(writePromise))
|
||||
}
|
||||
|
||||
let newFinalAction: HTTP1ConnectionStateMachine.Action.FinalSuccessfulStreamAction
|
||||
@@ -449,14 +439,48 @@ extension HTTP1ConnectionStateMachine.State {
|
||||
case .close:
|
||||
self = .closing
|
||||
newFinalAction = .close
|
||||
case .sendRequestEnd(let writePromise):
|
||||
self = .idle
|
||||
newFinalAction = .sendRequestEnd(writePromise, shouldClose: close)
|
||||
case .requestDone:
|
||||
if close {
|
||||
self = .closing
|
||||
newFinalAction = .close
|
||||
} else {
|
||||
self = .idle
|
||||
newFinalAction = .informConnectionIsIdle
|
||||
}
|
||||
case .none:
|
||||
self = .idle
|
||||
newFinalAction = close ? .close : .informConnectionIsIdle
|
||||
newFinalAction = .none
|
||||
}
|
||||
return .succeedRequest(newFinalAction, finalParts)
|
||||
return .sendRequestEnd(writePromise, newFinalAction)
|
||||
|
||||
case .forwardResponseHead(let head, let pauseRequestBodyStream):
|
||||
return .forwardResponseHead(head, pauseRequestBodyStream: pauseRequestBodyStream)
|
||||
case .forwardResponseBodyParts(let parts):
|
||||
return .forwardResponseBodyParts(parts)
|
||||
case .forwardResponseEnd(let finalAction, let finalParts):
|
||||
guard case .inRequest(_, close: let close) = self else {
|
||||
assertionFailure("Invalid state: \(self)")
|
||||
self = .closing
|
||||
return .failRequest(HTTPClientError.internalStateFailure(), .close(nil))
|
||||
}
|
||||
|
||||
let newFinalAction: HTTP1ConnectionStateMachine.Action.FinalSuccessfulStreamAction
|
||||
switch finalAction {
|
||||
case .close:
|
||||
self = .closing
|
||||
newFinalAction = .close
|
||||
case .requestDone:
|
||||
if close {
|
||||
self = .closing
|
||||
newFinalAction = .close
|
||||
} else {
|
||||
self = .idle
|
||||
newFinalAction = .informConnectionIsIdle
|
||||
}
|
||||
case .none:
|
||||
// request is ongoing. request stream is still alive
|
||||
newFinalAction = .none
|
||||
}
|
||||
return .forwardResponseEnd(newFinalAction, finalParts)
|
||||
|
||||
case .failRequest(let error, let finalAction):
|
||||
switch self {
|
||||
|
||||
Reference in New Issue
Block a user