mirror of
https://github.com/swift-server/async-http-client.git
synced 2026-06-02 07:37:34 +00:00
Add an idle write timeout (#718)
This commit is contained in:
@@ -35,8 +35,16 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
|
||||
private var request: HTTPExecutableRequest? {
|
||||
didSet {
|
||||
if let newRequest = self.request, let idleReadTimeout = newRequest.requestOptions.idleReadTimeout {
|
||||
self.idleReadTimeoutStateMachine = .init(timeAmount: idleReadTimeout)
|
||||
if let newRequest = self.request {
|
||||
if let idleReadTimeout = newRequest.requestOptions.idleReadTimeout {
|
||||
self.idleReadTimeoutStateMachine = .init(timeAmount: idleReadTimeout)
|
||||
}
|
||||
if let idleWriteTimeout = newRequest.requestOptions.idleWriteTimeout {
|
||||
self.idleWriteTimeoutStateMachine = .init(
|
||||
timeAmount: idleWriteTimeout,
|
||||
isWritabilityEnabled: self.channelContext?.channel.isWritable ?? false
|
||||
)
|
||||
}
|
||||
} else {
|
||||
self.idleReadTimeoutStateMachine = nil
|
||||
}
|
||||
@@ -46,6 +54,9 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
private var idleReadTimeoutStateMachine: IdleReadStateMachine?
|
||||
private var idleReadTimeoutTimer: Scheduled<Void>?
|
||||
|
||||
private var idleWriteTimeoutStateMachine: IdleWriteStateMachine?
|
||||
private var idleWriteTimeoutTimer: Scheduled<Void>?
|
||||
|
||||
init(eventLoop: EventLoop) {
|
||||
self.eventLoop = eventLoop
|
||||
}
|
||||
@@ -77,6 +88,10 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
}
|
||||
|
||||
func channelWritabilityChanged(context: ChannelHandlerContext) {
|
||||
if let timeoutAction = self.idleWriteTimeoutStateMachine?.channelWritabilityChanged(context: context) {
|
||||
self.runTimeoutAction(timeoutAction, context: context)
|
||||
}
|
||||
|
||||
let action = self.state.writabilityChanged(writable: context.channel.isWritable)
|
||||
self.run(action, context: context)
|
||||
}
|
||||
@@ -110,6 +125,10 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
// a single request.
|
||||
self.request = request
|
||||
|
||||
if let timeoutAction = self.idleWriteTimeoutStateMachine?.write() {
|
||||
self.runTimeoutAction(timeoutAction, context: context)
|
||||
}
|
||||
|
||||
request.willExecuteRequest(self)
|
||||
|
||||
let action = self.state.startRequest(
|
||||
@@ -153,8 +172,12 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
request.resumeRequestBodyStream()
|
||||
}
|
||||
if startIdleTimer {
|
||||
if let timeoutAction = self.idleReadTimeoutStateMachine?.requestEndSent() {
|
||||
self.runTimeoutAction(timeoutAction, context: context)
|
||||
if let readTimeoutAction = self.idleReadTimeoutStateMachine?.requestEndSent() {
|
||||
self.runTimeoutAction(readTimeoutAction, context: context)
|
||||
}
|
||||
|
||||
if let writeTimeoutAction = self.idleWriteTimeoutStateMachine?.requestEndSent() {
|
||||
self.runTimeoutAction(writeTimeoutAction, context: context)
|
||||
}
|
||||
}
|
||||
case .pauseRequestBodyStream:
|
||||
@@ -168,8 +191,12 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
case .sendRequestEnd(let writePromise):
|
||||
context.writeAndFlush(self.wrapOutboundOut(.end(nil)), promise: writePromise)
|
||||
|
||||
if let timeoutAction = self.idleReadTimeoutStateMachine?.requestEndSent() {
|
||||
self.runTimeoutAction(timeoutAction, context: context)
|
||||
if let readTimeoutAction = self.idleReadTimeoutStateMachine?.requestEndSent() {
|
||||
self.runTimeoutAction(readTimeoutAction, context: context)
|
||||
}
|
||||
|
||||
if let writeTimeoutAction = self.idleWriteTimeoutStateMachine?.requestEndSent() {
|
||||
self.runTimeoutAction(writeTimeoutAction, context: context)
|
||||
}
|
||||
|
||||
case .read:
|
||||
@@ -295,6 +322,36 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
}
|
||||
}
|
||||
|
||||
private func runTimeoutAction(_ action: IdleWriteStateMachine.Action, context: ChannelHandlerContext) {
|
||||
switch action {
|
||||
case .startIdleWriteTimeoutTimer(let timeAmount):
|
||||
assert(self.idleWriteTimeoutTimer == nil, "Expected there is no timeout timer so far.")
|
||||
|
||||
self.idleWriteTimeoutTimer = self.eventLoop.scheduleTask(in: timeAmount) {
|
||||
guard self.idleWriteTimeoutTimer != nil else { return }
|
||||
let action = self.state.idleWriteTimeoutTriggered()
|
||||
self.run(action, context: context)
|
||||
}
|
||||
case .resetIdleWriteTimeoutTimer(let timeAmount):
|
||||
if let oldTimer = self.idleWriteTimeoutTimer {
|
||||
oldTimer.cancel()
|
||||
}
|
||||
|
||||
self.idleWriteTimeoutTimer = self.eventLoop.scheduleTask(in: timeAmount) {
|
||||
guard self.idleWriteTimeoutTimer != nil else { return }
|
||||
let action = self.state.idleWriteTimeoutTriggered()
|
||||
self.run(action, context: context)
|
||||
}
|
||||
case .clearIdleWriteTimeoutTimer:
|
||||
if let oldTimer = self.idleWriteTimeoutTimer {
|
||||
self.idleWriteTimeoutTimer = nil
|
||||
oldTimer.cancel()
|
||||
}
|
||||
case .none:
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: Private HTTPRequestExecutor
|
||||
|
||||
private func writeRequestBodyPart0(_ data: IOData, request: HTTPExecutableRequest, promise: EventLoopPromise<Void>?) {
|
||||
@@ -308,6 +365,10 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
return
|
||||
}
|
||||
|
||||
if let timeoutAction = self.idleWriteTimeoutStateMachine?.write() {
|
||||
self.runTimeoutAction(timeoutAction, context: context)
|
||||
}
|
||||
|
||||
let action = self.state.requestStreamPartReceived(data, promise: promise)
|
||||
self.run(action, context: context)
|
||||
}
|
||||
@@ -338,6 +399,10 @@ final class HTTP2ClientRequestHandler: ChannelDuplexHandler {
|
||||
return
|
||||
}
|
||||
|
||||
if let timeoutAction = self.idleWriteTimeoutStateMachine?.cancelRequest() {
|
||||
self.runTimeoutAction(timeoutAction, context: context)
|
||||
}
|
||||
|
||||
let action = self.state.requestCancelled()
|
||||
self.run(action, context: context)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user