stream-transport-impl.hpp
Go to the documentation of this file.
1 /* -*- Mode:C++; c-file-style:"gnu"; indent-tabs-mode:nil; -*- */
2 /*
3  * Copyright (c) 2013-2019 Regents of the University of California.
4  *
5  * This file is part of ndn-cxx library (NDN C++ library with eXperimental eXtensions).
6  *
7  * ndn-cxx library is free software: you can redistribute it and/or modify it under the
8  * terms of the GNU Lesser General Public License as published by the Free Software
9  * Foundation, either version 3 of the License, or (at your option) any later version.
10  *
11  * ndn-cxx library is distributed in the hope that it will be useful, but WITHOUT ANY
12  * WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR A
13  * PARTICULAR PURPOSE. See the GNU Lesser General Public License for more details.
14  *
15  * You should have received copies of the GNU General Public License and GNU Lesser
16  * General Public License along with ndn-cxx, e.g., in COPYING.md file. If not, see
17  * <http://www.gnu.org/licenses/>.
18  *
19  * See AUTHORS.md for complete list of ndn-cxx authors and contributors.
20  */
21 
22 #ifndef NDN_TRANSPORT_DETAIL_STREAM_TRANSPORT_IMPL_HPP
23 #define NDN_TRANSPORT_DETAIL_STREAM_TRANSPORT_IMPL_HPP
24 
26 
27 #include <boost/asio/steady_timer.hpp>
28 #include <boost/asio/write.hpp>
29 
30 #include <list>
31 
32 namespace ndn {
33 namespace detail {
34 
40 template<typename BaseTransport, typename Protocol>
41 class StreamTransportImpl : public std::enable_shared_from_this<StreamTransportImpl<BaseTransport, Protocol>>
42 {
43 public:
45  using BlockSequence = std::list<Block>;
46  using TransmissionQueue = std::list<BlockSequence>;
47 
48  StreamTransportImpl(BaseTransport& transport, boost::asio::io_service& ioService)
49  : m_transport(transport)
50  , m_socket(ioService)
51  , m_connectTimer(ioService)
52  {
53  }
54 
55  void
56  connect(const typename Protocol::endpoint& endpoint)
57  {
58  if (m_isConnecting) {
59  return;
60  }
61  m_isConnecting = true;
62 
63  // Wait at most 4 seconds to connect
65  m_connectTimer.expires_from_now(std::chrono::seconds(4));
66  m_connectTimer.async_wait([self = this->shared_from_this()] (const auto& error) {
67  self->connectTimeoutHandler(error);
68  });
69 
70  m_socket.open();
71  m_socket.async_connect(endpoint, [self = this->shared_from_this()] (const auto& error) {
72  self->connectHandler(error);
73  });
74  }
75 
76  void
78  {
79  m_isConnecting = false;
80 
81  boost::system::error_code error; // to silently ignore all errors
82  m_connectTimer.cancel(error);
83  m_socket.cancel(error);
84  m_socket.close(error);
85 
86  m_transport.m_isConnected = false;
87  m_transport.m_isReceiving = false;
88  m_transmissionQueue.clear();
89  }
90 
91  void
93  {
94  if (m_isConnecting)
95  return;
96 
97  if (m_transport.m_isReceiving) {
98  m_transport.m_isReceiving = false;
99  m_socket.cancel();
100  }
101  }
102 
103  void
105  {
106  if (m_isConnecting)
107  return;
108 
109  if (!m_transport.m_isReceiving) {
110  m_transport.m_isReceiving = true;
111  m_inputBufferSize = 0;
112  asyncReceive();
113  }
114  }
115 
116  void
117  send(const Block& wire)
118  {
119  BlockSequence sequence;
120  sequence.push_back(wire);
121  send(std::move(sequence));
122  }
123 
124  void
125  send(const Block& header, const Block& payload)
126  {
127  BlockSequence sequence;
128  sequence.push_back(header);
129  sequence.push_back(payload);
130  send(std::move(sequence));
131  }
132 
133 protected:
134  void
135  connectHandler(const boost::system::error_code& error)
136  {
137  m_isConnecting = false;
138  m_connectTimer.cancel();
139 
140  if (error) {
141  m_transport.m_isConnected = false;
142  m_transport.close();
143  NDN_THROW(Transport::Error(error, "error while connecting to the forwarder"));
144  }
145 
146  m_transport.m_isConnected = true;
147 
148  if (!m_transmissionQueue.empty()) {
149  resume();
150  asyncWrite();
151  }
152  }
153 
154  void
155  connectTimeoutHandler(const boost::system::error_code& error)
156  {
157  if (error) // e.g., cancelled timer
158  return;
159 
160  m_transport.close();
161  NDN_THROW(Transport::Error(error, "error while connecting to the forwarder"));
162  }
163 
164  void
165  send(BlockSequence&& sequence)
166  {
167  m_transmissionQueue.emplace_back(sequence);
168 
169  if (m_transport.m_isConnected && m_transmissionQueue.size() == 1) {
170  asyncWrite();
171  }
172 
173  // if not connected or there is transmission in progress (m_transmissionQueue.size() > 1),
174  // next write will be scheduled either in connectHandler or in asyncWriteHandler
175  }
176 
177  void
179  {
180  BOOST_ASSERT(!m_transmissionQueue.empty());
181  boost::asio::async_write(m_socket, m_transmissionQueue.front(),
182  bind(&Impl::handleAsyncWrite, this->shared_from_this(), _1,
183  m_transmissionQueue.begin()));
184  }
185 
186  void
187  handleAsyncWrite(const boost::system::error_code& error, TransmissionQueue::iterator queueItem)
188  {
189  if (error) {
190  if (error == boost::system::errc::operation_canceled) {
191  // async receive has been explicitly cancelled (e.g., socket close)
192  return;
193  }
194  m_transport.close();
195  NDN_THROW(Transport::Error(error, "error while writing data to socket"));
196  }
197 
198  if (!m_transport.m_isConnected) {
199  return; // queue has been already cleared
200  }
201 
202  m_transmissionQueue.erase(queueItem);
203 
204  if (!m_transmissionQueue.empty()) {
205  asyncWrite();
206  }
207  }
208 
209  void
211  {
212  m_socket.async_receive(boost::asio::buffer(m_inputBuffer + m_inputBufferSize,
214  bind(&Impl::handleAsyncReceive, this->shared_from_this(), _1, _2));
215  }
216 
217  void
218  handleAsyncReceive(const boost::system::error_code& error, std::size_t nBytesRecvd)
219  {
220  if (error) {
221  if (error == boost::system::errc::operation_canceled) {
222  // async receive has been explicitly cancelled (e.g., socket close)
223  return;
224  }
225  m_transport.close();
226  NDN_THROW(Transport::Error(error, "error while receiving data from socket"));
227  }
228 
229  m_inputBufferSize += nBytesRecvd;
230  // do magic
231 
232  std::size_t offset = 0;
233  bool hasProcessedSome = processAllReceived(m_inputBuffer, offset, m_inputBufferSize);
234  if (!hasProcessedSome && m_inputBufferSize == MAX_NDN_PACKET_SIZE && offset == 0) {
235  m_transport.close();
236  NDN_THROW(Transport::Error("input buffer full, but a valid TLV cannot be decoded"));
237  }
238 
239  if (offset > 0) {
240  if (offset != m_inputBufferSize) {
242  m_inputBufferSize -= offset;
243  }
244  else {
245  m_inputBufferSize = 0;
246  }
247  }
248 
249  asyncReceive();
250  }
251 
252  bool
253  processAllReceived(uint8_t* buffer, size_t& offset, size_t nBytesAvailable)
254  {
255  while (offset < nBytesAvailable) {
256  bool isOk = false;
257  Block element;
258  std::tie(isOk, element) = Block::fromBuffer(buffer + offset, nBytesAvailable - offset);
259  if (!isOk)
260  return false;
261 
262  m_transport.m_receiveCallback(element);
263  offset += element.size();
264  }
265  return true;
266  }
267 
268 protected:
269  BaseTransport& m_transport;
270 
271  typename Protocol::socket m_socket;
273  size_t m_inputBufferSize = 0;
274 
276  boost::asio::steady_timer m_connectTimer;
277  bool m_isConnecting = false;
278 };
279 
280 } // namespace detail
281 } // namespace ndn
282 
283 #endif // NDN_TRANSPORT_DETAIL_STREAM_TRANSPORT_IMPL_HPP
void connectHandler(const boost::system::error_code &error)
Definition: data.cpp:26
static std::tuple< bool, Block > fromBuffer(ConstBufferPtr buffer, size_t offset)
Try to parse Block from a wire buffer.
Definition: block.cpp:194
void handleAsyncWrite(const boost::system::error_code &error, TransmissionQueue::iterator queueItem)
StreamTransportImpl(BaseTransport &transport, boost::asio::io_service &ioService)
void send(const Block &header, const Block &payload)
Represents a TLV element of NDN packet format.
Definition: block.hpp:42
Implementation detail of a Boost.Asio-based stream-oriented transport.
#define NDN_THROW(e)
Definition: exception.hpp:61
void send(BlockSequence &&sequence)
size_t size() const
Return the size of the encoded wire, i.e.
Definition: block.cpp:290
boost::asio::steady_timer m_connectTimer
void handleAsyncReceive(const boost::system::error_code &error, std::size_t nBytesRecvd)
void connectTimeoutHandler(const boost::system::error_code &error)
std::list< BlockSequence > TransmissionQueue
uint8_t m_inputBuffer[MAX_NDN_PACKET_SIZE]
bool processAllReceived(uint8_t *buffer, size_t &offset, size_t nBytesAvailable)
const size_t MAX_NDN_PACKET_SIZE
practical limit of network layer packet size
Definition: tlv.hpp:41
void connect(const typename Protocol::endpoint &endpoint)