xrpld
Loading...
Searching...
No Matches
ProtocolMessage.h
1#pragma once
2
3#include <xrpld/overlay/Compression.h>
4#include <xrpld/overlay/Message.h>
5#include <xrpld/overlay/detail/ZeroCopyStream.h>
6
7#include <xrpl/beast/utility/instrumentation.h>
8
9#include <boost/asio/buffer.hpp>
10#include <boost/asio/buffers_iterator.hpp>
11
12#include <google/protobuf/message.h>
13
14#include <xrpl.pb.h>
15
16#include <cstddef>
17#include <cstdint>
18#include <memory>
19#include <optional>
20#include <string>
21#include <type_traits>
22#include <utility>
23#include <vector>
24
25namespace xrpl {
26
27inline protocol::MessageType
28protocolMessageType(protocol::TMGetLedger const&)
29{
30 return protocol::mtGET_LEDGER;
31}
32
33inline protocol::MessageType
34protocolMessageType(protocol::TMReplayDeltaRequest const&)
35{
36 return protocol::mtREPLAY_DELTA_REQ;
37}
38
39inline protocol::MessageType
40protocolMessageType(protocol::TMProofPathRequest const&)
41{
42 return protocol::mtPROOF_PATH_REQ;
43}
44
48template <class = void>
51{
52 switch (type)
53 {
54 case protocol::mtMANIFESTS:
55 return "manifests";
56 case protocol::mtPING:
57 return "ping";
58 case protocol::mtCLUSTER:
59 return "cluster";
60 case protocol::mtENDPOINTS:
61 return "endpoints";
62 case protocol::mtTRANSACTION:
63 return "tx";
64 case protocol::mtGET_LEDGER:
65 return "get_ledger";
66 case protocol::mtLEDGER_DATA:
67 return "ledger_data";
68 case protocol::mtPROPOSE_LEDGER:
69 return "propose";
70 case protocol::mtSTATUS_CHANGE:
71 return "status";
72 case protocol::mtHAVE_SET:
73 return "have_set";
74 case protocol::mtVALIDATOR_LIST_COLLECTION:
75 return "validator_list_collection";
76 case protocol::mtVALIDATION:
77 return "validation";
78 case protocol::mtGET_OBJECTS:
79 return "get_objects";
80 case protocol::mtHAVE_TRANSACTIONS:
81 return "have_transactions";
82 case protocol::mtTRANSACTIONS:
83 return "transactions";
84 case protocol::mtSQUELCH:
85 return "squelch";
86 case protocol::mtPROOF_PATH_REQ:
87 return "proof_path_request";
88 case protocol::mtPROOF_PATH_RESPONSE:
89 return "proof_path_response";
90 case protocol::mtREPLAY_DELTA_REQ:
91 return "replay_delta_request";
92 case protocol::mtREPLAY_DELTA_RESPONSE:
93 return "replay_delta_response";
94 default:
95 break;
96 }
97 return "unknown";
98}
99
100namespace detail {
101
138
139template <typename BufferSequence>
140auto
141buffersBegin(BufferSequence const& bufs)
142{
143 return boost::asio::buffers_iterator<BufferSequence, std::uint8_t>::begin(bufs);
144}
145
146template <typename BufferSequence>
147auto
148buffersEnd(BufferSequence const& bufs)
149{
150 return boost::asio::buffers_iterator<BufferSequence, std::uint8_t>::end(bufs);
151}
152
163template <class BufferSequence>
165parseMessageHeader(boost::system::error_code& ec, BufferSequence const& bufs, std::size_t size)
166{
167 using namespace xrpl::compression;
168
169 MessageHeader hdr;
170 auto iter = buffersBegin(bufs);
171 XRPL_ASSERT(iter != buffersEnd(bufs), "xrpl::detail::parseMessageHeader : non-empty buffer");
172
173 // Check valid header compressed message:
174 // - 4 bits are the compression algorithm, 1st bit is always set to 1
175 // - 2 bits are always set to 0
176 // - 26 bits are the payload size
177 // - 32 bits are the uncompressed data size
178 if (*iter & 0x80)
179 {
181
182 // not enough bytes to parse the header
183 if (size < hdr.headerSize)
184 {
185 ec = make_error_code(boost::system::errc::success);
186 return std::nullopt;
187 }
188
189 if (*iter & 0x0C)
190 {
191 ec = make_error_code(boost::system::errc::protocol_error);
192 return std::nullopt;
193 }
194
195 hdr.algorithm = static_cast<compression::Algorithm>(*iter & 0xF0);
196
198 {
199 ec = make_error_code(boost::system::errc::protocol_error);
200 return std::nullopt;
201 }
202
203 for (int i = 0; i != 4; ++i)
204 hdr.payloadWireSize = (hdr.payloadWireSize << 8) + *iter++;
205
206 // clear the top four bits (the compression bits).
207 hdr.payloadWireSize &= 0x0FFFFFFF;
208
210
211 for (int i = 0; i != 2; ++i)
212 hdr.messageType = (hdr.messageType << 8) + *iter++;
213
214 for (int i = 0; i != 4; ++i)
215 hdr.uncompressedSize = (hdr.uncompressedSize << 8) + *iter++;
216
217 return hdr;
218 }
219
220 // Check valid header uncompressed message:
221 // - 6 bits are set to 0
222 // - 26 bits are the payload size
223 if ((*iter & 0xFC) == 0)
224 {
226
227 if (size < hdr.headerSize)
228 {
229 ec = make_error_code(boost::system::errc::success);
230 return std::nullopt;
231 }
232
234
235 for (int i = 0; i != 4; ++i)
236 hdr.payloadWireSize = (hdr.payloadWireSize << 8) + *iter++;
237
240
241 for (int i = 0; i != 2; ++i)
242 hdr.messageType = (hdr.messageType << 8) + *iter++;
243
244 return hdr;
245 }
246
247 ec = make_error_code(boost::system::errc::no_message);
248 return std::nullopt;
249}
250
251template <class T, class Buffers>
253parseMessageContent(MessageHeader const& header, Buffers const& buffers)
255{
256 auto m = std::make_shared<T>();
257
258 ZeroCopyInputStream<Buffers> stream(buffers);
259 stream.Skip(header.headerSize);
260
261 if (header.algorithm != compression::Algorithm::None)
262 {
264 payload.resize(header.uncompressedSize);
265
266 auto const payloadSize = xrpl::compression::decompress(
267 stream,
268 header.payloadWireSize,
269 payload.data(),
270 header.uncompressedSize,
271 header.algorithm);
272
273 if (payloadSize == 0 || !m->ParseFromArray(payload.data(), payloadSize))
274 return {};
275 }
276 else if (!m->ParseFromZeroCopyStream(&stream))
277 {
278 return {};
279 }
280
281 return m;
282}
283
284template <class T, class Buffers, class Handler>
285bool
286invoke(MessageHeader const& header, Buffers const& buffers, Handler& handler)
288{
289 auto const m = parseMessageContent<T>(header, buffers);
290 if (!m)
291 return false;
292
293 using namespace xrpl::compression;
294 handler.onMessageBegin(
295 header.messageType,
296 m,
297 header.payloadWireSize,
298 header.uncompressedSize,
299 header.algorithm != Algorithm::None);
300 handler.onMessage(m);
301 handler.onMessageEnd(header.messageType, m);
302
303 return true;
304}
305
306} // namespace detail
307
321template <class Buffers, class Handler>
323invokeProtocolMessage(Buffers const& buffers, Handler& handler, std::size_t& hint)
324{
326
327 auto const size = boost::asio::buffer_size(buffers);
328
329 if (size == 0)
330 return result;
331
332 auto header = detail::parseMessageHeader(result.second, buffers, size);
333
334 // If we can't parse the header then it may be that we don't have enough
335 // bytes yet, or because the message was cut off (if error_code is success).
336 // Otherwise we failed to match the header's marker (error_code is set to
337 // no_message) or the compression algorithm is invalid (error_code is
338 // protocol_error) and signal an error.
339 if (!header)
340 return result;
341
342 // We implement a maximum size for protocol messages. Sending a message
343 // whose size exceeds this may result in the connection being dropped. A
344 // larger message size may be supported in the future or negotiated as
345 // part of a protocol upgrade.
346 if (header->payloadWireSize > kMaximumMessageSize ||
347 header->uncompressedSize > kMaximumMessageSize)
348 {
349 result.second = make_error_code(boost::system::errc::message_size);
350 return result;
351 }
352
353 // We requested uncompressed messages from the peer but received compressed.
354 if (!handler.compressionEnabled() && header->algorithm != compression::Algorithm::None)
355 {
356 result.second = make_error_code(boost::system::errc::protocol_error);
357 return result;
358 }
359
360 if (header->messageType == protocol::mtPING &&
361 header->uncompressedSize + header->headerSize > kMaximumPingMessageSize)
362 {
363 result.second = make_error_code(boost::system::errc::message_size);
364 return result;
365 }
366
367 // We don't have the whole message yet. This isn't an error but we have
368 // nothing to do.
369 if (header->totalWireSize > size)
370 {
371 hint = header->totalWireSize - size;
372 return result;
373 }
374
375 // Drop an oversized TMManifests without penalty: consume the bytes and
376 // return no error, so the connection is preserved. The limit follows this
377 // node's configured manifests-per-message count.
378 if (header->messageType == protocol::mtMANIFESTS)
379 {
380 auto const maxSize = handler.maxManifestsMessageSize();
381 if (header->payloadWireSize > maxSize || header->uncompressedSize > maxSize)
382 {
383 result.first = header->totalWireSize;
384 return result;
385 }
386 }
387
388 bool success = false;
389
390 switch (header->messageType)
391 {
392 case protocol::mtMANIFESTS:
393 success = detail::invoke<protocol::TMManifests>(*header, buffers, handler);
394 break;
395 case protocol::mtPING:
396 success = detail::invoke<protocol::TMPing>(*header, buffers, handler);
397 break;
398 case protocol::mtCLUSTER:
399 success = detail::invoke<protocol::TMCluster>(*header, buffers, handler);
400 break;
401 case protocol::mtENDPOINTS:
402 success = detail::invoke<protocol::TMEndpoints>(*header, buffers, handler);
403 break;
404 case protocol::mtTRANSACTION:
405 success = detail::invoke<protocol::TMTransaction>(*header, buffers, handler);
406 break;
407 case protocol::mtGET_LEDGER:
408 success = detail::invoke<protocol::TMGetLedger>(*header, buffers, handler);
409 break;
410 case protocol::mtLEDGER_DATA:
411 success = detail::invoke<protocol::TMLedgerData>(*header, buffers, handler);
412 break;
413 case protocol::mtPROPOSE_LEDGER:
414 success = detail::invoke<protocol::TMProposeSet>(*header, buffers, handler);
415 break;
416 case protocol::mtSTATUS_CHANGE:
417 success = detail::invoke<protocol::TMStatusChange>(*header, buffers, handler);
418 break;
419 case protocol::mtHAVE_SET:
420 success = detail::invoke<protocol::TMHaveTransactionSet>(*header, buffers, handler);
421 break;
422 case protocol::mtVALIDATION:
423 success = detail::invoke<protocol::TMValidation>(*header, buffers, handler);
424 break;
425 case protocol::mtVALIDATOR_LIST_COLLECTION:
426 success =
428 break;
429 case protocol::mtGET_OBJECTS:
430 success = detail::invoke<protocol::TMGetObjectByHash>(*header, buffers, handler);
431 break;
432 case protocol::mtHAVE_TRANSACTIONS:
433 success = detail::invoke<protocol::TMHaveTransactions>(*header, buffers, handler);
434 break;
435 case protocol::mtTRANSACTIONS:
436 success = detail::invoke<protocol::TMTransactions>(*header, buffers, handler);
437 break;
438 case protocol::mtSQUELCH:
439 success = detail::invoke<protocol::TMSquelch>(*header, buffers, handler);
440 break;
441 case protocol::mtPROOF_PATH_REQ:
442 success = detail::invoke<protocol::TMProofPathRequest>(*header, buffers, handler);
443 break;
444 case protocol::mtPROOF_PATH_RESPONSE:
445 success = detail::invoke<protocol::TMProofPathResponse>(*header, buffers, handler);
446 break;
447 case protocol::mtREPLAY_DELTA_REQ:
448 success = detail::invoke<protocol::TMReplayDeltaRequest>(*header, buffers, handler);
449 break;
450 case protocol::mtREPLAY_DELTA_RESPONSE:
451 success = detail::invoke<protocol::TMReplayDeltaResponse>(*header, buffers, handler);
452 break;
453 default:
454 handler.onMessageUnknown(header->messageType);
455 success = true;
456 break;
457 }
458
459 result.first = header->totalWireSize;
460
461 if (!success)
462 result.second = make_error_code(boost::system::errc::bad_message);
463
464 return result;
465}
466
467} // namespace xrpl
Implements ZeroCopyInputStream around a buffer sequence.
T data(T... args)
T is_base_of_v
T make_shared(T... args)
constexpr std::size_t kHeaderBytesCompressed
Definition Compression.h:13
constexpr std::size_t kHeaderBytes
Definition Compression.h:12
std::size_t decompress(InputStream &in, std::size_t inSize, std::uint8_t *decompressed, std::size_t decompressedSize, Algorithm algorithm=Algorithm::LZ4)
Decompress input stream.
Definition Compression.h:32
bool invoke(MessageHeader const &header, Buffers const &buffers, Handler &handler)
std::optional< MessageHeader > parseMessageHeader(boost::system::error_code &ec, BufferSequence const &bufs, std::size_t size)
Parse a message header.
auto buffersBegin(BufferSequence const &bufs)
std::shared_ptr< T > parseMessageContent(MessageHeader const &header, Buffers const &buffers)
auto buffersEnd(BufferSequence const &bufs)
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition algorithm.h:5
constexpr std::size_t kMaximumMessageSize
Definition Message.h:22
constexpr std::size_t kMaximumPingMessageSize
Definition Message.h:25
std::error_code make_error_code(xrpl::TokenCodecErrc e)
protocol::MessageType protocolMessageType(protocol::TMGetLedger const &)
std::pair< std::size_t, boost::system::error_code > invokeProtocolMessage(Buffers const &buffers, Handler &handler, std::size_t &hint)
Calls the handler for up to one protocol message in the passed buffers.
std::string protocolMessageName(int type)
Returns the name of a protocol message given its type.
T resize(T... args)
std::uint32_t uncompressedSize
Uncompressed message size if the message is compressed.
std::uint16_t messageType
The type of the message.
std::uint32_t headerSize
The size of the header associated with this message.
compression::Algorithm algorithm
Indicates which compression algorithm the payload is compressed with.
std::uint32_t totalWireSize
The size of the message on the wire.
std::uint32_t payloadWireSize
The size of the payload on the wire.