/src/kea/src/lib/dhcp_ddns/ncr_io.cc
Line | Count | Source |
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> |
8 | | #include <asiolink/asio_wrapper.h> |
9 | | #include <dhcp_ddns/dhcp_ddns_log.h> |
10 | | #include <dhcp_ddns/ncr_io.h> |
11 | | #include <util/multi_threading_mgr.h> |
12 | | |
13 | | #include <boost/algorithm/string/predicate.hpp> |
14 | | |
15 | | #include <mutex> |
16 | | |
17 | | namespace isc { |
18 | | namespace dhcp_ddns { |
19 | | |
20 | | using namespace isc::util; |
21 | | using namespace std; |
22 | | |
23 | 159k | NameChangeProtocol stringToNcrProtocol(const std::string& protocol_str) { |
24 | 159k | if (boost::iequals(protocol_str, "UDP")) { |
25 | 159k | return (NCR_UDP); |
26 | 159k | } |
27 | | |
28 | 11 | if (boost::iequals(protocol_str, "TCP")) { |
29 | 1 | return (NCR_TCP); |
30 | 1 | } |
31 | | |
32 | 10 | isc_throw(BadValue, |
33 | 10 | "Invalid NameChangeRequest protocol: " << protocol_str); |
34 | 10 | } |
35 | | |
36 | 23.9k | std::string ncrProtocolToString(NameChangeProtocol protocol) { |
37 | 23.9k | switch (protocol) { |
38 | 23.9k | case NCR_UDP: |
39 | 23.9k | return ("UDP"); |
40 | 1 | case NCR_TCP: |
41 | 1 | return ("TCP"); |
42 | 0 | default: |
43 | 0 | break; |
44 | 23.9k | } |
45 | | |
46 | 0 | std::ostringstream stream; |
47 | 0 | stream << "UNKNOWN(" << protocol << ")"; |
48 | 0 | return (stream.str()); |
49 | 23.9k | } |
50 | | |
51 | | //************************** NameChangeListener *************************** |
52 | | |
53 | | NameChangeListener::NameChangeListener(RequestReceiveHandlerPtr recv_handler) |
54 | 0 | : listening_(false), io_pending_(false), recv_handler_(recv_handler) { |
55 | 0 | }; |
56 | | |
57 | | void |
58 | 0 | NameChangeListener::startListening(const isc::asiolink::IOServicePtr& io_service) { |
59 | 0 | if (amListening()) { |
60 | | // This amounts to a programmatic error. |
61 | 0 | isc_throw(NcrListenerError, "NameChangeListener is already listening"); |
62 | 0 | } |
63 | | |
64 | | // Call implementation dependent open. |
65 | 0 | try { |
66 | 0 | open(io_service); |
67 | 0 | } catch (const isc::Exception& ex) { |
68 | 0 | stopListening(); |
69 | 0 | isc_throw(NcrListenerOpenError, "Open failed: " << ex.what()); |
70 | 0 | } |
71 | | |
72 | | // Set our status to listening. |
73 | 0 | setListening(true); |
74 | | |
75 | | // Start the first asynchronous receive. |
76 | 0 | try { |
77 | 0 | receiveNext(); |
78 | 0 | } catch (const isc::Exception& ex) { |
79 | 0 | stopListening(); |
80 | 0 | isc_throw(NcrListenerReceiveError, "doReceive failed: " << ex.what()); |
81 | 0 | } |
82 | 0 | } |
83 | | |
84 | | void |
85 | 0 | NameChangeListener::receiveNext() { |
86 | 0 | io_pending_ = true; |
87 | 0 | doReceive(); |
88 | 0 | } |
89 | | |
90 | | void |
91 | 0 | NameChangeListener::stopListening() { |
92 | 0 | try { |
93 | | // Call implementation dependent close. |
94 | 0 | close(); |
95 | 0 | } 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. |
98 | 0 | LOG_ERROR(dhcp_ddns_logger, DHCP_DDNS_NCR_LISTEN_CLOSE_ERROR) |
99 | 0 | .arg(ex.what()); |
100 | 0 | } |
101 | | |
102 | | // Set it false, no matter what. This allows us to at least try to |
103 | | // re-open via startListening(). |
104 | 0 | setListening(false); |
105 | 0 | } |
106 | | |
107 | | void |
108 | | NameChangeListener::invokeRecvHandler(const Result result, |
109 | 0 | NameChangeRequestPtr& ncr) { |
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 | 0 | try { |
115 | 0 | io_pending_ = false; |
116 | 0 | (*recv_handler_)(result, ncr); |
117 | 0 | } catch (const std::exception& ex) { |
118 | 0 | LOG_ERROR(dhcp_ddns_logger, DHCP_DDNS_UNCAUGHT_NCR_RECV_HANDLER_ERROR) |
119 | 0 | .arg(ex.what()); |
120 | 0 | } |
121 | | |
122 | | // Start the next IO layer asynchronous receive. |
123 | 0 | scheduleNextReceive(); |
124 | 0 | } |
125 | | |
126 | | void |
127 | 0 | NameChangeListener::scheduleNextReceive() { |
128 | | // In the event the application handler decided to stop listening |
129 | | // we need to check that first. |
130 | 0 | if (!amListening()) { |
131 | 0 | return; |
132 | 0 | } |
133 | | |
134 | 0 | try { |
135 | 0 | receiveNext(); |
136 | 0 | } 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. |
143 | 0 | LOG_ERROR(dhcp_ddns_logger, DHCP_DDNS_NCR_RECV_NEXT_ERROR) |
144 | 0 | .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. |
150 | 0 | NameChangeRequestPtr empty; |
151 | 0 | try { |
152 | 0 | io_pending_ = false; |
153 | 0 | (*recv_handler_)(ERROR, empty); |
154 | 0 | } catch (const std::exception& std_ex) { |
155 | 0 | LOG_ERROR(dhcp_ddns_logger, |
156 | 0 | DHCP_DDNS_UNCAUGHT_NCR_RECV_HANDLER_ERROR) |
157 | 0 | .arg(std_ex.what()); |
158 | 0 | } |
159 | 0 | } |
160 | 0 | } |
161 | | |
162 | | //************************* NameChangeSender ****************************** |
163 | | |
164 | | NameChangeSender::NameChangeSender(RequestSendHandlerPtr send_handler, |
165 | | size_t send_queue_max) |
166 | 0 | : sending_(false), send_handler_(send_handler), |
167 | 0 | send_queue_max_(send_queue_max), mutex_(new mutex()) { |
168 | | |
169 | | // Queue size must be big enough to hold at least 1 entry. |
170 | 0 | setQueueMaxSize(send_queue_max); |
171 | 0 | } |
172 | | |
173 | | void |
174 | 0 | NameChangeSender::startSending(const isc::asiolink::IOServicePtr& io_service) { |
175 | 0 | if (amSending()) { |
176 | | // This amounts to a programmatic error. |
177 | 0 | isc_throw(NcrSenderError, "NameChangeSender is already sending"); |
178 | 0 | } |
179 | | |
180 | | // Call implementation dependent open. |
181 | 0 | try { |
182 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
183 | 0 | lock_guard<mutex> lock(*mutex_); |
184 | 0 | startSendingInternal(io_service); |
185 | 0 | } else { |
186 | 0 | startSendingInternal(io_service); |
187 | 0 | } |
188 | 0 | } catch (const isc::Exception& ex) { |
189 | 0 | stopSending(); |
190 | 0 | isc_throw(NcrSenderOpenError, "Open failed: " << ex.what()); |
191 | 0 | } |
192 | 0 | } |
193 | | |
194 | | void |
195 | 0 | NameChangeSender::startSendingInternal(const isc::asiolink::IOServicePtr& io_service) { |
196 | | // Clear send marker. |
197 | 0 | ncr_to_send_.reset(); |
198 | | |
199 | | // Remember io service we're given. |
200 | 0 | io_service_ = io_service; |
201 | 0 | open(io_service); |
202 | | |
203 | | // Set our status to sending. |
204 | 0 | setSending(true); |
205 | | |
206 | | // If there's any queued already.. we'll start sending. |
207 | 0 | sendNext(); |
208 | 0 | } |
209 | | |
210 | | void |
211 | 0 | NameChangeSender::stopSending() { |
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 | 0 | setSending(false); |
216 | | |
217 | | // If there is an outstanding IO to complete, attempt to process it. |
218 | 0 | if (ioReady() && io_service_) { |
219 | 0 | try { |
220 | 0 | runReadyIO(); |
221 | 0 | } 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. |
224 | 0 | LOG_ERROR(dhcp_ddns_logger, |
225 | 0 | DHCP_DDNS_NCR_FLUSH_IO_ERROR).arg(ex.what()); |
226 | 0 | } |
227 | 0 | } |
228 | |
|
229 | 0 | try { |
230 | | // Call implementation dependent close. |
231 | 0 | close(); |
232 | 0 | } 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. |
235 | 0 | LOG_ERROR(dhcp_ddns_logger, |
236 | 0 | DHCP_DDNS_NCR_SEND_CLOSE_ERROR).arg(ex.what()); |
237 | 0 | } |
238 | |
|
239 | 0 | if (io_service_) { |
240 | 0 | try { |
241 | 0 | io_service_->stopAndPoll(false); |
242 | 0 | } 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. |
245 | 0 | LOG_ERROR(dhcp_ddns_logger, |
246 | 0 | DHCP_DDNS_NCR_FLUSH_IO_ERROR).arg(ex.what()); |
247 | 0 | } |
248 | 0 | } |
249 | |
|
250 | 0 | io_service_.reset(); |
251 | 0 | } |
252 | | |
253 | | void |
254 | 0 | NameChangeSender::sendRequest(NameChangeRequestPtr& ncr) { |
255 | 0 | if (!amSending()) { |
256 | 0 | isc_throw(NcrSenderError, "sender is not ready to send"); |
257 | 0 | } |
258 | | |
259 | 0 | if (!ncr) { |
260 | 0 | isc_throw(NcrSenderError, "request to send is empty"); |
261 | 0 | } |
262 | | |
263 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
264 | 0 | lock_guard<mutex> lock(*mutex_); |
265 | 0 | sendRequestInternal(ncr); |
266 | 0 | } else { |
267 | 0 | sendRequestInternal(ncr); |
268 | 0 | } |
269 | 0 | } |
270 | | |
271 | | void |
272 | 0 | NameChangeSender::sendRequestInternal(NameChangeRequestPtr& ncr) { |
273 | 0 | if (send_queue_.size() >= send_queue_max_) { |
274 | 0 | isc_throw(NcrSenderQueueFull, |
275 | 0 | "send queue has reached maximum capacity: " |
276 | 0 | << send_queue_max_); |
277 | 0 | } |
278 | | |
279 | | // Put it on the queue. |
280 | 0 | send_queue_.push_back(ncr); |
281 | | |
282 | | // Call sendNext to schedule the next one to go. |
283 | 0 | sendNext(); |
284 | 0 | } |
285 | | |
286 | | void |
287 | 0 | NameChangeSender::sendNext() { |
288 | 0 | 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 | 0 | return; |
293 | 0 | } |
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 | 0 | if (!send_queue_.empty()) { |
298 | 0 | 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 | 0 | try { |
309 | 0 | doSend(ncr_to_send_); |
310 | 0 | } catch (const std::exception& ex) { |
311 | 0 | LOG_ERROR(dhcp_ddns_logger, DHCP_DDNS_NCR_SEND_NEXT_ERROR) |
312 | 0 | .arg(ex.what()); |
313 | 0 | send_queue_.pop_front(); |
314 | | // Use the internal path: sendNext() may already run under lock. |
315 | 0 | invokeSendHandlerInternal(ERROR); |
316 | 0 | } |
317 | 0 | } |
318 | 0 | } |
319 | | |
320 | | void |
321 | 0 | NameChangeSender::invokeSendHandler(const NameChangeSender::Result result) { |
322 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
323 | 0 | lock_guard<mutex> lock(*mutex_); |
324 | 0 | invokeSendHandlerInternal(result); |
325 | 0 | } else { |
326 | 0 | invokeSendHandlerInternal(result); |
327 | 0 | } |
328 | 0 | } |
329 | | |
330 | | void |
331 | 0 | NameChangeSender::invokeSendHandlerInternal(const NameChangeSender::Result result) { |
332 | | // @todo reset defense timer |
333 | 0 | if (result == SUCCESS) { |
334 | | // It shipped so pull it off the queue. |
335 | 0 | send_queue_.pop_front(); |
336 | 0 | } |
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 | 0 | try { |
344 | 0 | (*send_handler_)(result, ncr_to_send_); |
345 | 0 | } catch (const std::exception& ex) { |
346 | 0 | LOG_ERROR(dhcp_ddns_logger, DHCP_DDNS_UNCAUGHT_NCR_SEND_HANDLER_ERROR) |
347 | 0 | .arg(ex.what()); |
348 | 0 | } |
349 | | |
350 | | // Clear the pending ncr pointer. |
351 | 0 | ncr_to_send_.reset(); |
352 | | |
353 | | // Set up the next send |
354 | 0 | try { |
355 | 0 | if (amSending()) { |
356 | 0 | sendNext(); |
357 | 0 | } |
358 | 0 | } 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. |
365 | 0 | LOG_ERROR(dhcp_ddns_logger, DHCP_DDNS_NCR_SEND_NEXT_ERROR) |
366 | 0 | .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 | 0 | try { |
373 | 0 | (*send_handler_)(ERROR, ncr_to_send_); |
374 | 0 | } catch (const std::exception& std_ex) { |
375 | 0 | LOG_ERROR(dhcp_ddns_logger, |
376 | 0 | DHCP_DDNS_UNCAUGHT_NCR_SEND_HANDLER_ERROR).arg(std_ex.what()); |
377 | 0 | } |
378 | 0 | } |
379 | 0 | } |
380 | | |
381 | | void |
382 | 0 | NameChangeSender::skipNext() { |
383 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
384 | 0 | lock_guard<mutex> lock(*mutex_); |
385 | 0 | skipNextInternal(); |
386 | 0 | } else { |
387 | 0 | skipNextInternal(); |
388 | 0 | } |
389 | 0 | } |
390 | | |
391 | | void |
392 | 0 | NameChangeSender::skipNextInternal() { |
393 | 0 | if (!send_queue_.empty()) { |
394 | | // Discards the request at the front of the queue. |
395 | 0 | send_queue_.pop_front(); |
396 | 0 | } |
397 | 0 | } |
398 | | |
399 | | void |
400 | 0 | NameChangeSender::clearSendQueue() { |
401 | 0 | if (amSending()) { |
402 | 0 | isc_throw(NcrSenderError, "Cannot clear queue while sending"); |
403 | 0 | } |
404 | | |
405 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
406 | 0 | lock_guard<mutex> lock(*mutex_); |
407 | 0 | send_queue_.clear(); |
408 | 0 | } else { |
409 | 0 | send_queue_.clear(); |
410 | 0 | } |
411 | 0 | } |
412 | | |
413 | | void |
414 | 0 | NameChangeSender::setQueueMaxSize(const size_t new_max) { |
415 | 0 | if (new_max == 0) { |
416 | 0 | isc_throw(NcrSenderError, "NameChangeSender:" |
417 | 0 | " queue size must be greater than zero"); |
418 | 0 | } |
419 | | |
420 | 0 | send_queue_max_ = new_max; |
421 | 0 | } |
422 | | |
423 | | size_t |
424 | 0 | NameChangeSender::getQueueSize() const { |
425 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
426 | 0 | lock_guard<mutex> lock(*mutex_); |
427 | 0 | return (getQueueSizeInternal()); |
428 | 0 | } else { |
429 | 0 | return (getQueueSizeInternal()); |
430 | 0 | } |
431 | 0 | } |
432 | | |
433 | | size_t |
434 | 0 | NameChangeSender::getQueueSizeInternal() const { |
435 | 0 | return (send_queue_.size()); |
436 | 0 | } |
437 | | |
438 | | const NameChangeRequestPtr& |
439 | 0 | NameChangeSender::peekAt(const size_t index) const { |
440 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
441 | 0 | lock_guard<mutex> lock(*mutex_); |
442 | 0 | return (peekAtInternal(index)); |
443 | 0 | } else { |
444 | 0 | return (peekAtInternal(index)); |
445 | 0 | } |
446 | 0 | } |
447 | | |
448 | | const NameChangeRequestPtr& |
449 | 0 | NameChangeSender::peekAtInternal(const size_t index) const { |
450 | 0 | auto size = getQueueSizeInternal(); |
451 | 0 | if (index >= size) { |
452 | 0 | isc_throw(NcrSenderError, |
453 | 0 | "NameChangeSender::peekAt peek beyond end of queue attempted" |
454 | 0 | << " index: " << index << " queue size: " << size); |
455 | 0 | } |
456 | | |
457 | 0 | return (send_queue_.at(index)); |
458 | 0 | } |
459 | | |
460 | | bool |
461 | 0 | NameChangeSender::isSendInProgress() const { |
462 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
463 | 0 | lock_guard<mutex> lock(*mutex_); |
464 | 0 | return ((ncr_to_send_) ? true : false); |
465 | 0 | } else { |
466 | 0 | return ((ncr_to_send_) ? true : false); |
467 | 0 | } |
468 | 0 | } |
469 | | |
470 | | void |
471 | 0 | NameChangeSender::assumeQueue(NameChangeSender& source_sender) { |
472 | 0 | if (source_sender.amSending()) { |
473 | 0 | isc_throw(NcrSenderError, "Cannot assume queue:" |
474 | 0 | " source sender is actively sending"); |
475 | 0 | } |
476 | | |
477 | 0 | if (amSending()) { |
478 | 0 | isc_throw(NcrSenderError, "Cannot assume queue:" |
479 | 0 | " target sender is actively sending"); |
480 | 0 | } |
481 | | |
482 | 0 | if (getQueueMaxSize() < source_sender.getQueueSize()) { |
483 | 0 | isc_throw(NcrSenderError, "Cannot assume queue:" |
484 | 0 | " source queue count exceeds target queue max"); |
485 | 0 | } |
486 | | |
487 | 0 | if (MultiThreadingMgr::instance().getMode()) { |
488 | 0 | lock_guard<mutex> lock(*mutex_); |
489 | 0 | assumeQueueInternal(source_sender); |
490 | 0 | } else { |
491 | 0 | assumeQueueInternal(source_sender); |
492 | 0 | } |
493 | 0 | } |
494 | | |
495 | | void |
496 | 0 | NameChangeSender::assumeQueueInternal(NameChangeSender& source_sender) { |
497 | 0 | if (!send_queue_.empty()) { |
498 | 0 | isc_throw(NcrSenderError, "Cannot assume queue:" |
499 | 0 | " target queue is not empty"); |
500 | 0 | } |
501 | | |
502 | 0 | send_queue_.swap(source_sender.getSendQueue()); |
503 | 0 | } |
504 | | |
505 | | int |
506 | 0 | NameChangeSender::getSelectFd() { |
507 | 0 | isc_throw(NotImplemented, "NameChangeSender::getSelectFd is not supported"); |
508 | 0 | } |
509 | | |
510 | | void |
511 | 0 | NameChangeSender::runReadyIO() { |
512 | 0 | if (!io_service_) { |
513 | 0 | isc_throw(NcrSenderError, "NameChangeSender::runReadyIO" |
514 | 0 | " sender io service is null"); |
515 | 0 | } |
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 | 0 | io_service_->pollOne(); |
520 | 0 | } |
521 | | |
522 | | } // namespace dhcp_ddns |
523 | | } // namespace isc |