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-2018 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/deadline_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  typedef std::list<Block> BlockSequence;
46  typedef std::list<BlockSequence> TransmissionQueue;
47 
48  StreamTransportImpl(BaseTransport& transport, boost::asio::io_service& ioService)
49  : m_transport(transport)
50  , m_socket(ioService)
52  , m_isConnecting(false)
53  , m_connectTimer(ioService)
54  {
55  }
56 
57  void
58  connect(const typename Protocol::endpoint& endpoint)
59  {
60  if (!m_isConnecting) {
61  m_isConnecting = true;
62 
63  // Wait at most 4 seconds to connect
65  m_connectTimer.expires_from_now(boost::posix_time::seconds(4));
66  m_connectTimer.async_wait(bind(&Impl::connectTimeoutHandler, this->shared_from_this(), _1));
67 
68  m_socket.open();
69  m_socket.async_connect(endpoint, bind(&Impl::connectHandler, this->shared_from_this(), _1));
70  }
71  }
72 
73  void
75  {
76  m_isConnecting = false;
77 
78  boost::system::error_code error; // to silently ignore all errors
79  m_connectTimer.cancel(error);
80  m_socket.cancel(error);
81  m_socket.close(error);
82 
83  m_transport.m_isConnected = false;
84  m_transport.m_isReceiving = false;
85  m_transmissionQueue.clear();
86  }
87 
88  void
90  {
91  if (m_isConnecting)
92  return;
93 
94  if (m_transport.m_isReceiving) {
95  m_transport.m_isReceiving = false;
96  m_socket.cancel();
97  }
98  }
99 
100  void
102  {
103  if (m_isConnecting)
104  return;
105 
106  if (!m_transport.m_isReceiving) {
107  m_transport.m_isReceiving = true;
108  m_inputBufferSize = 0;
109  asyncReceive();
110  }
111  }
112 
113  void
114  send(const Block& wire)
115  {
116  BlockSequence sequence;
117  sequence.push_back(wire);
118  send(std::move(sequence));
119  }
120 
121  void
122  send(const Block& header, const Block& payload)
123  {
124  BlockSequence sequence;
125  sequence.push_back(header);
126  sequence.push_back(payload);
127  send(std::move(sequence));
128  }
129 
130 protected:
131  void
132  connectHandler(const boost::system::error_code& error)
133  {
134  m_isConnecting = false;
135  m_connectTimer.cancel();
136 
137  if (!error) {
138  m_transport.m_isConnected = true;
139 
140  if (!m_transmissionQueue.empty()) {
141  resume();
142  asyncWrite();
143  }
144  }
145  else {
146  m_transport.m_isConnected = false;
147  m_transport.close();
148  BOOST_THROW_EXCEPTION(Transport::Error(error, "error while connecting to the forwarder"));
149  }
150  }
151 
152  void
153  connectTimeoutHandler(const boost::system::error_code& error)
154  {
155  if (error) // e.g., cancelled timer
156  return;
157 
158  m_transport.close();
159  BOOST_THROW_EXCEPTION(Transport::Error(error, "error while connecting to the forwarder"));
160  }
161 
162  void
163  send(BlockSequence&& sequence)
164  {
165  m_transmissionQueue.emplace_back(sequence);
166 
167  if (m_transport.m_isConnected && m_transmissionQueue.size() == 1) {
168  asyncWrite();
169  }
170 
171  // if not connected or there is transmission in progress (m_transmissionQueue.size() > 1),
172  // next write will be scheduled either in connectHandler or in asyncWriteHandler
173  }
174 
175  void
177  {
178  BOOST_ASSERT(!m_transmissionQueue.empty());
179  boost::asio::async_write(m_socket, m_transmissionQueue.front(),
180  bind(&Impl::handleAsyncWrite, this->shared_from_this(), _1, m_transmissionQueue.begin()));
181  }
182 
183  void
184  handleAsyncWrite(const boost::system::error_code& error, TransmissionQueue::iterator queueItem)
185  {
186  if (error) {
187  if (error == boost::system::errc::operation_canceled) {
188  // async receive has been explicitly cancelled (e.g., socket close)
189  return;
190  }
191 
192  m_transport.close();
193  BOOST_THROW_EXCEPTION(Transport::Error(error, "error while sending data to socket"));
194  }
195 
196  if (!m_transport.m_isConnected) {
197  return; // queue has been already cleared
198  }
199 
200  m_transmissionQueue.erase(queueItem);
201 
202  if (!m_transmissionQueue.empty()) {
203  asyncWrite();
204  }
205  }
206 
207  void
209  {
210  m_socket.async_receive(boost::asio::buffer(m_inputBuffer + m_inputBufferSize,
212  bind(&Impl::handleAsyncReceive, this->shared_from_this(), _1, _2));
213  }
214 
215  void
216  handleAsyncReceive(const boost::system::error_code& error, std::size_t nBytesRecvd)
217  {
218  if (error) {
219  if (error == boost::system::errc::operation_canceled) {
220  // async receive has been explicitly cancelled (e.g., socket close)
221  return;
222  }
223 
224  m_transport.close();
225  BOOST_THROW_EXCEPTION(Transport::Error(error, "error while receiving data from socket"));
226  }
227 
228  m_inputBufferSize += nBytesRecvd;
229  // do magic
230 
231  std::size_t offset = 0;
232  bool hasProcessedSome = processAllReceived(m_inputBuffer, offset, m_inputBufferSize);
233  if (!hasProcessedSome && m_inputBufferSize == MAX_NDN_PACKET_SIZE && offset == 0) {
234  m_transport.close();
235  BOOST_THROW_EXCEPTION(Transport::Error(boost::system::error_code(),
236  "input buffer full, but a valid TLV cannot be "
237  "decoded"));
238  }
239 
240  if (offset > 0) {
241  if (offset != m_inputBufferSize) {
243  m_inputBufferSize -= offset;
244  }
245  else {
246  m_inputBufferSize = 0;
247  }
248  }
249 
250  asyncReceive();
251  }
252 
253  bool
254  processAllReceived(uint8_t* buffer, size_t& offset, size_t nBytesAvailable)
255  {
256  while (offset < nBytesAvailable) {
257  bool isOk = false;
258  Block element;
259  std::tie(isOk, element) = Block::fromBuffer(buffer + offset, nBytesAvailable - offset);
260  if (!isOk)
261  return false;
262 
263  m_transport.receive(element);
264  offset += element.size();
265  }
266  return true;
267  }
268 
269 protected:
270  BaseTransport& m_transport;
271 
272  typename Protocol::socket m_socket;
275 
276  TransmissionQueue m_transmissionQueue;
278 
279  boost::asio::deadline_timer m_connectTimer;
280 };
281 
282 } // namespace detail
283 } // namespace ndn
284 
285 #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:193
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)
std::list< BlockSequence > TransmissionQueue
Represents a TLV element of NDN packet format.
Definition: block.hpp:42
Implementation detail of a Boost.Asio-based stream-oriented transport.
void send(BlockSequence &&sequence)
size_t size() const
Get size of encoded wire, including Type-Length-Value.
Definition: block.cpp:298
void handleAsyncReceive(const boost::system::error_code &error, std::size_t nBytesRecvd)
void connectTimeoutHandler(const boost::system::error_code &error)
uint8_t m_inputBuffer[MAX_NDN_PACKET_SIZE]
boost::asio::deadline_timer m_connectTimer
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)
StreamTransportImpl< BaseTransport, Protocol > Impl