Coverage Report

Created: 2026-09-17 07:23

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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