xrpld
Loading...
Searching...
No Matches
short_read_test.cpp
1#include <test/jtx/envconfig.h>
2
3#include <xrpl/basics/make_SSLContext.h>
4#include <xrpl/beast/core/CurrentThreadName.h>
5#include <xrpl/beast/unit_test/suite.h>
6
7#include <boost/asio/basic_waitable_timer.hpp>
8#include <boost/asio/bind_executor.hpp>
9#include <boost/asio/buffer.hpp>
10#include <boost/asio/error.hpp>
11#include <boost/asio/executor_work_guard.hpp>
12#include <boost/asio/io_context.hpp>
13#include <boost/asio/ip/address.hpp>
14#include <boost/asio/ip/tcp.hpp>
15#include <boost/asio/post.hpp>
16#include <boost/asio/read_until.hpp>
17#include <boost/asio/ssl/context.hpp>
18#include <boost/asio/ssl/stream.hpp>
19#include <boost/asio/strand.hpp>
20#include <boost/asio/streambuf.hpp>
21#include <boost/asio/write.hpp>
22#include <boost/optional/optional.hpp> // IWYU pragma: keep
23#include <boost/system/detail/error_code.hpp>
24
25#include <cassert>
26#include <chrono>
27#include <condition_variable>
28#include <cstddef>
29#include <map>
30#include <memory>
31#include <mutex>
32#include <ostream>
33#include <string>
34#include <thread>
35#include <utility>
36#include <vector>
37
38namespace xrpl {
39/*
40
41Findings from the test:
42
43If the remote host calls async_shutdown then the local host's
44async_read will complete with eof.
45
46If both hosts call async_shutdown then the calls to async_shutdown
47will complete with eof.
48
49*/
50
52{
53private:
54 using IoContextType = boost::asio::io_context;
55 using StrandType = boost::asio::strand<IoContextType::executor_type>;
56 using TimerType = boost::asio::basic_waitable_timer<std::chrono::steady_clock>;
57 using AcceptorType = boost::asio::ip::tcp::acceptor;
58 using SocketType = boost::asio::ip::tcp::socket;
59 using StreamType = boost::asio::ssl::stream<SocketType&>;
60 using ErrorCode = boost::system::error_code;
61 using EndpointType = boost::asio::ip::tcp::endpoint;
62 using AddressType = boost::asio::ip::address;
63
65 boost::optional<boost::asio::executor_work_guard<boost::asio::io_context::executor_type>> work_;
68
69 template <class Streambuf>
70 static void
71 write(Streambuf& sb, std::string const& s)
72 {
73 using boost::asio::buffer;
74 using boost::asio::buffer_copy;
75 using boost::asio::buffer_size;
76 boost::asio::const_buffer const buf(s.data(), s.size());
77 sb.commit(buffer_copy(sb.prepare(buffer_size(buf)), buf));
78 }
79
80 //--------------------------------------------------------------------------
81
82 class Base
83 {
84 protected:
85 class Child
86 {
87 private:
89
90 public:
91 explicit Child(Base& base) : base_(base)
92 {
93 }
94
95 virtual ~Child()
96 {
97 base_.remove(this);
98 }
99
100 virtual void
101 close() = 0;
102 };
103
104 private:
108 bool closed_ = false;
109
110 public:
112 {
113 // Derived class must call wait() in the destructor
114 assert(list_.empty());
115 }
116
117 void
119 {
120 std::scoped_lock const lock(mutex_);
121 list_.emplace(child.get(), child);
122 }
123
124 void
125 remove(Child* child)
126 {
127 std::scoped_lock const lock(mutex_);
128 list_.erase(child);
129 if (list_.empty())
130 cond_.notify_one();
131 }
132
133 void
135 {
137 {
138 std::scoped_lock const lock(mutex_);
139 v.reserve(list_.size());
140 if (closed_)
141 return;
142 closed_ = true;
143 for (auto const& c : list_)
144 {
145 if (auto p = c.second.lock())
146 {
147 p->close();
148 // Must destroy shared_ptr outside the
149 // lock otherwise deadlock from the
150 // managed object's destructor.
151 v.emplace_back(std::move(p));
152 }
153 }
154 }
155 }
156
157 void
159 {
161 while (!list_.empty())
162 cond_.wait(lock);
163 }
164 };
165
166 //--------------------------------------------------------------------------
167
168 class Server : public Base
169 {
170 private:
173
175 {
181
183 : Child(server)
184 , server(server)
185 , test(server.test_)
186 , acceptor(
188 EndpointType(boost::asio::ip::make_address(test::getEnvLocalhostAddr()), 0))
190 , strand(boost::asio::make_strand(test.ioContext_))
191 {
192 acceptor.listen();
193 server.endpoint_ = acceptor.local_endpoint();
194 }
195
196 void
197 close() override
198 {
199 if (!strand.running_in_this_thread())
200 {
201 post(strand, [self = shared_from_this()] { self->close(); });
202 return;
203 }
204 acceptor.close();
205 }
206
207 void
209 {
210 acceptor.async_accept(
211 socket, bind_executor(strand, [self = shared_from_this()](ErrorCode const& ec) {
212 self->onAccept(ec);
213 }));
214 }
215
216 void
217 fail(std::string const& what, ErrorCode ec)
218 {
219 if (acceptor.is_open())
220 {
221 if (ec != boost::asio::error::operation_aborted)
222 test.log << what << ": " << ec.message() << std::endl;
223 acceptor.close();
224 }
225 }
226
227 void
229 {
230 if (ec)
231 {
232 fail("accept", ec);
233 return;
234 }
235 auto const p = std::make_shared<Connection>(server, std::move(socket));
236 server.add(p);
237 p->run();
238 acceptor.async_accept(
239 socket, bind_executor(strand, [self = shared_from_this()](ErrorCode const& ec) {
240 self->onAccept(ec);
241 }));
242 }
243 };
244
246 {
253 boost::asio::streambuf buf;
254
255 Connection(Server& inServer, SocketType&& inSocket)
256 : Child(inServer)
257 , server(inServer)
258 , test(server.test_)
259 , socket(std::move(inSocket))
261 , strand(boost::asio::make_strand(test.ioContext_))
263 {
264 }
265
266 void
267 close() override
268 {
269 if (!strand.running_in_this_thread())
270 {
271 post(strand, [self = shared_from_this()] { self->close(); });
272 return;
273 }
274 if (socket.is_open())
275 {
276 socket.close();
277 timer.cancel();
278 }
279 }
280
281 void
283 {
284 timer.expires_after(std::chrono::seconds(3));
285 timer.async_wait(bind_executor(
286 strand,
287 [self = shared_from_this()](ErrorCode const& ec) { self->onTimer(ec); }));
288 stream.async_handshake(
289 StreamType::server,
290 bind_executor(strand, [self = shared_from_this()](ErrorCode const& ec) {
291 self->onHandshake(ec);
292 }));
293 }
294
295 void
296 fail(std::string const& what, ErrorCode ec)
297 {
298 if (socket.is_open())
299 {
300 if (ec != boost::asio::error::operation_aborted)
301 test.log << "[server] " << what << ": " << ec.message() << std::endl;
302 socket.close();
303 timer.cancel();
304 }
305 }
306
307 void
309 {
310 if (ec == boost::asio::error::operation_aborted)
311 return;
312 if (ec)
313 {
314 fail("timer", ec);
315 return;
316 }
317 test.log << "[server] timeout" << std::endl;
318 socket.close();
319 }
320
321 void
323 {
324 if (ec)
325 {
326 fail("handshake", ec);
327 return;
328 }
329 boost::asio::async_read_until(
330 stream,
331 buf,
332 "\n",
333 bind_executor(
334 strand,
335 [self = shared_from_this()](
336 ErrorCode const& ec, std::size_t bytesTransferred) {
337 self->onRead(ec, bytesTransferred);
338 }));
339 }
340
341 void
342 onRead(ErrorCode ec, std::size_t bytesTransferred)
343 {
344 if (ec == boost::asio::error::eof)
345 {
346 server.test_.log << "[server] read: EOF" << std::endl;
347 stream.async_shutdown(
348 bind_executor(strand, [self = shared_from_this()](ErrorCode const& ec) {
349 self->onShutdown(ec);
350 }));
351 return;
352 }
353 if (ec)
354 {
355 fail("read", ec);
356 return;
357 }
358
359 buf.commit(bytesTransferred);
360 buf.consume(bytesTransferred);
361 write(buf, "BYE\n");
362 boost::asio::async_write(
363 stream,
364 buf.data(),
365 bind_executor(
366 strand,
367 [self = shared_from_this()](
368 ErrorCode const& ec, std::size_t bytesTransferred) {
369 self->onWrite(ec, bytesTransferred);
370 }));
371 }
372
373 void
374 onWrite(ErrorCode ec, std::size_t bytesTransferred)
375 {
376 buf.consume(bytesTransferred);
377 if (ec)
378 {
379 fail("write", ec);
380 return;
381 }
382 stream.async_shutdown(bind_executor(
383 strand,
384 [self = shared_from_this()](ErrorCode const& ec) { self->onShutdown(ec); }));
385 }
386
387 void
389 {
390 if (ec)
391 {
392 fail("shutdown", ec);
393 return;
394 }
395 socket.close();
396 timer.cancel();
397 }
398 };
399
400 public:
402 {
403 auto const p = std::make_shared<Acceptor>(*this);
404 add(p);
405 p->run();
406 }
407
409 {
410 close();
411 wait();
412 }
413
414 [[nodiscard]] EndpointType const&
415 endpoint() const
416 {
417 return endpoint_;
418 }
419 };
420
421 //--------------------------------------------------------------------------
422 class Client : public Base
423
424 {
425 private:
427
429 {
436 boost::asio::streambuf buf;
438
440 : Child(client)
441 , client(client)
442 , test(client.test_)
445 , strand(boost::asio::make_strand(test.ioContext_))
447 , ep(ep)
448 {
449 }
450
451 void
452 close() override
453 {
454 if (!strand.running_in_this_thread())
455 {
456 post(strand, [self = shared_from_this()] { self->close(); });
457 return;
458 }
459 if (socket.is_open())
460 {
461 socket.close();
462 timer.cancel();
463 }
464 }
465
466 void
468 {
469 timer.expires_after(std::chrono::seconds(3));
470 timer.async_wait(bind_executor(
471 strand,
472 [self = shared_from_this()](ErrorCode const& ec) { self->onTimer(ec); }));
473 socket.async_connect(
474 ep, bind_executor(strand, [self = shared_from_this()](ErrorCode const& ec) {
475 self->onConnect(ec);
476 }));
477 }
478
479 void
480 fail(std::string const& what, ErrorCode ec)
481 {
482 if (socket.is_open())
483 {
484 if (ec != boost::asio::error::operation_aborted)
485 test.log << "[client] " << what << ": " << ec.message() << std::endl;
486 socket.close();
487 timer.cancel();
488 }
489 }
490
491 void
493 {
494 if (ec == boost::asio::error::operation_aborted)
495 return;
496 if (ec)
497 {
498 fail("timer", ec);
499 return;
500 }
501 test.log << "[client] timeout";
502 socket.close();
503 }
504
505 void
507 {
508 if (ec)
509 {
510 fail("connect", ec);
511 return;
512 }
513 stream.async_handshake(
514 StreamType::client,
515 bind_executor(strand, [self = shared_from_this()](ErrorCode const& ec) {
516 self->onHandshake(ec);
517 }));
518 }
519
520 void
522 {
523 if (ec)
524 {
525 fail("handshake", ec);
526 return;
527 }
528 write(buf, "HELLO\n");
529
530 boost::asio::async_write(
531 stream,
532 buf.data(),
533 bind_executor(
534 strand,
535 [self = shared_from_this()](
536 ErrorCode const& ec, std::size_t bytesTransferred) {
537 self->onWrite(ec, bytesTransferred);
538 }));
539 }
540
541 void
542 onWrite(ErrorCode ec, std::size_t bytesTransferred)
543 {
544 buf.consume(bytesTransferred);
545 if (ec)
546 {
547 fail("write", ec);
548 return;
549 }
550 boost::asio::async_read_until(
551 stream,
552 buf,
553 "\n",
554 bind_executor(
555 strand,
556 [self = shared_from_this()](
557 ErrorCode const& ec, std::size_t bytesTransferred) {
558 self->onRead(ec, bytesTransferred);
559 }));
560 }
561
562 void
563 onRead(ErrorCode ec, std::size_t bytesTransferred)
564 {
565 if (ec)
566 {
567 fail("read", ec);
568 return;
569 }
570 buf.commit(bytesTransferred);
571 stream.async_shutdown(bind_executor(
572 strand,
573 [self = shared_from_this()](ErrorCode const& ec) { self->onShutdown(ec); }));
574 }
575
576 void
578 {
579 if (ec)
580 {
581 fail("shutdown", ec);
582 return;
583 }
584 socket.close();
585 timer.cancel();
586 }
587 };
588
589 public:
591 {
592 auto const p = std::make_shared<Connection>(*this, ep);
593 add(p);
594 p->run(ep);
595 }
596
598 {
599 close();
600 wait();
601 }
602 };
603
604public:
606 : work_(ioContext_.get_executor())
607 , thread_(std::thread([this]() {
608 beast::setCurrentThreadName("io_context");
609 this->ioContext_.run();
610 }))
612 {
613 }
614
616 {
617 work_.reset();
618 thread_.join();
619 }
620
621 void
622 run() override
623 {
624 Server const s(*this);
625 Client c(*this, s.endpoint());
626 c.wait();
627 pass();
628 }
629};
630
631BEAST_DEFINE_TESTSUITE(short_read, overlay, xrpl);
632
633} // namespace xrpl
A testsuite class.
Definition suite.h:52
void pass()
Record a successful test condition.
Definition suite.h:532
std::map< Child *, std::weak_ptr< Child > > list_
std::condition_variable cond_
void add(std::shared_ptr< Child > const &child)
Client(short_read_test &test, EndpointType const &ep)
Server(short_read_test &test)
EndpointType const & endpoint() const
void run() override
Runs the suite.
boost::asio::ip::tcp::endpoint EndpointType
boost::asio::ip::address AddressType
boost::asio::basic_waitable_timer< std::chrono::steady_clock > TimerType
boost::system::error_code ErrorCode
std::shared_ptr< boost::asio::ssl::context > context_
boost::asio::ip::tcp::acceptor AcceptorType
boost::asio::ip::tcp::socket SocketType
boost::asio::strand< IoContextType::executor_type > StrandType
static void write(Streambuf &sb, std::string const &s)
boost::asio::io_context IoContextType
boost::optional< boost::asio::executor_work_guard< boost::asio::io_context::executor_type > > work_
boost::asio::ssl::stream< SocketType & > StreamType
T data(T... args)
T emplace_back(T... args)
T endl(T... args)
T make_shared(T... args)
void setCurrentThreadName(std::string_view newThreadName)
Changes the name of the caller thread.
STL namespace.
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
Definition algorithm.h:5
std::shared_ptr< boost::asio::ssl::context > makeSslContext(std::string const &cipherList)
Create a self-signed SSL context that allows anonymous Diffie Hellman.
BEAST_DEFINE_TESTSUITE(AccountTxPaging, app, xrpl)
T reserve(T... args)
T size(T... args)
void fail(std::string const &what, ErrorCode ec)
void onRead(ErrorCode ec, std::size_t bytesTransferred)
Connection(Client &client, EndpointType const &ep)
void onWrite(ErrorCode ec, std::size_t bytesTransferred)
void fail(std::string const &what, ErrorCode ec)
void fail(std::string const &what, ErrorCode ec)
Connection(Server &inServer, SocketType &&inSocket)
void onWrite(ErrorCode ec, std::size_t bytesTransferred)
void onRead(ErrorCode ec, std::size_t bytesTransferred)