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 m->DiscardUnknownFields();
282
283 return m;
284}
285
286template <class T, class Buffers, class Handler>
287bool
288invoke(MessageHeader const& header, Buffers const& buffers, Handler& handler)
290{
291 auto const m = parseMessageContent<T>(header, buffers);
292 if (!m)
293 return false;
294
295 using namespace xrpl::compression;
296 handler.onMessageBegin(
297 header.messageType,
298 m,
299 header.payloadWireSize,
300 header.uncompressedSize,
301 header.algorithm != Algorithm::None);
302 handler.onMessage(m);
303 handler.onMessageEnd(header.messageType, m);
304
305 return true;
306}
307
308} // namespace detail
309
323template <class Buffers, class Handler>
325invokeProtocolMessage(Buffers const& buffers, Handler& handler, std::size_t& hint)
326{
328
329 auto const size = boost::asio::buffer_size(buffers);
330
331 if (size == 0)
332 return result;
333
334 auto header = detail::parseMessageHeader(result.second, buffers, size);
335
336 // If we can't parse the header then it may be that we don't have enough
337 // bytes yet, or because the message was cut off (if error_code is success).
338 // Otherwise we failed to match the header's marker (error_code is set to
339 // no_message) or the compression algorithm is invalid (error_code is
340 // protocol_error) and signal an error.
341 if (!header)
342 return result;
343
344 // We implement a maximum size for protocol messages. Sending a message
345 // whose size exceeds this may result in the connection being dropped. A
346 // larger message size may be supported in the future or negotiated as
347 // part of a protocol upgrade.
348 if (header->payloadWireSize > kMaximumMessageSize ||
349 header->uncompressedSize > kMaximumMessageSize)
350 {
351 result.second = make_error_code(boost::system::errc::message_size);
352 return result;
353 }
354
355 // We requested uncompressed messages from the peer but received compressed.
356 if (!handler.compressionEnabled() && header->algorithm != compression::Algorithm::None)
357 {
358 result.second = make_error_code(boost::system::errc::protocol_error);
359 return result;
360 }
361
362 if (header->messageType == protocol::mtPING &&
363 header->uncompressedSize + header->headerSize > kMaximumPingMessageSize)
364 {
365 result.second = make_error_code(boost::system::errc::message_size);
366 return result;
367 }
368
369 // We don't have the whole message yet. This isn't an error but we have
370 // nothing to do.
371 if (header->totalWireSize > size)
372 {
373 hint = header->totalWireSize - size;
374 return result;
375 }
376
377 // Drop an oversized TMManifests without penalty: consume the bytes and
378 // return no error, so the connection is preserved. The limit follows this
379 // node's configured manifests-per-message count.
380 if (header->messageType == protocol::mtMANIFESTS)
381 {
382 auto const maxSize = handler.maxManifestsMessageSize();
383 if (header->payloadWireSize > maxSize || header->uncompressedSize > maxSize)
384 {
385 result.first = header->totalWireSize;
386 return result;
387 }
388 }
389
390 bool success = false;
391
392 switch (header->messageType)
393 {
394 case protocol::mtMANIFESTS:
395 success = detail::invoke<protocol::TMManifests>(*header, buffers, handler);
396 break;
397 case protocol::mtPING:
398 success = detail::invoke<protocol::TMPing>(*header, buffers, handler);
399 break;
400 case protocol::mtCLUSTER:
401 success = detail::invoke<protocol::TMCluster>(*header, buffers, handler);
402 break;
403 case protocol::mtENDPOINTS:
404 success = detail::invoke<protocol::TMEndpoints>(*header, buffers, handler);
405 break;
406 case protocol::mtTRANSACTION:
407 success = detail::invoke<protocol::TMTransaction>(*header, buffers, handler);
408 break;
409 case protocol::mtGET_LEDGER:
410 success = detail::invoke<protocol::TMGetLedger>(*header, buffers, handler);
411 break;
412 case protocol::mtLEDGER_DATA:
413 success = detail::invoke<protocol::TMLedgerData>(*header, buffers, handler);
414 break;
415 case protocol::mtPROPOSE_LEDGER:
416 success = detail::invoke<protocol::TMProposeSet>(*header, buffers, handler);
417 break;
418 case protocol::mtSTATUS_CHANGE:
419 success = detail::invoke<protocol::TMStatusChange>(*header, buffers, handler);
420 break;
421 case protocol::mtHAVE_SET:
422 success = detail::invoke<protocol::TMHaveTransactionSet>(*header, buffers, handler);
423 break;
424 case protocol::mtVALIDATION:
425 success = detail::invoke<protocol::TMValidation>(*header, buffers, handler);
426 break;
427 case protocol::mtVALIDATOR_LIST_COLLECTION:
428 success =
430 break;
431 case protocol::mtGET_OBJECTS:
432 success = detail::invoke<protocol::TMGetObjectByHash>(*header, buffers, handler);
433 break;
434 case protocol::mtHAVE_TRANSACTIONS:
435 success = detail::invoke<protocol::TMHaveTransactions>(*header, buffers, handler);
436 break;
437 case protocol::mtTRANSACTIONS:
438 success = detail::invoke<protocol::TMTransactions>(*header, buffers, handler);
439 break;
440 case protocol::mtSQUELCH:
441 success = detail::invoke<protocol::TMSquelch>(*header, buffers, handler);
442 break;
443 case protocol::mtPROOF_PATH_REQ:
444 success = detail::invoke<protocol::TMProofPathRequest>(*header, buffers, handler);
445 break;
446 case protocol::mtPROOF_PATH_RESPONSE:
447 success = detail::invoke<protocol::TMProofPathResponse>(*header, buffers, handler);
448 break;
449 case protocol::mtREPLAY_DELTA_REQ:
450 success = detail::invoke<protocol::TMReplayDeltaRequest>(*header, buffers, handler);
451 break;
452 case protocol::mtREPLAY_DELTA_RESPONSE:
453 success = detail::invoke<protocol::TMReplayDeltaResponse>(*header, buffers, handler);
454 break;
455 default:
456 handler.onMessageUnknown(header->messageType);
457 success = true;
458 break;
459 }
460
461 result.first = header->totalWireSize;
462
463 if (!success)
464 result.second = make_error_code(boost::system::errc::bad_message);
465
466 return result;
467}
468
469} // 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.