Kea 3.3.3
tcp_connection.cc
Go to the documentation of this file.
1// Copyright (C) 2022-2026 Internet Systems Consortium, Inc. ("ISC")
2//
3// This Source Code Form is subject to the terms of the Mozilla Public
4// License, v. 2.0. If a copy of the MPL was not distributed with this
5// file, You can obtain one at http://mozilla.org/MPL/2.0/.
6
7#include <config.h>
8
10#include <tcp/tcp_connection.h>
12#include <tcp/tcp_log.h>
13#include <tcp/tcp_messages.h>
14#include <boost/make_shared.hpp>
15
16#include <iomanip>
17#include <sstream>
18#include <functional>
19
20using namespace isc::asiolink;
21namespace ph = std::placeholders;
22
23namespace isc {
24namespace tcp {
25
26void
27TcpConnection::
28SocketCallback::operator()(boost::system::error_code ec, size_t length) {
29 if (ec.value() == boost::asio::error::operation_aborted) {
30 return;
31 }
32 callback_(ec, length);
33}
34
36 const TcpConnectionAcceptorPtr& acceptor,
37 const TlsContextPtr& tls_context,
38 TcpConnectionPool& connection_pool,
39 const TcpConnectionAcceptorCallback& acceptor_callback,
40 const TcpConnectionFilterCallback& connection_filter,
41 const long idle_timeout,
42 const size_t read_max /* = 32768 */)
43 : io_service_(io_service),
44 tls_context_(tls_context),
45 idle_timeout_(idle_timeout),
46 idle_timer_(io_service),
49 acceptor_(acceptor),
50 connection_pool_(connection_pool),
51 acceptor_callback_(acceptor_callback),
52 connection_filter_(connection_filter),
53 read_max_(read_max),
54 input_buf_(read_max) {
55 if (!tls_context) {
57 } else {
59 tls_context));
60 }
61}
62
66
67void
68TcpConnection::shutdownCallback(const boost::system::error_code&) {
69 tls_socket_->close();
70}
71
72void
74 idle_timer_.cancel();
75 if (tcp_socket_) {
76 tcp_socket_->close();
77 return;
78 }
79
80 if (tls_socket_) {
81 // Create instance of the callback to close the socket.
82 SocketCallback cb(std::bind(&TcpConnection::shutdownCallback,
83 shared_from_this(),
84 ph::_1)); // error_code
85 tls_socket_->shutdown(cb);
86 return;
87 }
88
89 // Not reachable?
90 isc_throw(Unexpected, "internal error: unable to shutdown the socket");
91}
92
93void
95 idle_timer_.cancel();
96 if (tcp_socket_) {
97 tcp_socket_->close();
98 return;
99 }
100
101 if (tls_socket_) {
102 tls_socket_->close();
103 return;
104 }
105
106 // Not reachable?
107 isc_throw(Unexpected, "internal error: unable to close the socket");
108}
109
110void
121
122void
133
134void
136 // Create instance of the callback. It is safe to pass the local instance
137 // of the callback, because the underlying boost functions make copies
138 // as needed.
140 shared_from_this(),
141 ph::_1);
142 try {
143 TlsConnectionAcceptorPtr tls_acceptor =
144 boost::dynamic_pointer_cast<TlsConnectionAcceptor>(acceptor_);
145 if (!tls_acceptor) {
146 if (!tcp_socket_) {
147 isc_throw(Unexpected, "internal error: TCP socket is null");
148 }
149 acceptor_->asyncAccept(*tcp_socket_, cb);
150 } else {
151 if (!tls_socket_) {
152 isc_throw(Unexpected, "internal error: TLS socket is null");
153 }
154 tls_acceptor->asyncAccept(*tls_socket_, cb);
155 }
156 } catch (const std::exception& ex) {
157 isc_throw(TcpConnectionError, "unable to start accepting TCP "
158 "connections: " << ex.what());
159 }
160}
161
162void
164 // Skip the handshake if the socket is not a TLS one.
165 if (!tls_socket_) {
166 doRead();
167 return;
168 }
169
171
172 // Create instance of the callback. It is safe to pass the local instance
173 // of the callback, because the underlying boost functions make copies
174 // as needed.
175 SocketCallback cb(std::bind(&TcpConnection::handshakeCallback,
176 shared_from_this(),
177 ph::_1)); // error
178 try {
179 tls_socket_->handshake(cb);
180
181 } catch (const std::exception& ex) {
182 isc_throw(TcpConnectionError, "unable to perform TLS handshake: "
183 << ex.what());
184 }
185}
186
187void
189 try {
190 TCPEndpoint endpoint;
191
193
194 // Request hasn't been created if we are starting to read the
195 // new request.
196 if (!request) {
197 request = createRequest();
198 }
199
200 // Create instance of the callback. It is safe to pass the local instance
201 // of the callback, because the underlying std functions make copies
202 // as needed.
203 SocketCallback cb(std::bind(&TcpConnection::socketReadCallback,
204 shared_from_this(),
205 request,
206 ph::_1, // error
207 ph::_2)); // bytes_transferred
208 if (tcp_socket_) {
209 tcp_socket_->asyncReceive(static_cast<void*>(getInputBufData()),
210 getInputBufSize(), 0, &endpoint, cb);
211 return;
212 }
213
214 if (tls_socket_) {
215 tls_socket_->asyncReceive(static_cast<void*>(getInputBufData()),
216 getInputBufSize(), 0, &endpoint, cb);
217 return;
218 }
219 } catch (...) {
221 }
222}
223
224void
226 try {
227 if (response->wireDataAvail()) {
228 // Create instance of the callback. It is safe to pass the
229 // local instance of the callback, because the underlying
230 // std functions make copies as needed.
231 SocketCallback cb(std::bind(&TcpConnection::socketWriteCallback,
232 shared_from_this(),
233 response,
234 ph::_1, // error
235 ph::_2)); // bytes_transferred
236 if (tcp_socket_) {
240 tcp_socket_->asyncSend(response->getWireData(),
241 response->getWireDataSize(),
242 cb);
243 return;
244 }
245 if (tls_socket_) {
249 tls_socket_->asyncSend(response->getWireData(),
250 response->getWireDataSize(),
251 cb);
252 return;
253 }
254 } else {
255 // The connection remains open and we are done sending the response.
256 // If the response sent handler returns true then we should start the
257 // idle timer.
258 if (responseSent(response)) {
260 }
261 }
262 } catch (...) {
263 // The connection is dead and there can't be a pending write as
264 // they are in sequence.
266 }
267}
268
269void
273
274
275void
276TcpConnection::acceptorCallback(const boost::system::error_code& ec) {
277 if (!acceptor_->isOpen()) {
278 return;
279 }
280
281 if (ec) {
283 }
284
285 // Stage a new connection to listen for next client.
287
288 if (!ec) {
289 try {
290 if (tcp_socket_ && tcp_socket_->getASIOSocket().is_open()) {
292 tcp_socket_->getASIOSocket().remote_endpoint();
293 } else if (tls_socket_ && tls_socket_->getASIOSocket().is_open()) {
295 tls_socket_->getASIOSocket().remote_endpoint();
296 }
297 } catch (...) {
298 // Let's it to fail later.
299 }
300
301 // In theory, we should not get here with an unopened socket
302 // but just in case, we'll check for NO_ENDPOINT.
303 if ((remote_endpoint_ == NO_ENDPOINT()) ||
310 return;
311 }
312
313 if (!tls_context_) {
317 .arg(static_cast<unsigned>(idle_timeout_/1000));
318 } else {
322 .arg(static_cast<unsigned>(idle_timeout_/1000));
323 }
324
325 doHandshake();
326 }
327}
328
329void
330TcpConnection::handshakeCallback(const boost::system::error_code& ec) {
331 if (ec) {
334 .arg(ec.message());
336 } else {
340 .arg(static_cast<unsigned>(idle_timeout_/1000));
341 doRead();
342 }
343}
344
345void
347 boost::system::error_code ec, size_t length) {
348 if (ec) {
349 // IO service has been stopped and the connection is probably
350 // going to be shutting down.
351 if (ec.value() == boost::asio::error::operation_aborted) {
352 return;
353
354 // EWOULDBLOCK and EAGAIN are special cases. Everything else is
355 // treated as fatal error.
356 } else if ((ec.value() != boost::asio::error::try_again) &&
357 (ec.value() != boost::asio::error::would_block)) {
359 return;
360
361 // We got EWOULDBLOCK or EAGAIN which indicate that we may be able to
362 // read something from the socket on the next attempt. Just make sure
363 // we don't try to read anything now in case there is any garbage
364 // passed in length.
365 } else {
366 length = 0;
367 }
368 }
369
370 // Data received, Restart the request timer.
372
373 TcpRequestPtr next_request = request;
374 if (length) {
377 .arg(length)
379 WireData input_data(input_buf_.begin(), input_buf_.begin() + length);
380 next_request = postData(request, input_data);
381 }
382
383 // Start next read.
384 doRead(next_request);
385}
386
389 size_t bytes_left = 0;
390 size_t length = input_data.size();
391 if (length) {
392 // Add data to the current request.
393 size_t bytes_used = request->postBuffer(static_cast<void*>(input_data.data()), length);
394 // Remove only the bytes consumed; leftover data may belong to the
395 // next pipelined message in the same read.
396 bytes_left = length - bytes_used;
397 input_data.erase(input_data.begin(), input_data.begin() + bytes_used);
398 }
399
400 if (request->needData()) {
401 // Current request is incomplete and we're out of data
402 // return the incomplete request and we'll read again.
403 return (request);
404 }
405
406 try {
410
411 // Request complete, stop the timer.
412 idle_timer_.cancel();
413
414 // Process the completed request.
415 requestReceived(request);
416 } catch (const std::exception& ex) {
419 .arg(ex.what());
420 }
421
422 // Create a new, empty request.
423 request = createRequest();
424 if (bytes_left) {
425 // The input buffer spanned messages. Recurse to post the remainder to the
426 // new request.
427 request = postData(request, input_data);
428 }
429
430 return (request);
431}
432
433void
435 boost::system::error_code ec, size_t length) {
436 if (ec) {
437 // IO service has been stopped and the connection is probably
438 // going to be shutting down.
439 if (ec.value() == boost::asio::error::operation_aborted) {
440 return;
441
442 // EWOULDBLOCK and EAGAIN are special cases. Everything else is
443 // treated as fatal error.
444 } else if ((ec.value() != boost::asio::error::try_again) &&
445 (ec.value() != boost::asio::error::would_block)) {
446 // The connection is dead and there can't be a pending write as
447 // they are in sequence.
449 return;
450
451 // We got EWOULDBLOCK or EAGAIN which indicate that we may be able to
452 // write something to the socket on the next attempt. Just make sure
453 // we don't consume wire data now in case there is any garbage
454 // passed in length.
455 } else {
456 length = 0;
457 }
458 }
459
461 .arg(length)
463
464 // Eat the 'length' number of bytes from the output buffer and only
465 // leave the part of the response that hasn't been sent.
466 response->consumeWireData(length);
467
468 // Schedule the write of the unsent data.
469 doWrite(response);
470}
471
472void
477
478void
483 // In theory we should shutdown first and stop/close after but
484 // it is better to put the connection management responsibility
485 // on the client... so simply drop idle connections.
487}
488
489std::string
491 if (remote_endpoint_ != NO_ENDPOINT()) {
492 return (remote_endpoint_.address().to_string());
493 }
494
495 return ("(unknown address)");
496}
497
498void
499TcpConnection::setReadMax(const size_t read_max) {
500 if (!read_max) {
501 isc_throw(BadValue, "TcpConnection read_max must be > 0");
502 }
503
504 read_max_ = read_max;
505 input_buf_.resize(read_max);
506}
507
508} // end of namespace isc::tcp
509} // end of namespace isc
A generic exception that is thrown if a parameter given to a method is considered invalid in that con...
A generic exception that is thrown when an unexpected error condition occurs.
Generic error reported within TcpConnection class.
Pool of active TCP connections.
static std::atomic< uint64_t > rejected_counter_
Class/static rejected (by the accept filter) connection counter.
void doWrite(TcpResponsePtr response)
Starts asynchronous write to the socket.
virtual TcpRequestPtr createRequest()=0
Creates a new, empty request.
void shutdownCallback(const boost::system::error_code &ec)
Callback invoked when TLS shutdown is performed.
unsigned char * getInputBufData()
Returns pointer to the first byte of the input buffer.
void acceptorCallback(const boost::system::error_code &ec)
Local callback invoked when new connection is accepted.
void asyncAccept()
Asynchronously accepts new connection.
size_t getInputBufSize() const
Returns input buffer size.
virtual void shutdown()
Shutdown the socket.
boost::asio::ip::tcp::endpoint remote_endpoint_
Remote endpoint.
virtual void requestReceived(TcpRequestPtr request)=0
Processes a request once it has been completely received.
virtual void shutdownConnection()
Shuts down current connection.
virtual void stopThisConnection()
Stops current connection.
void setReadMax(const size_t read_max)
Sets the maximum number of bytes read during single socket read.
TcpConnectionAcceptorPtr acceptor_
Pointer to the TCP acceptor used to accept new connections.
virtual void close()
Closes the socket.
asiolink::TlsContextPtr tls_context_
TLS context.
WireData input_buf_
Buffer for a single socket read.
void socketReadCallback(TcpRequestPtr request, boost::system::error_code ec, size_t length)
Callback invoked when new data is received over the socket.
void setupIdleTimer()
Reset timer for detecting idle timeout in connections.
virtual void socketWriteCallback(TcpResponsePtr request, boost::system::error_code ec, size_t length)
Callback invoked when data is sent over the socket.
void asyncSendResponse(TcpResponsePtr response)
Sends TCP response asynchronously.
asiolink::IOServicePtr io_service_
The IO service used to handle events.
std::string getRemoteEndpointAddressAsText() const
returns remote address in textual form
TcpConnectionFilterCallback connection_filter_
External callback for filtering connections by IP address.
std::unique_ptr< asiolink::TCPSocket< SocketCallback > > tcp_socket_
TCP socket used by this connection.
void doRead(TcpRequestPtr request=TcpRequestPtr())
Starts asynchronous read from the socket.
static const boost::asio::ip::tcp::endpoint & NO_ENDPOINT()
Returns an empty end point.
virtual ~TcpConnection()
Destructor.
asiolink::IntervalTimer idle_timer_
Timer used to detect idle Timeout.
virtual bool responseSent(TcpResponsePtr response)=0
Determines behavior after a response has been sent.
TcpConnectionAcceptorCallback acceptor_callback_
External TCP acceptor callback.
size_t read_max_
Maximum bytes to read in a single socket read.
TcpConnectionPool & connection_pool_
Connection pool holding this connection.
void doHandshake()
Asynchronously performs TLS handshake.
TcpRequestPtr postData(TcpRequestPtr request, WireData &input_data)
Appends newly received raw data to the given request.
std::unique_ptr< asiolink::TLSSocket< SocketCallback > > tls_socket_
TLS socket used by this connection.
void handshakeCallback(const boost::system::error_code &ec)
Local callback invoked when TLS handshake is performed.
void idleTimeoutCallback()
Callback invoked when the client has been idle.
TcpConnection(const asiolink::IOServicePtr &io_service, const TcpConnectionAcceptorPtr &acceptor, const asiolink::TlsContextPtr &tls_context, TcpConnectionPool &connection_pool, const TcpConnectionAcceptorCallback &acceptor_callback, const TcpConnectionFilterCallback &connection_filter, const long idle_timeout, const size_t read_max=32768)
Constructor.
long idle_timeout_
Timeout after which the a TCP connection is shut down by the server.
#define isc_throw(type, stream)
A shortcut macro to insert known values into exception arguments.
#define LOG_ERROR(LOGGER, MESSAGE)
Macro to conveniently test error output and log it.
Definition macros.h:32
#define LOG_INFO(LOGGER, MESSAGE)
Macro to conveniently test info output and log it.
Definition macros.h:20
#define LOG_DEBUG(LOGGER, LEVEL, MESSAGE)
Macro to conveniently test debug output and log it.
Definition macros.h:14
const int DBGLVL_TRACE_BASIC
Trace basic operations.
const int DBGLVL_TRACE_DETAIL_DATA
Trace data associated with detailed operations.
const int DBGLVL_TRACE_DETAIL
Trace detailed operations.
std::function< bool(const boost::asio::ip::tcp::endpoint &)> TcpConnectionFilterCallback
Type of the callback for filtering new connections by ip address.
boost::shared_ptr< TlsConnectionAcceptor > TlsConnectionAcceptorPtr
Type of shared pointer to TLS acceptors.
const isc::log::MessageID TCP_CONNECTION_STOP_FAILED
const isc::log::MessageID TCP_IDLE_CONNECTION_TIMEOUT_OCCURRED
boost::shared_ptr< TcpConnectionAcceptor > TcpConnectionAcceptorPtr
Type of shared pointer to TCP acceptors.
const isc::log::MessageID TLS_REQUEST_RECEIVE_START
const isc::log::MessageID TLS_CONNECTION_HANDSHAKE_FAILED
const isc::log::MessageID TCP_CONNECTION_STOP
const isc::log::MessageID TCP_DATA_SENT
boost::shared_ptr< TcpRequest > TcpRequestPtr
Defines a smart pointer to a TcpRequest.
const isc::log::MessageID TCP_CONNECTION_REJECTED_BY_FILTER
const isc::log::MessageID TLS_CONNECTION_HANDSHAKE_START
boost::shared_ptr< TcpResponse > TcpResponsePtr
const isc::log::MessageID TCP_SERVER_CLIENT_REQUEST_RECEIVED
const isc::log::MessageID TLS_SERVER_RESPONSE_SEND
const isc::log::MessageID TCP_DATA_RECEIVED
const isc::log::MessageID TCP_REQUEST_RECEIVE_START
std::vector< uint8_t > WireData
Defines a data structure for storing raw bytes of data on the wire.
Definition wire_data.h:17
isc::log::Logger tcp_logger("tcp")
Defines the logger used within libkea-tcp library.
Definition tcp_log.h:18
const isc::log::MessageID TCP_SERVER_RESPONSE_SEND
const isc::log::MessageID TCP_REQUEST_RECEIVED_FAILED
const isc::log::MessageID TCP_CONNECTION_SHUTDOWN
const isc::log::MessageID TCP_CONNECTION_SHUTDOWN_FAILED
std::function< void(const boost::system::error_code &)> TcpConnectionAcceptorCallback
Type of the callback for the TCP acceptor used in this library.
Defines the logger used by the top-level component of kea-lfc.