-
Notifications
You must be signed in to change notification settings - Fork 5
/
MessageReader.java
157 lines (140 loc) · 6.09 KB
/
MessageReader.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
/*
* Java Reliable Event Logging Protocol Library Server Implementation RLP-03
* Copyright (C) 2021 Suomen Kanuuna Oy
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU Affero General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU Affero General Public License for more details.
*
* You should have received a copy of the GNU Affero General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*
*
* Additional permission under GNU Affero General Public License version 3
* section 7
*
* If you modify this Program, or any covered work, by linking or combining it
* with other code, such other code is not for that reason alone subject to any
* of the requirements of the GNU Affero GPL version 3 as long as this Program
* is the same Program as licensed from Suomen Kanuuna Oy without any additional
* modifications.
*
* Supplemented terms under GNU Affero General Public License version 3
* section 7
*
* Origin of the software must be attributed to Suomen Kanuuna Oy. Any modified
* versions must be marked as "Modified version of" The Program.
*
* Names of the licensors and authors may not be used for publicity purposes.
*
* No rights are granted for use of trade names, trademarks, or service marks
* which are in The Program if any.
*
* Licensee must indemnify licensors and authors for any liability that these
* contractual assumptions impose on licensors and authors.
*
* To the extent this program is licensed as part of the Commercial versions of
* Teragrep, the applicable Commercial License may apply to this file if you as
* a licensee so wish it.
*/
package com.teragrep.rlp_03;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.ArrayDeque;
import java.util.Deque;
import com.teragrep.rlp_01.RelpFrameTX;
import com.teragrep.rlp_01.RelpParser;
import com.teragrep.rlp_01.TxID;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import tlschannel.NeedsReadException;
import tlschannel.NeedsWriteException;
/*
* Request reader class that reads incoming requests and sends them out for processing.
*/
class MessageReader {
private static final Logger LOGGER = LoggerFactory.getLogger(MessageReader.class);
private final RelpServerSocket relpServerSocket;
private final Deque<RelpFrameTX> txDeque;
private final ByteBuffer readBuffer;
private final FrameProcessor frameProcessor;
private final TxID txIdChecker = new TxID();
private RelpParser relpParser;
/**
* Constructor.
*/
MessageReader(RelpServerSocket relpServerSocket, Deque<RelpFrameTX> txDeque,
FrameProcessor frameProcessor) {
this.frameProcessor = frameProcessor;
this.relpServerSocket = relpServerSocket;
this.txDeque = txDeque;
this.readBuffer = ByteBuffer.allocateDirect(MAX_HEADER_CAPACITY + 1024*256);
this.relpParser = new RelpParser(false);
}
// Maximum capacity for HEADER part of RELP message frames:
// (MAX TXNR = 999999999).length + SP.length + (MAX COMMAND = "serverclose").length + SP.length + DATALEN.length + NL.length
private final int MAX_HEADER_CAPACITY = Integer.toString(txIdChecker.MAX_ID).length() +
" ".length() +
"serverclose".length() +
" ".length() +
Long.toString(Long.MAX_VALUE).length() +
" ".length();
/**
* Reads incoming requests from the associated RelpServerSocket, parses each incoming
* byte until there is a complete message, creates a frame for the parsed message, adds
* it to the to be processed RelpFrameServerRX queue and calls on frameProcessor to process it.
*
* @return READ state.
*/
ConnectionOperation readRequest() throws IOException {
LOGGER.trace("messageReader.readRequest> entry with parser: " + relpParser + " and parser state: " + relpParser.getState());
int readBytes = relpServerSocket.read(readBuffer);
while (readBytes > 0) {
readBuffer.flip(); // for reading
while (readBuffer.hasRemaining()) {
relpParser.parse(readBuffer.get());
if (relpParser.isComplete()) {
LOGGER.trace("messageReader.readRequest> read entire message complete ");
// TODO read long as we can to process batches
Deque<RelpFrameServerRX> rxFrames = new ArrayDeque<>();
RelpFrameServerRX rxFrame = new RelpFrameServerRX(
relpParser.getTxnId(),
relpParser.getCommandString(),
relpParser.getLength(),
relpParser.getData(),
relpServerSocket.getTransportInfo()
);
rxFrames.addLast(rxFrame);
txDeque.addAll(frameProcessor.process(rxFrames));
// reset parser state, TODO improve performance by having clear
relpParser = new RelpParser(false);
}
}
readBuffer.compact();
readBuffer.flip(); // for writing
try {
// read until there is no more data available
readBytes = relpServerSocket.read(readBuffer);
}
catch (NeedsReadException | NeedsWriteException tlsException) {
break;
}
}
if (readBytes < 0) {
// problem with socket, closing
LOGGER.trace("messageReader.readRequest> closing. " +
"readBytes: " + readBytes);
return ConnectionOperation.CLOSE;
}
else {
LOGGER.trace("messageReader.readRequest> exit with readBuffer: " + readBuffer);
return ConnectionOperation.READ;
}
}
}