Kea 3.2.1-git
d2_queue_mgr.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>
8#include <d2/d2_queue_mgr.h>
9#include <d2srv/d2_log.h>
10#include <dhcp_ddns/ncr_udp.h>
11#include <stats/stats_mgr.h>
12
13using namespace isc::stats;
14
15namespace isc {
16namespace d2 {
17
18// Makes constant visible to Google test macros.
20
21D2QueueMgr::D2QueueMgr(asiolink::IOServicePtr& io_service, const size_t max_queue_size)
22 : io_service_(io_service), max_queue_size_(max_queue_size),
23 mgr_state_(NOT_INITTED), target_stop_state_(NOT_INITTED) {
24 if (!io_service_) {
25 isc_throw(D2QueueMgrError, "IOServicePtr cannot be null");
26 }
27
28 // Use setter to do validation.
29 setMaxQueueSize(max_queue_size);
30}
31
34
35void
38 try {
39 // Note that error conditions must be handled here without throwing
40 // exceptions. Remember this is the application level "link" in the
41 // callback chain. Throwing an exception here will "break" the
42 // io_service "run" we are operating under. With that in mind,
43 // if we hit a problem, we will stop the listener transition to
44 // the appropriate stopped state. Upper layer(s) must monitor our
45 // state as well as our queue size.
46 switch (result) {
48 {
49 // Receive was successful, attempt to queue the request(s).
50 // If it is a chain, so how many there are.
51 size_t chain_length = 0;
52 auto tmp = ncr;
53 while (tmp) {
54 ++chain_length;
55 tmp = tmp->getNextNcr();
56 }
57
58 if ((getQueueSize() + chain_length) > getMaxQueueSize()) {
59 // Won't fit in the queue, stop the listener.
60 // Note that we can move straight to a STOPPED state as there
61 // is no receive in progress.
63 .arg(max_queue_size_);
64 StatsMgr::instance().addValue("queue-mgr-queue-full", static_cast<int64_t>(1));
66 break;
67 }
68
69 bool nested = false;
70 do {
71 // Add the NCR to the queue.
72 enqueue(ncr);
73
74 // Log that we got the request
78 .arg(ncr->getRequestId());
79
80 if (nested) {
81 StatsMgr::instance().addValue("ncr-received", static_cast<int64_t>(1));
82 }
83
84 ncr = ncr->getNextNcr();
85 nested = !!ncr;
86 } while (ncr);
87
88 break;
89 }
91 if (mgr_state_ == STOPPING) {
92 // This is confirmation that the listener has stopped and its
93 // callback will not be called again, unless its restarted.
94 updateStopState();
95 } else {
96 // We should not get a receive complete status of stopped
97 // unless we canceled the read as part of stopping. Therefore
98 // this is unexpected so we will treat it as a receive error.
99 // This is most likely an unforeseen programmatic issue.
101 .arg(D2QueueMgr::stateToText(mgr_state_));
103 }
104
105 break;
106
107 default:
108 // Receive failed, stop the listener.
109 // Note that we can move straight to a STOPPED state as there
110 // is no receive in progress.
113 break;
114 }
115 } catch (const std::exception& ex) {
116 // On the outside chance a throw occurs, let's log it and swallow it.
118 .arg(ex.what());
119 }
120}
121
122void
124 const uint32_t port,
125 const dhcp_ddns::NameChangeFormat format,
126 const bool reuse_address) {
127
128 if (listener_) {
130 "D2QueueMgr listener is already initialized");
131 }
132
133 // Instantiate a UDP listener and set state to INITTED.
134 // Note UDP listener constructor does not throw.
135 listener_.reset(new dhcp_ddns::NameChangeUDPListener(ip_address, port, format,
136 shared_from_this(), reuse_address));
137 mgr_state_ = INITTED;
138}
139
140void
142 // We can't listen if we haven't initialized the listener yet.
143 if (!listener_) {
144 isc_throw(D2QueueMgrError, "D2QueueMgr "
145 "listener is not initialized, cannot start listening");
146 }
147
148 // If we are already listening, we do not want to "reopen" the listener
149 // and really we shouldn't be trying.
150 if (mgr_state_ == RUNNING) {
151 isc_throw(D2QueueMgrError, "D2QueueMgr "
152 "cannot call startListening from the RUNNING state");
153 }
154
155 // Instruct the listener to start listening and set state accordingly.
156 try {
157 listener_->startListening(io_service_);
158 mgr_state_ = RUNNING;
159 } catch (const isc::Exception& ex) {
160 isc_throw(D2QueueMgrError, "D2QueueMgr listener start failed: "
161 << ex.what());
162 }
163
166}
167
168void
169D2QueueMgr::stopListening(const State target_stop_state) {
170 if (listener_) {
171 // Enforce only valid "stop" states.
172 // This is purely a programmatic error and should never happen.
173 if (target_stop_state != STOPPED &&
174 target_stop_state != STOPPED_QUEUE_FULL &&
175 target_stop_state != STOPPED_RECV_ERROR) {
177 "D2QueueMgr invalid value for stop state: "
178 << target_stop_state);
179 }
180
181 // Remember the state we want to achieve.
182 target_stop_state_ = target_stop_state;
183
184 // Instruct the listener to stop. If the listener reports that it
185 // has IO pending, then we transition to STOPPING to wait for the
186 // cancellation event. Otherwise, we can move directly to the targeted
187 // state.
188 listener_->stopListening();
189 if (listener_->isIoPending()) {
190 mgr_state_ = STOPPING;
191 } else {
192 updateStopState();
193 }
194 }
195}
196
197void
198D2QueueMgr::updateStopState() {
199 mgr_state_ = target_stop_state_;
202}
203
204void
206 // Force our managing layer(s) to stop us properly first.
207 if (mgr_state_ == RUNNING) {
209 "D2QueueMgr cannot delete listener while state is RUNNING");
210 }
211
212 listener_.reset();
213 mgr_state_ = NOT_INITTED;
214}
215
218 if (getQueueSize() == 0) {
220 "D2QueueMgr peek attempted on an empty queue");
221 }
222
223 return (ncr_queue_.front());
224}
225
227D2QueueMgr::peekAt(const size_t index) const {
228 if (index >= getQueueSize()) {
230 "D2QueueMgr peek beyond end of queue attempted"
231 << " index: " << index << " queue size: " << getQueueSize());
232 }
233
234 return (ncr_queue_.at(index));
235}
236
237void
238D2QueueMgr::dequeueAt(const size_t index) {
239 if (index >= getQueueSize()) {
241 "D2QueueMgr dequeue beyond end of queue attempted"
242 << " index: " << index << " queue size: " << getQueueSize());
243 }
244
245 RequestQueue::iterator pos = ncr_queue_.begin() + index;
246 ncr_queue_.erase(pos);
247}
248
249void
251 if (getQueueSize() == 0) {
253 "D2QueueMgr dequeue attempted on an empty queue");
254 }
255
256 ncr_queue_.pop_front();
257}
258
259void
261 ncr_queue_.push_back(ncr);
262}
263
264void
266 ncr_queue_.clear();
267}
268
269void
270D2QueueMgr::setMaxQueueSize(const size_t new_queue_max) {
271 if (new_queue_max < 1) {
273 "D2QueueMgr maximum queue size must be greater than zero");
274 }
275
276 if (new_queue_max < getQueueSize()) {
277 isc_throw(D2QueueMgrError, "D2QueueMgr maximum queue size value cannot"
278 " be less than the current queue size :" << getQueueSize());
279 }
280
281 max_queue_size_ = new_queue_max;
282}
283
284} // namespace isc::d2
285} // namespace isc
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.
Thrown if the queue manager encounters a general error.
Thrown if a queue index is beyond the end of the queue.
Thrown if the request queue empty and a read is attempted.
const dhcp_ddns::NameChangeRequestPtr & peek() const
Returns the entry at the front of the queue.
D2QueueMgr(asiolink::IOServicePtr &io_service, const size_t max_queue_size=MAX_QUEUE_DEFAULT)
Constructor.
State
Defines the list of possible states for D2QueueMgr.
virtual ~D2QueueMgr()
Destructor.
size_t getMaxQueueSize() const
Returns the maximum number of entries allowed in the queue.
static const size_t MAX_QUEUE_DEFAULT
Maximum number of entries allowed in the request queue.
void dequeue()
Removes the entry at the front of the queue.
const dhcp_ddns::NameChangeRequestPtr & peekAt(const size_t index) const
Returns the entry at a given position in the queue.
static std::string const & stateToText(State const &state)
Convert enum to string.
void removeListener()
Deletes the current listener.
void enqueue(dhcp_ddns::NameChangeRequestPtr &ncr)
Adds a request to the end of the queue.
void startListening()
Starts actively listening for requests.
void setMaxQueueSize(const size_t max_queue_size)
Sets the maximum number of entries allowed in the queue.
size_t getQueueSize() const
Returns the number of entries in the queue.
virtual void operator()(const dhcp_ddns::NameChangeListener::Result result, dhcp_ddns::NameChangeRequestPtr &ncr)
Function operator implementing the NCR receive callback.
void clearQueue()
Removes all entries from the queue.
void initUDPListener(const isc::asiolink::IOAddress &ip_address, const uint32_t port, const dhcp_ddns::NameChangeFormat format, const bool reuse_address=false)
Initializes the listener as a UDP listener.
void dequeueAt(const size_t index)
Removes the entry at a given position in the queue.
void stopListening(const State target_stop_state=STOPPED)
Stops listening for requests.
Result
Defines the outcome of an asynchronous NCR receive.
Definition ncr_io.h:171
Provides the ability to receive NameChangeRequests via UDP socket.
Definition ncr_udp.h:317
static StatsMgr & instance()
Statistics Manager accessor method.
This file defines the class D2QueueMgr.
#define isc_throw(type, stream)
A shortcut macro to insert known values into exception arguments.
void addValue(const std::string &name, const int64_t value)
Records incremental integer observation.
#define LOG_ERROR(LOGGER, MESSAGE)
Macro to conveniently test error output and log it.
Definition macros.h:32
#define LOG_DEBUG(LOGGER, LEVEL, MESSAGE)
Macro to conveniently test debug output and log it.
Definition macros.h:14
const isc::log::MessageID DHCP_DDNS_QUEUE_MGR_RECV_ERROR
Definition d2_messages.h:57
const isc::log::MessageID DHCP_DDNS_QUEUE_MGR_QUEUE_RECEIVE
Definition d2_messages.h:54
const isc::log::MessageID DHCP_DDNS_QUEUE_MGR_UNEXPECTED_HANDLER_ERROR
Definition d2_messages.h:65
const isc::log::MessageID DHCP_DDNS_QUEUE_MGR_UNEXPECTED_STOP
Definition d2_messages.h:66
const isc::log::MessageID DHCP_DDNS_QUEUE_MGR_STARTED
Definition d2_messages.h:60
isc::log::Logger dhcp_to_d2_logger("dhcp-to-d2")
Definition d2_log.h:19
isc::log::Logger d2_logger("dhcpddns")
Defines the logger used within D2.
Definition d2_log.h:18
const isc::log::MessageID DHCP_DDNS_QUEUE_MGR_QUEUE_FULL
Definition d2_messages.h:53
const isc::log::MessageID DHCP_DDNS_QUEUE_MGR_STOPPED
Definition d2_messages.h:62
NameChangeFormat
Defines the list of data wire formats supported.
Definition ncr_msg.h:60
boost::shared_ptr< NameChangeRequest > NameChangeRequestPtr
Defines a pointer to a NameChangeRequest.
Definition ncr_msg.h:242
const int DBGLVL_TRACE_BASIC
Trace basic operations.
const int DBGLVL_START_SHUT
This is given a value of 0 as that is the level selected if debugging is enabled without giving a lev...
const int DBGLVL_TRACE_DETAIL_DATA
Trace data associated with detailed operations.
Defines the logger used by the top-level component of kea-lfc.
This file provides UDP socket based implementation for sending and receiving NameChangeRequests.