Kea 3.3.3
ncr_io.cc
Go to the documentation of this file.
1// Copyright (C) 2013-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>
10#include <dhcp_ddns/ncr_io.h>
12
13#include <boost/algorithm/string/predicate.hpp>
14
15#include <mutex>
16
17namespace isc {
18namespace dhcp_ddns {
19
20using namespace isc::util;
21using namespace std;
22
23NameChangeProtocol stringToNcrProtocol(const std::string& protocol_str) {
24 if (boost::iequals(protocol_str, "UDP")) {
25 return (NCR_UDP);
26 }
27
28 if (boost::iequals(protocol_str, "TCP")) {
29 return (NCR_TCP);
30 }
31
33 "Invalid NameChangeRequest protocol: " << protocol_str);
34}
35
37 switch (protocol) {
38 case NCR_UDP:
39 return ("UDP");
40 case NCR_TCP:
41 return ("TCP");
42 default:
43 break;
44 }
45
46 std::ostringstream stream;
47 stream << "UNKNOWN(" << protocol << ")";
48 return (stream.str());
49}
50
51//************************** NameChangeListener ***************************
52
54 : listening_(false), io_pending_(false), recv_handler_(recv_handler) {
55};
56
57void
59 if (amListening()) {
60 // This amounts to a programmatic error.
61 isc_throw(NcrListenerError, "NameChangeListener is already listening");
62 }
63
64 // Call implementation dependent open.
65 try {
66 open(io_service);
67 } catch (const isc::Exception& ex) {
69 isc_throw(NcrListenerOpenError, "Open failed: " << ex.what());
70 }
71
72 // Set our status to listening.
73 setListening(true);
74
75 // Start the first asynchronous receive.
76 try {
78 } catch (const isc::Exception& ex) {
80 isc_throw(NcrListenerReceiveError, "doReceive failed: " << ex.what());
81 }
82}
83
84void
86 io_pending_ = true;
87 doReceive();
88}
89
90void
92 try {
93 // Call implementation dependent close.
94 close();
95 } catch (const isc::Exception &ex) {
96 // Swallow exceptions. If we have some sort of error we'll log
97 // it but we won't propagate the throw.
99 .arg(ex.what());
100 }
101
102 // Set it false, no matter what. This allows us to at least try to
103 // re-open via startListening().
104 setListening(false);
105}
106
107void
110 // Call the registered application layer handler.
111 // Surround the invocation with a try-catch. The invoked handler is
112 // not supposed to throw, but in the event it does we will at least
113 // report it.
114 try {
115 io_pending_ = false;
116 (*recv_handler_)(result, ncr);
117 } catch (const std::exception& ex) {
119 .arg(ex.what());
120 }
121
122 // Start the next IO layer asynchronous receive.
124}
125
126void
128 // In the event the application handler decided to stop listening
129 // we need to check that first.
130 if (!amListening()) {
131 return;
132 }
133
134 try {
135 receiveNext();
136 } catch (const isc::Exception& isc_ex) {
137 // It is possible though unlikely, for doReceive to fail without
138 // scheduling the read. While, unlikely, it does mean the callback
139 // will not get called with a failure. A throw here would surface
140 // at the IOService::run (or run variant) invocation. So we will
141 // close the window by invoking the application handler with
142 // a failed result, and let the application layer sort it out.
144 .arg(isc_ex.what());
145
146 // Call the registered application layer handler.
147 // Surround the invocation with a try-catch. The invoked handler is
148 // not supposed to throw, but in the event it does we will at least
149 // report it.
151 try {
152 io_pending_ = false;
153 (*recv_handler_)(ERROR, empty);
154 } catch (const std::exception& std_ex) {
157 .arg(std_ex.what());
158 }
159 }
160}
161
162//************************* NameChangeSender ******************************
163
165 size_t send_queue_max)
166 : sending_(false), send_handler_(send_handler),
167 send_queue_max_(send_queue_max), mutex_(new mutex()) {
168
169 // Queue size must be big enough to hold at least 1 entry.
170 setQueueMaxSize(send_queue_max);
171}
172
173void
175 if (amSending()) {
176 // This amounts to a programmatic error.
177 isc_throw(NcrSenderError, "NameChangeSender is already sending");
178 }
179
180 // Call implementation dependent open.
181 try {
182 if (MultiThreadingMgr::instance().getMode()) {
183 lock_guard<mutex> lock(*mutex_);
184 startSendingInternal(io_service);
185 } else {
186 startSendingInternal(io_service);
187 }
188 } catch (const isc::Exception& ex) {
189 stopSending();
190 isc_throw(NcrSenderOpenError, "Open failed: " << ex.what());
191 }
192}
193
194void
195NameChangeSender::startSendingInternal(const isc::asiolink::IOServicePtr& io_service) {
196 // Clear send marker.
197 ncr_to_send_.reset();
198
199 // Remember io service we're given.
200 io_service_ = io_service;
201 open(io_service);
202
203 // Set our status to sending.
204 setSending(true);
205
206 // If there's any queued already.. we'll start sending.
207 sendNext();
208}
209
210void
212 // Set it send indicator to false, no matter what. This allows us to at
213 // least try to re-open via startSending(). Also, setting it false now,
214 // allows us to break sendNext() chain in invokeSendHandler.
215 setSending(false);
216
217 // If there is an outstanding IO to complete, attempt to process it.
218 if (ioReady() && io_service_) {
219 try {
220 runReadyIO();
221 } catch (const std::exception& ex) {
222 // Swallow exceptions. If we have some sort of error we'll log
223 // it but we won't propagate the throw.
225 DHCP_DDNS_NCR_FLUSH_IO_ERROR).arg(ex.what());
226 }
227 }
228
229 try {
230 // Call implementation dependent close.
231 close();
232 } catch (const isc::Exception &ex) {
233 // Swallow exceptions. If we have some sort of error we'll log
234 // it but we won't propagate the throw.
237 }
238
239 if (io_service_) {
240 try {
241 io_service_->stopAndPoll(false);
242 } catch (const std::exception& ex) {
243 // Swallow exceptions. If we have some sort of error we'll log
244 // it but we won't propagate the throw.
246 DHCP_DDNS_NCR_FLUSH_IO_ERROR).arg(ex.what());
247 }
248 }
249
250 io_service_.reset();
251}
252
253void
255 if (!amSending()) {
256 isc_throw(NcrSenderError, "sender is not ready to send");
257 }
258
259 if (!ncr) {
260 isc_throw(NcrSenderError, "request to send is empty");
261 }
262
263 if (MultiThreadingMgr::instance().getMode()) {
264 lock_guard<mutex> lock(*mutex_);
265 sendRequestInternal(ncr);
266 } else {
267 sendRequestInternal(ncr);
268 }
269}
270
271void
272NameChangeSender::sendRequestInternal(NameChangeRequestPtr& ncr) {
273 if (send_queue_.size() >= send_queue_max_) {
275 "send queue has reached maximum capacity: "
276 << send_queue_max_);
277 }
278
279 // Put it on the queue.
280 send_queue_.push_back(ncr);
281
282 // Call sendNext to schedule the next one to go.
283 sendNext();
284}
285
286void
288 if (ncr_to_send_) {
289 // @todo Not sure if there is any risk of getting stuck here but
290 // an interval timer to defend would be good.
291 // In reality, the derivation should ensure they timeout themselves
292 return;
293 }
294
295 // If queue isn't empty, then get one from the front. Note we leave
296 // it on the front of the queue until we successfully send it.
297 if (!send_queue_.empty()) {
298 ncr_to_send_ = send_queue_.front();
299
300 // @todo start defense timer
301 // If a send were to hang and we timed it out, then timeout
302 // handler need to cycle thru open/close ?
303
304 // Call implementation dependent send. If doSend throws before an
305 // asynchronous send is started (for example because the serialized
306 // NCR exceeds the UDP send buffer), clear the in-progress marker and
307 // discard the request so the queue cannot permanently stall.
308 try {
309 doSend(ncr_to_send_);
310 } catch (const std::exception& ex) {
312 .arg(ex.what());
313 send_queue_.pop_front();
314 // Use the internal path: sendNext() may already run under lock.
315 invokeSendHandlerInternal(ERROR);
316 }
317 }
318}
319
320void
322 if (MultiThreadingMgr::instance().getMode()) {
323 lock_guard<mutex> lock(*mutex_);
324 invokeSendHandlerInternal(result);
325 } else {
326 invokeSendHandlerInternal(result);
327 }
328}
329
330void
331NameChangeSender::invokeSendHandlerInternal(const NameChangeSender::Result result) {
332 // @todo reset defense timer
333 if (result == SUCCESS) {
334 // It shipped so pull it off the queue.
335 send_queue_.pop_front();
336 }
337
338 // Invoke the completion handler passing in the result and a pointer
339 // the request involved.
340 // Surround the invocation with a try-catch. The invoked handler is
341 // not supposed to throw, but in the event it does we will at least
342 // report it.
343 try {
344 (*send_handler_)(result, ncr_to_send_);
345 } catch (const std::exception& ex) {
347 .arg(ex.what());
348 }
349
350 // Clear the pending ncr pointer.
351 ncr_to_send_.reset();
352
353 // Set up the next send
354 try {
355 if (amSending()) {
356 sendNext();
357 }
358 } catch (const isc::Exception& isc_ex) {
359 // It is possible though unlikely, for sendNext to fail without
360 // scheduling the send. While, unlikely, it does mean the callback
361 // will not get called with a failure. A throw here would surface
362 // at the IOService::run (or run variant) invocation. So we will
363 // close the window by invoking the application handler with
364 // a failed result, and let the application layer sort it out.
366 .arg(isc_ex.what());
367
368 // Invoke the completion handler passing in failed result.
369 // Surround the invocation with a try-catch. The invoked handler is
370 // not supposed to throw, but in the event it does we will at least
371 // report it.
372 try {
373 (*send_handler_)(ERROR, ncr_to_send_);
374 } catch (const std::exception& std_ex) {
377 }
378 }
379}
380
381void
383 if (MultiThreadingMgr::instance().getMode()) {
384 lock_guard<mutex> lock(*mutex_);
385 skipNextInternal();
386 } else {
387 skipNextInternal();
388 }
389}
390
391void
392NameChangeSender::skipNextInternal() {
393 if (!send_queue_.empty()) {
394 // Discards the request at the front of the queue.
395 send_queue_.pop_front();
396 }
397}
398
399void
401 if (amSending()) {
402 isc_throw(NcrSenderError, "Cannot clear queue while sending");
403 }
404
405 if (MultiThreadingMgr::instance().getMode()) {
406 lock_guard<mutex> lock(*mutex_);
407 send_queue_.clear();
408 } else {
409 send_queue_.clear();
410 }
411}
412
413void
415 if (new_max == 0) {
416 isc_throw(NcrSenderError, "NameChangeSender:"
417 " queue size must be greater than zero");
418 }
419
420 send_queue_max_ = new_max;
421}
422
423size_t
425 if (MultiThreadingMgr::instance().getMode()) {
426 lock_guard<mutex> lock(*mutex_);
427 return (getQueueSizeInternal());
428 } else {
429 return (getQueueSizeInternal());
430 }
431}
432
433size_t
434NameChangeSender::getQueueSizeInternal() const {
435 return (send_queue_.size());
436}
437
439NameChangeSender::peekAt(const size_t index) const {
440 if (MultiThreadingMgr::instance().getMode()) {
441 lock_guard<mutex> lock(*mutex_);
442 return (peekAtInternal(index));
443 } else {
444 return (peekAtInternal(index));
445 }
446}
447
449NameChangeSender::peekAtInternal(const size_t index) const {
450 auto size = getQueueSizeInternal();
451 if (index >= size) {
453 "NameChangeSender::peekAt peek beyond end of queue attempted"
454 << " index: " << index << " queue size: " << size);
455 }
456
457 return (send_queue_.at(index));
458}
459
460bool
462 if (MultiThreadingMgr::instance().getMode()) {
463 lock_guard<mutex> lock(*mutex_);
464 return ((ncr_to_send_) ? true : false);
465 } else {
466 return ((ncr_to_send_) ? true : false);
467 }
468}
469
470void
472 if (source_sender.amSending()) {
473 isc_throw(NcrSenderError, "Cannot assume queue:"
474 " source sender is actively sending");
475 }
476
477 if (amSending()) {
478 isc_throw(NcrSenderError, "Cannot assume queue:"
479 " target sender is actively sending");
480 }
481
482 if (getQueueMaxSize() < source_sender.getQueueSize()) {
483 isc_throw(NcrSenderError, "Cannot assume queue:"
484 " source queue count exceeds target queue max");
485 }
486
487 if (MultiThreadingMgr::instance().getMode()) {
488 lock_guard<mutex> lock(*mutex_);
489 assumeQueueInternal(source_sender);
490 } else {
491 assumeQueueInternal(source_sender);
492 }
493}
494
495void
496NameChangeSender::assumeQueueInternal(NameChangeSender& source_sender) {
497 if (!send_queue_.empty()) {
498 isc_throw(NcrSenderError, "Cannot assume queue:"
499 " target queue is not empty");
500 }
501
502 send_queue_.swap(source_sender.getSendQueue());
503}
504
505int
507 isc_throw(NotImplemented, "NameChangeSender::getSelectFd is not supported");
508}
509
510void
512 if (!io_service_) {
513 isc_throw(NcrSenderError, "NameChangeSender::runReadyIO"
514 " sender io service is null");
515 }
516
517 // We shouldn't be here if IO isn't ready to execute.
518 // By running poll we're guaranteed not to hang.
519 io_service_->pollOne();
520}
521
522} // namespace dhcp_ddns
523} // namespace isc
A generic exception that is thrown if a parameter given to a method is considered invalid in that con...
This is a base class for exceptions thrown from the DNS library module.
virtual const char * what() const
Returns a C-style character string of the cause of the exception.
A generic exception that is thrown when a function is not implemented.
virtual void open(const isc::asiolink::IOServicePtr &io_service)=0
Abstract method which opens the IO source for reception.
boost::shared_ptr< RequestReceiveHandler > RequestReceiveHandlerPtr
Defines a smart pointer to an instance of a request receive handler.
Definition ncr_io.h:207
void stopListening()
Closes the IO source and stops listen logic.
Definition ncr_io.cc:91
virtual void close()=0
Abstract method which closes the IO source.
NameChangeListener(RequestReceiveHandlerPtr recv_handler)
Constructor.
Definition ncr_io.cc:53
Result
Defines the outcome of an asynchronous NCR receive.
Definition ncr_io.h:172
virtual void doReceive()=0
Initiates an IO layer asynchronous read.
void invokeRecvHandler(const Result result, NameChangeRequestPtr &ncr)
Calls the NCR receive handler registered with the listener.
Definition ncr_io.cc:108
bool amListening() const
Returns true if the listener is listening, false otherwise.
Definition ncr_io.h:327
void receiveNext()
Initiates an asynchronous receive.
Definition ncr_io.cc:85
void startListening(const isc::asiolink::IOServicePtr &io_service)
Prepares the IO for reception and initiates the first receive.
Definition ncr_io.cc:58
void scheduleNextReceive()
Schedules the next asynchronous receive if still listening.
Definition ncr_io.cc:127
Abstract interface for sending NameChangeRequests.
Definition ncr_io.h:479
asiolink::IOServicePtr io_service_
Pointer to the IOService currently being used by the sender.
Definition ncr_io.h:860
void stopSending()
Closes the IO sink and stops send logic.
Definition ncr_io.cc:211
virtual int getSelectFd()=0
Returns a file descriptor suitable for use with select.
Definition ncr_io.cc:506
void startSending(const isc::asiolink::IOServicePtr &io_service)
Prepares the IO for transmission.
Definition ncr_io.cc:174
NameChangeSender(RequestSendHandlerPtr send_handler, size_t send_queue_max=MAX_QUEUE_DEFAULT)
Constructor.
Definition ncr_io.cc:164
void assumeQueue(NameChangeSender &source_sender)
Move all queued requests from a given sender into the send queue.
Definition ncr_io.cc:471
Result
Defines the outcome of an asynchronous NCR send.
Definition ncr_io.h:489
size_t getQueueMaxSize() const
Returns the maximum number of entries allowed in the send queue.
Definition ncr_io.h:780
size_t getQueueSize() const
Returns the number of entries currently in the send queue.
Definition ncr_io.cc:424
const NameChangeRequestPtr & peekAt(const size_t index) const
Returns the entry at a given position in the queue.
Definition ncr_io.cc:439
virtual bool ioReady()=0
Returns whether or not the sender has IO ready to process.
void skipNext()
Removes the request at the front of the send queue.
Definition ncr_io.cc:382
boost::shared_ptr< RequestSendHandler > RequestSendHandlerPtr
Defines a smart pointer to an instance of a request send handler.
Definition ncr_io.h:539
void clearSendQueue()
Flushes all entries in the send queue.
Definition ncr_io.cc:400
bool amSending() const
Returns true if the sender is in send mode, false otherwise.
Definition ncr_io.h:765
virtual void doSend(NameChangeRequestPtr &ncr)=0
Initiates an IO layer asynchronous send.
void setQueueMaxSize(const size_t new_max)
Sets the maximum queue size to the given value.
Definition ncr_io.cc:414
void invokeSendHandler(const NameChangeSender::Result result)
Calls the NCR send completion handler registered with the sender.
Definition ncr_io.cc:321
virtual void open(const isc::asiolink::IOServicePtr &io_service)=0
Abstract method which opens the IO sink for transmission.
virtual void close()=0
Abstract method which closes the IO sink.
void sendRequest(NameChangeRequestPtr &ncr)
Queues the given request to be sent.
Definition ncr_io.cc:254
virtual void runReadyIO()
Processes sender IO events.
Definition ncr_io.cc:511
SendQueue & getSendQueue()
Returns a reference to the send queue.
Definition ncr_io.h:838
bool isSendInProgress() const
Returns true when a send is in progress.
Definition ncr_io.cc:461
void sendNext()
Dequeues and sends the next request on the send queue in a thread safe context.
Definition ncr_io.cc:287
Exception thrown if an NcrListenerError encounters a general error.
Definition ncr_io.h:95
Exception thrown if an error occurs during IO source open.
Definition ncr_io.h:102
Exception thrown if an error occurs initiating an IO receive.
Definition ncr_io.h:109
Thrown when a NameChangeSender encounters an error.
Definition ncr_io.h:371
Exception thrown if an error occurs during IO source open.
Definition ncr_io.h:378
Exception thrown if an error occurs initiating an IO send.
Definition ncr_io.h:385
static MultiThreadingMgr & instance()
Returns a single instance of Multi Threading Manager.
#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
isc::log::Logger dhcp_ddns_logger("libdhcp-ddns")
Defines the logger used within lib dhcp_ddns.
NameChangeProtocol stringToNcrProtocol(const std::string &protocol_str)
Function which converts text labels to NameChangeProtocol enums.
Definition ncr_io.cc:23
const isc::log::MessageID DHCP_DDNS_NCR_SEND_CLOSE_ERROR
const isc::log::MessageID DHCP_DDNS_NCR_SEND_NEXT_ERROR
NameChangeProtocol
Defines the list of socket protocols supported.
Definition ncr_io.h:70
const isc::log::MessageID DHCP_DDNS_NCR_FLUSH_IO_ERROR
std::string ncrProtocolToString(NameChangeProtocol protocol)
Function which converts NameChangeProtocol enums to text labels.
Definition ncr_io.cc:36
const isc::log::MessageID DHCP_DDNS_UNCAUGHT_NCR_SEND_HANDLER_ERROR
const isc::log::MessageID DHCP_DDNS_UNCAUGHT_NCR_RECV_HANDLER_ERROR
boost::shared_ptr< NameChangeRequest > NameChangeRequestPtr
Defines a pointer to a NameChangeRequest.
Definition ncr_msg.h:245
const isc::log::MessageID DHCP_DDNS_NCR_LISTEN_CLOSE_ERROR
const isc::log::MessageID DHCP_DDNS_NCR_RECV_NEXT_ERROR
Defines the logger used by the top-level component of kea-lfc.
This file defines abstract classes for exchanging NameChangeRequests.