fastcgi — cleanup

This commit is contained in:
Vladimir Krivosheev
2017-07-13 12:58:59 +02:00
parent 84aa38ffee
commit c5c90dffac
@@ -10,13 +10,13 @@ import org.jetbrains.io.Decoder
internal const val HEADER_LENGTH = 8
internal class FastCgiDecoder(private val errorOutputConsumer: Consumer<String>, private val responseHandler: FastCgiService) : Decoder(), Decoder.FullMessageConsumer<Void> {
private enum class State {
HEADER,
CONTENT
}
private enum class DecodeRecordState {
HEADER,
CONTENT
}
private var state = State.HEADER
internal class FastCgiDecoder(private val errorOutputConsumer: Consumer<String>, private val responseHandler: FastCgiService) : Decoder(), Decoder.FullMessageConsumer<Void> {
private var state = DecodeRecordState.HEADER
private enum class ProtocolStatus {
REQUEST_COMPLETE,
@@ -31,8 +31,8 @@ internal class FastCgiDecoder(private val errorOutputConsumer: Consumer<String>,
val STDERR = 7
}
private var type: Int = 0
private var id: Int = 0
private var type = 0
private var id = 0
private var contentLength: Int = 0
private var paddingLength: Int = 0
@@ -41,7 +41,7 @@ internal class FastCgiDecoder(private val errorOutputConsumer: Consumer<String>,
override fun messageReceived(context: ChannelHandlerContext, input: ByteBuf) {
while (true) {
when (state) {
FastCgiDecoder.State.HEADER -> {
DecodeRecordState.HEADER -> {
if (paddingLength > 0) {
if (input.readableBytes() > paddingLength) {
input.skipBytes(paddingLength)
@@ -55,21 +55,20 @@ internal class FastCgiDecoder(private val errorOutputConsumer: Consumer<String>,
}
val buffer = getBufferIfSufficient(input, HEADER_LENGTH, context) ?: return
decodeHeader(buffer)
state = State.CONTENT
if (contentLength > 0) {
readContent(input, context, contentLength, this)
}
state = State.HEADER
buffer.skipBytes(1)
type = buffer.readUnsignedByte().toInt()
id = buffer.readUnsignedShort()
contentLength = buffer.readUnsignedShort()
paddingLength = buffer.readUnsignedByte().toInt()
buffer.skipBytes(1)
state = DecodeRecordState.CONTENT
}
FastCgiDecoder.State.CONTENT -> {
DecodeRecordState.CONTENT -> {
if (contentLength > 0) {
readContent(input, context, contentLength, this)
}
state = State.HEADER
state = DecodeRecordState.HEADER
}
}
}
@@ -95,52 +94,30 @@ internal class FastCgiDecoder(private val errorOutputConsumer: Consumer<String>,
}
}
private fun decodeHeader(buffer: ByteBuf) {
buffer.skipBytes(1)
type = buffer.readUnsignedByte().toInt()
id = buffer.readUnsignedShort()
contentLength = buffer.readUnsignedShort()
paddingLength = buffer.readUnsignedByte().toInt()
buffer.skipBytes(1)
}
override fun contentReceived(buffer: ByteBuf, context: ChannelHandlerContext, isCumulateBuffer: Boolean): Void? {
when (type) {
RecordType.END_REQUEST -> {
val appStatus = buffer.readInt()
val protocolStatus = buffer.readUnsignedByte().toInt()
if (appStatus != 0 || protocolStatus != ProtocolStatus.REQUEST_COMPLETE.ordinal) {
LOG.warn("Protocol status $protocolStatus")
dataBuffers.remove(id)
responseHandler.responseReceived(id, null)
}
else if (protocolStatus == ProtocolStatus.REQUEST_COMPLETE.ordinal) {
responseHandler.responseReceived(id, dataBuffers.remove(id))
}
}
RecordType.STDOUT -> {
var data = dataBuffers.get(id)
val sliced = if (isCumulateBuffer) buffer else buffer.slice(buffer.readerIndex(), contentLength)
if (data == null) {
dataBuffers.put(id, sliced)
}
else if (data is CompositeByteBuf) {
data.addComponent(sliced)
data.writerIndex(data.writerIndex() + sliced.readableBytes())
}
else {
if (sliced is CompositeByteBuf) {
data = sliced.addComponent(0, data)
data.writerIndex(data.writerIndex() + data.readableBytes())
when (data) {
null -> dataBuffers.put(id, sliced)
is CompositeByteBuf -> {
data.addComponent(sliced)
data.writerIndex(data.writerIndex() + sliced.readableBytes())
}
else {
// must be computed here before we set data to new composite buffer
val newLength = data.readableBytes() + sliced.readableBytes()
data = context.alloc().compositeBuffer(Decoder.DEFAULT_MAX_COMPOSITE_BUFFER_COMPONENTS).addComponents(data, sliced)
data.writerIndex(data.writerIndex() + newLength)
else -> {
if (sliced is CompositeByteBuf) {
data = sliced.addComponent(0, data)
data.writerIndex(data.writerIndex() + data.readableBytes())
}
else {
// must be computed here before we set data to new composite buffer
val newLength = data.readableBytes() + sliced.readableBytes()
data = context.alloc().compositeBuffer(Decoder.DEFAULT_MAX_COMPOSITE_BUFFER_COMPONENTS).addComponents(data, sliced)
data.writerIndex(data.writerIndex() + newLength)
}
dataBuffers.put(id, data)
}
dataBuffers.put(id, data)
}
sliced.retain()
}
@@ -154,7 +131,25 @@ internal class FastCgiDecoder(private val errorOutputConsumer: Consumer<String>,
}
}
else -> LOG.error("Unknown type $type")
RecordType.END_REQUEST -> {
val appStatus = buffer.readInt()
val protocolStatus = buffer.readUnsignedByte().toInt()
if (appStatus != 0 || protocolStatus != ProtocolStatus.REQUEST_COMPLETE.ordinal) {
LOG.warn("Protocol status $protocolStatus")
dataBuffers.remove(id)
responseHandler.responseReceived(id, null)
}
else if (protocolStatus == ProtocolStatus.REQUEST_COMPLETE.ordinal) {
responseHandler.responseReceived(id, dataBuffers.remove(id))
}
else {
LOG.warn("protocolStatus $protocolStatus")
}
}
else -> {
LOG.error("Unknown type $type")
}
}
return null
}