Kea 3.3.3
communication_state.cc
Go to the documentation of this file.
1// Copyright (C) 2018-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
9#include <cc/data.h>
10#include <dhcp/dhcp4.h>
11#include <dhcp/dhcp6.h>
12#include <dhcp/option_int.h>
13#include <dhcp/pkt4.h>
14#include <dhcp/pkt6.h>
18#include <http/date_time.h>
21
22#include <ctime>
23#include <functional>
24#include <limits>
25#include <sstream>
26#include <utility>
27
28#include <boost/pointer_cast.hpp>
29
30#include <ha_log.h>
31
32using namespace isc::asiolink;
33using namespace isc::data;
34using namespace isc::dhcp;
35using namespace isc::http;
36using namespace isc::log;
37using namespace isc::util;
38
39using namespace boost::posix_time;
40using namespace std;
41
42namespace {
43
45constexpr long WARN_CLOCK_SKEW = 30;
46
48constexpr long TERM_CLOCK_SKEW = 60;
49
51constexpr long MIN_TIME_SINCE_CLOCK_SKEW_WARN = 60;
52
53}
54
55namespace isc {
56namespace ha {
57
59 const HAConfigPtr& config)
60 : io_service_(io_service), config_(config), timer_(), interval_(0),
61 poke_time_(boost::posix_time::microsec_clock::universal_time()),
66 partner_unsent_update_count_{0, 0}, mutex_(new mutex()) {
67}
68
72
73void
75 if (MultiThreadingMgr::instance().getMode()) {
76 std::lock_guard<std::mutex> lk(*mutex_);
77 poke_time_ += boost::posix_time::seconds(secs);
78 } else {
79 poke_time_ += boost::posix_time::seconds(secs);
80 }
81}
82
83int
85 if (MultiThreadingMgr::instance().getMode()) {
86 std::lock_guard<std::mutex> lk(*mutex_);
87 return (partner_state_);
88 } else {
89 return (partner_state_);
90 }
91}
92
93void
94CommunicationState::setPartnerState(const std::string& state) {
95 if (MultiThreadingMgr::instance().getMode()) {
96 std::lock_guard<std::mutex> lk(*mutex_);
97 setPartnerStateInternal(state);
98 } else {
99 setPartnerStateInternal(state);
100 }
101}
102
103void
105 if (MultiThreadingMgr::instance().getMode()) {
106 std::lock_guard<std::mutex> lk(*mutex_);
107 setPartnerStateInternal("unavailable");
108 resetPartnerTimeInternal();
109 } else {
110 setPartnerStateInternal("unavailable");
111 resetPartnerTimeInternal();
112 }
113}
114
115void
116CommunicationState::setPartnerStateInternal(const std::string& state) {
117 try {
118 auto new_partner_state = stringToState(state);
119 if (new_partner_state != partner_state_) {
120 setCurrentPartnerStateTimeInternal();
121 }
122 partner_state_ = new_partner_state;
123 } catch (...) {
124 isc_throw(BadValue, "unsupported HA partner state returned "
125 << state);
126 }
127}
128
129time_duration
131 ptime now = boost::posix_time::microsec_clock::universal_time();
132 if (MultiThreadingMgr::instance().getMode()) {
133 std::lock_guard<std::mutex> lk(*mutex_);
134 return (now - partner_state_time_);
135 } else {
136 return (now - partner_state_time_);
137 }
138}
139
140void
141CommunicationState::setCurrentPartnerStateTimeInternal() {
142 partner_state_time_ = boost::posix_time::microsec_clock::universal_time();
143}
144
145std::set<std::string>
147 if (MultiThreadingMgr::instance().getMode()) {
148 std::lock_guard<std::mutex> lk(*mutex_);
149 return (partner_scopes_);
150 } else {
151 return (partner_scopes_);
152 }
153}
154
155void
157 if (MultiThreadingMgr::instance().getMode()) {
158 std::lock_guard<std::mutex> lk(*mutex_);
159 setPartnerScopesInternal(new_scopes);
160 } else {
161 setPartnerScopesInternal(new_scopes);
162 }
163}
164
165void
166CommunicationState::setPartnerScopesInternal(ConstElementPtr new_scopes) {
167 if (!new_scopes || (new_scopes->getType() != Element::list)) {
168 isc_throw(BadValue, "unable to record partner's HA scopes because"
169 " the received value is not a valid JSON list");
170 }
171
172 std::set<std::string> partner_scopes;
173 for (unsigned i = 0; i < new_scopes->size(); ++i) {
174 auto scope = new_scopes->get(i);
175 if (scope->getType() != Element::string) {
176 isc_throw(BadValue, "unable to record partner's HA scopes because"
177 " the received scope value is not a valid JSON string");
178 }
179 auto scope_str = scope->stringValue();
180 if (!scope_str.empty()) {
181 partner_scopes.insert(scope_str);
182 }
183 }
184 partner_scopes_ = partner_scopes;
185}
186
187void
189 const std::function<void()>& heartbeat_impl) {
190 if (MultiThreadingMgr::instance().getMode()) {
191 std::lock_guard<std::mutex> lk(*mutex_);
192 startHeartbeatInternal(interval, heartbeat_impl);
193 } else {
194 startHeartbeatInternal(interval, heartbeat_impl);
195 }
196}
197
198void
199CommunicationState::startHeartbeatInternal(const long interval,
200 const std::function<void()>& heartbeat_impl) {
201 bool settings_modified = false;
202
203 // If we're setting the heartbeat for the first time, it should
204 // be non-null.
205 if (heartbeat_impl) {
206 settings_modified = true;
207 heartbeat_impl_ = heartbeat_impl;
208
209 } else if (!heartbeat_impl_) {
210 // The heartbeat is re-scheduled but we have no historic implementation
211 // pointer we could re-use. This is a programmatic issue.
212 isc_throw(BadValue, "unable to start heartbeat when pointer"
213 " to the heartbeat implementation is not specified");
214 }
215
216 // If we're setting the heartbeat for the first time, the interval
217 // should be greater than 0.
218 if (interval != 0) {
219 settings_modified |= (interval_ != interval);
220 interval_ = interval;
221
222 } else if (interval_ <= 0) {
223 // The heartbeat is re-scheduled but we have no historic interval
224 // which we could re-use. This is a programmatic issue.
225 heartbeat_impl_ = 0;
226 isc_throw(BadValue, "unable to start heartbeat when interval"
227 " for the heartbeat timer is not specified");
228 }
229
230 if (!timer_) {
231 timer_.reset(new IntervalTimer(io_service_));
232 }
233
234 if (settings_modified) {
236 }
237}
238
239void
241 if (MultiThreadingMgr::instance().getMode()) {
242 std::lock_guard<std::mutex> lk(*mutex_);
243 stopHeartbeatInternal();
244 } else {
245 stopHeartbeatInternal();
246 }
247}
248
249void
250CommunicationState::stopHeartbeatInternal() {
251 if (timer_) {
252 timer_->cancel();
253 timer_.reset();
254 interval_ = 0;
255 heartbeat_impl_ = 0;
256 }
257}
258
259bool
261 if (MultiThreadingMgr::instance().getMode()) {
262 std::lock_guard<std::mutex> lk(*mutex_);
263 return (static_cast<bool>(timer_));
264 } else {
265 return (static_cast<bool>(timer_));
266 }
267}
268
269boost::posix_time::time_duration
271 if (MultiThreadingMgr::instance().getMode()) {
272 std::lock_guard<std::mutex> lk(*mutex_);
273 return (updatePokeTimeInternal());
274 } else {
275 return (updatePokeTimeInternal());
276 }
277}
278
279boost::posix_time::time_duration
280CommunicationState::updatePokeTimeInternal() {
281 // Remember previous poke time.
282 boost::posix_time::ptime prev_poke_time = poke_time_;
283 // Set poke time to the current time.
284 poke_time_ = boost::posix_time::microsec_clock::universal_time();
285 return (poke_time_ - prev_poke_time);
286}
287
288void
290 if (MultiThreadingMgr::instance().getMode()) {
291 std::lock_guard<std::mutex> lk(*mutex_);
292 pokeInternal();
293 } else {
294 pokeInternal();
295 }
296}
297
298void
299CommunicationState::pokeInternal() {
300 // Update poke time and compute duration.
301 boost::posix_time::time_duration duration_since_poke = updatePokeTimeInternal();
302
303 // If we have been tracking the DHCP messages directed to the partner,
304 // we need to clear any gathered information because the connection
305 // seems to be (re)established.
308
309 if (timer_) {
310 // Check the duration since last poke. If it is less than a second, we don't
311 // want to reschedule the timer. In order to avoid the overhead of
312 // re-scheduling the timer too frequently we reschedule it only if the
313 // duration is 1s or more. This matches the time resolution for heartbeats.
314 if (duration_since_poke.total_seconds() > 0) {
315 // A poke causes the timer to be re-scheduled to prevent it
316 // from triggering a heartbeat shortly after confirming the
317 // connection is ok.
318 startHeartbeatInternal();
319 }
320 }
321}
322
323int64_t
325 if (MultiThreadingMgr::instance().getMode()) {
326 std::lock_guard<std::mutex> lk(*mutex_);
327 return (getDurationInMillisecsInternal());
328 } else {
329 return (getDurationInMillisecsInternal());
330 }
331}
332
333int64_t
334CommunicationState::getDurationInMillisecsInternal() const {
335 ptime now = boost::posix_time::microsec_clock::universal_time();
336 time_duration duration = now - poke_time_;
337 return (duration.total_milliseconds());
338}
339
340bool
342 return (getDurationInMillisecs() > config_->getMaxResponseDelay());
343}
344
345std::vector<uint8_t>
347 const uint16_t option_type) {
348 std::vector<uint8_t> client_id;
349 OptionPtr opt_client_id = message->getOption(option_type);
350 if (opt_client_id) {
351 client_id = opt_client_id->getData();
352 }
353 return (client_id);
354}
355
356size_t
360
361size_t
363 if (MultiThreadingMgr::instance().getMode()) {
364 std::lock_guard<std::mutex> lk(*mutex_);
366 } else {
368 }
369}
370
371bool
373 const uint32_t lifetime) {
374 if (MultiThreadingMgr::instance().getMode()) {
375 std::lock_guard<std::mutex> lk(*mutex_);
376 return (reportRejectedLeaseUpdateInternal(message, lifetime));
377 } else {
378 return (reportRejectedLeaseUpdateInternal(message, lifetime));
379 }
380}
381
382bool
384 if (MultiThreadingMgr::instance().getMode()) {
385 std::lock_guard<std::mutex> lk(*mutex_);
386 return (reportSuccessfulLeaseUpdateInternal(message));
387 } else {
388 return (reportSuccessfulLeaseUpdateInternal(message));
389 }
390}
391
392void
394 if (MultiThreadingMgr::instance().getMode()) {
395 std::lock_guard<std::mutex> lk(*mutex_);
397 } else {
399 }
400}
401
402bool
404 if (MultiThreadingMgr::instance().getMode()) {
405 std::lock_guard<std::mutex> lk(*mutex_);
406 return (clockSkewShouldWarnInternal());
407 } else {
408 return (clockSkewShouldWarnInternal());
409 }
410}
411
412bool
413CommunicationState::clockSkewShouldWarnInternal() {
414 // First check if the clock skew is beyond the threshold.
415 if (isClockSkewGreater(WARN_CLOCK_SKEW)) {
416
417 // In order to prevent to frequent warnings we provide a gating mechanism
418 // which doesn't allow for issuing a warning earlier than 60 seconds after
419 // the previous one.
420
421 // Find the current time and the duration since last warning.
422 ptime now = boost::posix_time::microsec_clock::universal_time();
423 time_duration since_warn_duration = now - last_clock_skew_warn_;
424
425 // If the last warning was issued more than 60 seconds ago or it is a
426 // first warning, we need to update the last warning timestamp and return
427 // true to indicate that new warning should be issued.
428 if (last_clock_skew_warn_.is_not_a_date_time() ||
429 (since_warn_duration.total_seconds() > MIN_TIME_SINCE_CLOCK_SKEW_WARN)) {
432 .arg(config_->getThisServerName())
433 .arg(logFormatClockSkewInternal());
434 return (true);
435 }
436 }
437
438 // The warning should not be issued.
439 return (false);
440}
441
442bool
444 if (MultiThreadingMgr::instance().getMode()) {
445 std::lock_guard<std::mutex> lk(*mutex_);
446 // Issue a warning if the clock skew is greater than 60s.
447 return (clockSkewShouldTerminateInternal());
448 } else {
449 return (clockSkewShouldTerminateInternal());
450 }
451}
452
453bool
454CommunicationState::clockSkewShouldTerminateInternal() {
455 if (isClockSkewGreater(TERM_CLOCK_SKEW)) {
457 .arg(config_->getThisServerName())
458 .arg(logFormatClockSkewInternal());
459 return (true);
460 }
461 return (false);
462}
463
464bool
466 if (MultiThreadingMgr::instance().getMode()) {
467 std::lock_guard<std::mutex> lk(*mutex_);
468 return (rejectedLeaseUpdatesShouldTerminateInternal());
469 } else {
470 return (rejectedLeaseUpdatesShouldTerminateInternal());
471 }
472}
473
474bool
475CommunicationState::rejectedLeaseUpdatesShouldTerminateInternal() {
476 if (config_->getMaxRejectedLeaseUpdates() &&
477 (config_->getMaxRejectedLeaseUpdates() <= getRejectedLeaseUpdatesCountInternal())) {
479 .arg(config_->getThisServerName());
480 return (true);
481 }
482 return (false);
483}
484
485bool
486CommunicationState::isClockSkewGreater(const long seconds) const {
487 return ((clock_skew_.total_seconds() > seconds) ||
488 (clock_skew_.total_seconds() < -seconds));
489}
490
491void
492CommunicationState::setPartnerTime(const std::string& time_text) {
493 if (MultiThreadingMgr::instance().getMode()) {
494 std::lock_guard<std::mutex> lk(*mutex_);
495 setPartnerTimeInternal(time_text);
496 } else {
497 setPartnerTimeInternal(time_text);
498 }
499}
500
501void
502CommunicationState::setPartnerTimeInternal(const std::string& time_text) {
506}
507
508void
509CommunicationState::resetPartnerTimeInternal() {
510 clock_skew_ = boost::posix_time::time_duration(0, 0, 0, 0);
511 last_clock_skew_warn_ = boost::posix_time::ptime();
512 my_time_at_skew_ = boost::posix_time::ptime();
513 partner_time_at_skew_ = boost::posix_time::ptime();
514}
515
516std::string
518 if (MultiThreadingMgr::instance().getMode()) {
519 std::lock_guard<std::mutex> lk(*mutex_);
520 return (logFormatClockSkewInternal());
521 } else {
522 return (logFormatClockSkewInternal());
523 }
524}
525
526std::string
527CommunicationState::logFormatClockSkewInternal() const {
528 std::ostringstream os;
529
530 if ((my_time_at_skew_.is_not_a_date_time()) ||
531 (partner_time_at_skew_.is_not_a_date_time())) {
532 // Guard against being called before times have been set.
533 // Otherwise we'll get out-range exceptions.
534 return ("skew not initialized");
535 }
536
537 // Note HttpTime resolution is only to seconds, so we use fractional
538 // precision of zero when logging.
539 os << "my time: " << ptimeToText(my_time_at_skew_, 0)
540 << ", partner's time: " << ptimeToText(partner_time_at_skew_, 0)
541 << ", partner's clock is ";
542
543 if (clock_skew_.total_seconds() == 0) {
544 // Most common case.
545 os << "synchronized";
546 } else if (clock_skew_.is_negative()) {
547 // Partner's time is behind our time.
548 os << clock_skew_.invert_sign().total_seconds() << "s behind";
549 } else {
550 // Partner's time is ahead of ours.
551 os << clock_skew_.total_seconds() << "s ahead";
552 }
553
554 return (os.str());
555}
556
559 auto report = Element::createMap();
560
561 auto in_touch = (getPartnerState() > 0);
562 report->set("in-touch", Element::create(in_touch));
563
564 auto age = in_touch ? static_cast<long long int>(getDurationInMillisecs() / 1000) : 0;
565 report->set("age", Element::create(age));
566
567 try {
568 report->set("last-state", Element::create(stateToString(getPartnerState())));
569
570 } catch (...) {
571 report->set("last-state", Element::create(std::string()));
572 }
573
574 auto list = Element::createList();
575 for (auto const& scope : getPartnerScopes()) {
576 list->add(Element::create(scope));
577 }
578 report->set("last-scopes", list);
579 report->set("communication-interrupted",
581 report->set("connecting-clients", Element::create(static_cast<long long>(getConnectingClientsCount())));
582 report->set("unacked-clients", Element::create(static_cast<long long>(getUnackedClientsCount())));
583
584 long long unacked_clients_left = 0;
585 if (isCommunicationInterrupted() && (config_->getMaxUnackedClients() >= getUnackedClientsCount())) {
586 unacked_clients_left = static_cast<long long>(config_->getMaxUnackedClients() -
588 }
589 report->set("unacked-clients-left", Element::create(unacked_clients_left));
590 report->set("analyzed-packets", Element::create(static_cast<long long>(getAnalyzedMessagesCount())));
591 if (partner_time_at_skew_.is_not_a_date_time()) {
592 report->set("system-time", Element::create());
593 report->set("clock-skew", Element::create());
594 } else {
595 report->set("system-time", Element::create(ptimeToText(partner_time_at_skew_, 0)));
596 report->set("clock-skew", Element::create(clock_skew_.total_seconds()));
597 }
598
599 return (report);
600}
601
602uint64_t
604 if (MultiThreadingMgr::instance().getMode()) {
605 std::lock_guard<std::mutex> lk(*mutex_);
606 return (unsent_update_count_);
607 } else {
608 return (unsent_update_count_);
609 }
610}
611
612void
614 if (MultiThreadingMgr::instance().getMode()) {
615 std::lock_guard<std::mutex> lk(*mutex_);
616 increaseUnsentUpdateCountInternal();
617 } else {
618 increaseUnsentUpdateCountInternal();
619 }
620}
621
622void
623CommunicationState::increaseUnsentUpdateCountInternal() {
624 // Protect against setting the incremented value to zero.
625 // The zero value is reserved for a server startup.
626 if (unsent_update_count_ < std::numeric_limits<uint64_t>::max()) {
628 } else {
630 }
631}
632
633bool
635 if (MultiThreadingMgr::instance().getMode()) {
636 std::lock_guard<std::mutex> lk(*mutex_);
637 return (hasPartnerNewUnsentUpdatesInternal());
638 } else {
639 return (hasPartnerNewUnsentUpdatesInternal());
640 }
641}
642
643bool
644CommunicationState::hasPartnerNewUnsentUpdatesInternal() const {
645 return (partner_unsent_update_count_.second > 0 &&
647}
648
649void
651 if (MultiThreadingMgr::instance().getMode()) {
652 std::lock_guard<std::mutex> lk(*mutex_);
653 setPartnerUnsentUpdateCountInternal(unsent_update_count);
654 } else {
655 setPartnerUnsentUpdateCountInternal(unsent_update_count);
656 }
657}
658
659void
660CommunicationState::setPartnerUnsentUpdateCountInternal(uint64_t unsent_update_count) {
662 partner_unsent_update_count_.second = unsent_update_count;
663}
664
665boost::posix_time::ptime
669
670boost::posix_time::ptime
674
680
681void
682CommunicationState4::analyzeMessage(const boost::shared_ptr<dhcp::Pkt>& message) {
683 if (MultiThreadingMgr::instance().getMode()) {
684 std::lock_guard<std::mutex> lk(*mutex_);
685 analyzeMessageInternal(message);
686 } else {
687 analyzeMessageInternal(message);
688 }
689}
690
691void
693 // The DHCP message must successfully cast to a Pkt4 object.
694 Pkt4Ptr msg = boost::dynamic_pointer_cast<Pkt4>(message);
695 if (!msg) {
696 isc_throw(BadValue, "DHCP message to be analyzed is not a DHCPv4 message");
697 }
698
700
701 // Check value of the "secs" field by comparing it with the configured
702 // threshold.
703 uint16_t secs = msg->getSecs();
704
705 // It was observed that some Windows clients may send swapped bytes in the
706 // "secs" field. When the second byte is 0 and the first byte is non-zero
707 // we consider bytes to be swapped and so we correct them.
708 if ((secs > 255) && ((secs & 0xFF) == 0)) {
709 secs = ((secs >> 8) | (secs << 8));
710 }
711
712 // Check the value of the "secs" field. The "secs" field holds a value in
713 // seconds, hence we have to multiple by 1000 to get a value in milliseconds.
714 // If the secs value is above the threshold, it means that the current
715 // client should be considered unacked.
716 auto unacked = (secs * 1000 > config_->getMaxAckDelay());
717
718 // Client identifier will be stored together with the hardware address. It
719 // may remain empty if the client hasn't specified it.
720 auto client_id = getClientId(message, DHO_DHCP_CLIENT_IDENTIFIER);
721 bool log_unacked = false;
722
723 // Check if the given client was already recorded.
724 auto& idx = connecting_clients_.get<0>();
725 auto existing_request = idx.find(boost::make_tuple(msg->getHWAddr()->hwaddr_, client_id));
726 if (existing_request != idx.end()) {
727 // If the client was recorded and was not considered unacked
728 // but it should be considered unacked as a result of processing
729 // this packet, let's update the recorded request to mark the
730 // client unacked.
731 if (!existing_request->unacked_ && unacked) {
732 ConnectingClient4 connecting_client{ msg->getHWAddr()->hwaddr_, client_id, unacked };
733 idx.replace(existing_request, connecting_client);
734 log_unacked = true;
735 }
736
737 } else {
738 // This is the first time we see the packet from this client. Let's
739 // record it.
740 ConnectingClient4 connecting_client{ msg->getHWAddr()->hwaddr_, client_id, unacked };
741 idx.insert(connecting_client);
742 log_unacked = unacked;
743
744 if (!unacked) {
745 // This is the first time we see this client after getting into the
746 // communication interrupted state. But, this client hasn't been
747 // yet trying log enough to be considered unacked.
749 .arg(config_->getThisServerName())
750 .arg(message->getLabel());
751 }
752 }
753
754 // Only log the first time we detect a client is unacked.
755 if (log_unacked) {
756 unsigned unacked_left = 0;
757 unsigned unacked_total = connecting_clients_.get<1>().count(true);
758 if (config_->getMaxUnackedClients() >= unacked_total) {
759 unacked_left = config_->getMaxUnackedClients() - unacked_total + 1;
760 }
762 .arg(config_->getThisServerName())
763 .arg(message->getLabel())
764 .arg(unacked_total)
765 .arg(unacked_left);
766 }
767}
768
769bool
771 if (MultiThreadingMgr::instance().getMode()) {
772 std::lock_guard<std::mutex> lk(*mutex_);
773 return (failureDetectedInternal());
774 } else {
775 return (failureDetectedInternal());
776 }
777}
778
779bool
781 return ((config_->getMaxUnackedClients() == 0) ||
782 (connecting_clients_.get<1>().count(true) >
783 config_->getMaxUnackedClients()));
784}
785
786size_t
788 if (MultiThreadingMgr::instance().getMode()) {
789 std::lock_guard<std::mutex> lk(*mutex_);
790 return (connecting_clients_.size());
791 } else {
792 return (connecting_clients_.size());
793 }
794}
795
796size_t
798 if (MultiThreadingMgr::instance().getMode()) {
799 std::lock_guard<std::mutex> lk(*mutex_);
800 return (connecting_clients_.get<1>().count(true));
801 } else {
802 return (connecting_clients_.get<1>().count(true));
803 }
804}
805
806void
810
811size_t
815
816bool
817CommunicationState4::reportRejectedLeaseUpdateInternal(const PktPtr& message, const uint32_t lifetime) {
818 Pkt4Ptr msg = boost::dynamic_pointer_cast<Pkt4>(message);
819 if (!msg) {
820 isc_throw(BadValue, "DHCP message for which the lease update was rejected is not a DHCPv4 message");
821 }
822 auto client_id = getClientId(message, DHO_DHCP_CLIENT_IDENTIFIER);
823 RejectedClient4 client{ msg->getHWAddr()->hwaddr_, client_id, time(NULL) + lifetime };
824 auto existing_client = rejected_clients_.find(boost::make_tuple(msg->getHWAddr()->hwaddr_, client_id));
825 if (existing_client == rejected_clients_.end()) {
826 rejected_clients_.insert(client);
827 return (true);
828 }
829 rejected_clients_.replace(existing_client, client);
830 return (false);
831}
832
833bool
835 // Early check if there is anything to do.
837 return (false);
838 }
839 Pkt4Ptr msg = boost::dynamic_pointer_cast<Pkt4>(message);
840 if (!msg) {
841 isc_throw(BadValue, "DHCP message for which the lease update was successful is not a DHCPv4 message");
842 }
843 auto client_id = getClientId(msg, DHO_DHCP_CLIENT_IDENTIFIER);
844 auto existing_client = rejected_clients_.find(boost::make_tuple(msg->getHWAddr()->hwaddr_, client_id));
845 if (existing_client != rejected_clients_.end()) {
846 rejected_clients_.erase(existing_client);
847 return (true);
848 }
849 return (false);
850}
851
852void
856
862
863void
864CommunicationState6::analyzeMessage(const boost::shared_ptr<dhcp::Pkt>& message) {
865 if (MultiThreadingMgr::instance().getMode()) {
866 std::lock_guard<std::mutex> lk(*mutex_);
867 analyzeMessageInternal(message);
868 } else {
869 analyzeMessageInternal(message);
870 }
871}
872
873void
874CommunicationState6::analyzeMessageInternal(const boost::shared_ptr<dhcp::Pkt>& message) {
875 // The DHCP message must successfully cast to a Pkt6 object.
876 Pkt6Ptr msg = boost::dynamic_pointer_cast<Pkt6>(message);
877 if (!msg) {
878 isc_throw(BadValue, "DHCP message to be analyzed is not a DHCPv6 message");
879 }
880
882
883 // Check the value of the "elapsed time" option. If it is below the threshold
884 // there is nothing to do. The "elapsed time" option holds the time in
885 // 1/100 of second, hence we have to multiply by 10 to get a value in milliseconds.
886 OptionUint16Ptr elapsed_time = boost::dynamic_pointer_cast<
887 OptionUint16>(msg->getOption(D6O_ELAPSED_TIME));
888 auto unacked = (elapsed_time && elapsed_time->getValue() * 10 > config_->getMaxAckDelay());
889
890 // Get the DUID of the client to see if it hasn't been recorded already.
891 auto duid = getClientId(msg, D6O_CLIENTID);
892 if (duid.empty()) {
893 return;
894 }
895
896 bool log_unacked = false;
897
898 // Check if the given client was already recorded.
899 auto& idx = connecting_clients_.get<0>();
900 auto existing_request = idx.find(duid);
901 if (existing_request != idx.end()) {
902 // If the client was recorded and was not considered unacked
903 // but it should be considered unacked as a result of processing
904 // this packet, let's update the recorded request to mark the
905 // client unacked.
906 if (!existing_request->unacked_ && unacked) {
907 ConnectingClient6 connecting_client{ duid, unacked };
908 idx.replace(existing_request, connecting_client);
909 log_unacked = true;
910 }
911
912 } else {
913 // This is the first time we see the packet from this client. Let's
914 // record it.
915 ConnectingClient6 connecting_client{ duid, unacked };
916 idx.insert(connecting_client);
917 log_unacked = unacked;
918
919 if (!unacked) {
920 // This is the first time we see this client after getting into the
921 // communication interrupted state. But, this client hasn't been
922 // yet trying log enough to be considered unacked.
924 .arg(config_->getThisServerName())
925 .arg(message->getLabel());
926 }
927 }
928
929 // Only log the first time we detect a client is unacked.
930 if (log_unacked) {
931 unsigned unacked_left = 0;
932 unsigned unacked_total = connecting_clients_.get<1>().count(true);
933 if (config_->getMaxUnackedClients() >= unacked_total) {
934 unacked_left = config_->getMaxUnackedClients() - unacked_total + 1;
935 }
937 .arg(config_->getThisServerName())
938 .arg(message->getLabel())
939 .arg(unacked_total)
940 .arg(unacked_left);
941 }
942}
943
944bool
946 if (MultiThreadingMgr::instance().getMode()) {
947 std::lock_guard<std::mutex> lk(*mutex_);
948 return (failureDetectedInternal());
949 } else {
950 return (failureDetectedInternal());
951 }
952}
953
954bool
956 return ((config_->getMaxUnackedClients() == 0) ||
957 (connecting_clients_.get<1>().count(true) >
958 config_->getMaxUnackedClients()));
959}
960
961size_t
963 if (MultiThreadingMgr::instance().getMode()) {
964 std::lock_guard<std::mutex> lk(*mutex_);
965 return (connecting_clients_.size());
966 } else {
967 return (connecting_clients_.size());
968 }
969}
970
971size_t
973 if (MultiThreadingMgr::instance().getMode()) {
974 std::lock_guard<std::mutex> lk(*mutex_);
975 return (connecting_clients_.get<1>().count(true));
976 } else {
977 return (connecting_clients_.get<1>().count(true));
978 }
979}
980
981void
985
986size_t
990
991bool
992CommunicationState6::reportRejectedLeaseUpdateInternal(const PktPtr& message, const uint32_t lifetime) {
993 Pkt6Ptr msg = boost::dynamic_pointer_cast<Pkt6>(message);
994 if (!msg) {
995 isc_throw(BadValue, "DHCP message for which the lease update was rejected is not a DHCPv6 message");
996 }
997 auto duid = getClientId(msg, D6O_CLIENTID);
998 if (duid.empty()) {
999 return (false);
1000 }
1001 RejectedClient6 client{ duid, time(NULL) + lifetime };
1002 auto existing_client = rejected_clients_.find(duid);
1003 if (existing_client == rejected_clients_.end()) {
1004 rejected_clients_.insert(client);
1005 return (true);
1006 }
1007 rejected_clients_.replace(existing_client, client);
1008 return (false);
1009}
1010
1011bool
1013 // Early check if there is anything to do.
1015 return (false);
1016 }
1017 Pkt6Ptr msg = boost::dynamic_pointer_cast<Pkt6>(message);
1018 if (!msg) {
1019 isc_throw(BadValue, "DHCP message for which the lease update was successful is not a DHCPv6 message");
1020 }
1021 auto duid = getClientId(msg, D6O_CLIENTID);
1022 if (duid.empty()) {
1023 return (false);
1024 }
1025 auto existing_client = rejected_clients_.find(duid);
1026 if (existing_client != rejected_clients_.end()) {
1027 rejected_clients_.erase(existing_client);
1028 return (true);
1029 }
1030 return (false);
1031}
1032
1033void
1037
1038} // end of namespace isc::ha
1039} // end of namespace isc
static ElementPtr create(const Position &pos=ZERO_POSITION())
Create a NullElement.
Definition data.cc:300
@ list
Definition data.h:159
@ string
Definition data.h:157
static ElementPtr createMap(const Position &pos=ZERO_POSITION())
Creates an empty MapElement type ElementPtr.
Definition data.cc:355
static ElementPtr createList(const Position &pos=ZERO_POSITION())
Creates an empty ListElement type ElementPtr.
Definition data.cc:350
A generic exception that is thrown if a parameter given to a method is considered invalid in that con...
virtual bool reportRejectedLeaseUpdateInternal(const dhcp::PktPtr &message, const uint32_t lifetime)
Marks that the lease update failed due to a conflict for the specified DHCP message.
virtual size_t getRejectedLeaseUpdatesCountInternal()
Returns the number of lease updates rejected by the partner.
virtual bool reportSuccessfulLeaseUpdateInternal(const dhcp::PktPtr &message)
Marks the lease update successful.
virtual size_t getUnackedClientsCount() const
Returns the current number of clients which haven't gotten a lease from the partner server.
virtual void clearRejectedLeaseUpdatesInternal()
Clears rejected client leases.
virtual void analyzeMessageInternal(const boost::shared_ptr< dhcp::Pkt > &message)
Checks if the DHCPv4 message appears to be unanswered.
virtual size_t getConnectingClientsCount() const
Returns the current number of clients which attempted to get a lease from the partner server.
virtual void analyzeMessage(const boost::shared_ptr< dhcp::Pkt > &message)
Checks if the DHCPv4 message appears to be unanswered.
RejectedClients4 rejected_clients_
Holds information about the clients for whom lease updates have been rejected by the partner.
virtual bool failureDetectedInternal() const
Checks if the partner failure has been detected based on the DHCP traffic analysis.
ConnectingClients4 connecting_clients_
Holds information about the clients attempting to contact the partner server while the servers are in...
virtual bool failureDetected() const
Checks if the partner failure has been detected based on the DHCP traffic analysis.
virtual void clearConnectingClients()
Removes information about the clients the partner server should respond to while communication with t...
CommunicationState4(const asiolink::IOServicePtr &io_service, const HAConfigPtr &config)
Constructor.
virtual void analyzeMessage(const boost::shared_ptr< dhcp::Pkt > &message)
Checks if the DHCPv6 message appears to be unanswered.
RejectedClients6 rejected_clients_
Holds information about the clients for whom lease updates have been rejected by the partner.
virtual size_t getRejectedLeaseUpdatesCountInternal()
Returns the number of lease updates rejected by the partner.
ConnectingClients6 connecting_clients_
Holds information about the clients attempting to contact the partner server while the servers are in...
CommunicationState6(const asiolink::IOServicePtr &io_service, const HAConfigPtr &config)
Constructor.
virtual bool reportSuccessfulLeaseUpdateInternal(const dhcp::PktPtr &message)
Marks the lease update successful.
virtual void clearConnectingClients()
Removes information about the clients the partner server should respond to while communication with t...
virtual bool failureDetected() const
Checks if the partner failure has been detected based on the DHCP traffic analysis.
virtual size_t getUnackedClientsCount() const
Returns the current number of clients which haven't gotten a lease from the partner server.
virtual bool failureDetectedInternal() const
Checks if the partner failure has been detected based on the DHCP traffic analysis.
virtual void analyzeMessageInternal(const boost::shared_ptr< dhcp::Pkt > &message)
Checks if the DHCPv6 message appears to be unanswered.
virtual bool reportRejectedLeaseUpdateInternal(const dhcp::PktPtr &message, const uint32_t lifetime=86400)
Marks that the lease update failed due to a conflict for the specified DHCP message.
virtual size_t getConnectingClientsCount() const
Returns the current number of clients which attempted to get a lease from the partner server.
virtual void clearRejectedLeaseUpdatesInternal()
Clears rejected client leases.
virtual size_t getConnectingClientsCount() const =0
Returns the current number of clients which attempted to get a lease from the partner server.
virtual bool reportRejectedLeaseUpdateInternal(const dhcp::PktPtr &message, const uint32_t lifetime)=0
Marks that the lease update failed due to a conflict for the specified DHCP message.
boost::posix_time::ptime partner_state_time_
Holds a time when partner was first seen in the current state.
virtual void clearRejectedLeaseUpdatesInternal()=0
Clears rejected client leases.
virtual size_t getUnackedClientsCount() const =0
Returns the current number of clients which haven't got the lease from the partner server.
virtual void clearConnectingClients()=0
Removes information about the clients the partner server should respond to while communication with t...
void clearRejectedLeaseUpdates()
Clears rejected client leases (MT safe).
void startHeartbeat(const long interval, const std::function< void()> &heartbeat_impl)
Starts recurring heartbeat (public interface).
uint64_t unsent_update_count_
Total number of unsent lease updates.
bool isCommunicationInterrupted() const
Checks if communication with the partner is interrupted.
void setPartnerScopes(data::ConstElementPtr new_scopes)
Sets partner scopes.
int getPartnerState() const
Returns last known state of the partner.
bool clockSkewShouldWarn()
Issues a warning about high clock skew between the active servers if one is warranted.
std::string logFormatClockSkew() const
Returns current clock skew value in the logger friendly format.
void setPartnerUnsentUpdateCount(uint64_t unsent_update_count)
Saves new total number of unsent lease updates from the partner.
void setPartnerState(const std::string &state)
Sets partner state.
bool clockSkewShouldTerminate()
Indicates whether the HA service should enter "terminated" state as a result of the clock skew exceed...
std::pair< uint64_t, uint64_t > partner_unsent_update_count_
Previous and current total number of unsent lease updates from the partner.
std::set< std::string > getPartnerScopes() const
Returns scopes served by the partner server.
virtual ~CommunicationState()
Destructor.
HAConfigPtr config_
High availability configuration.
bool isHeartbeatRunning() const
Checks if recurring heartbeat is running.
static size_t getRejectedLeaseUpdatesCountFromContainer(RejectedClientsType &rejected_clients)
Extracts the number of lease updates rejected by the partner from the specified container.
long interval_
Interval specified for the heartbeat.
void setPartnerUnavailable()
Sets partner state unavailable.
void stopHeartbeat()
Stops recurring heartbeat.
void increaseUnsentUpdateCount()
Increases a total number of unsent lease updates by 1.
void setPartnerTime(const std::string &time_text)
Provide partner's notion of time so the new clock skew can be calculated.
bool hasPartnerNewUnsentUpdates() const
Checks if the partner allocated new leases for which it hasn't sent any lease updates.
virtual bool reportSuccessfulLeaseUpdateInternal(const dhcp::PktPtr &message)=0
Marks the lease update successful.
asiolink::IOServicePtr io_service_
Pointer to the common IO service instance.
virtual size_t getRejectedLeaseUpdatesCountInternal()=0
Returns the number of lease updates rejected by the partner.
void modifyPokeTime(const long secs)
Modifies poke time by adding seconds to it.
const boost::scoped_ptr< std::mutex > mutex_
The mutex used to protect internal state.
data::ElementPtr getReport() const
Returns the report about current communication state.
boost::posix_time::ptime getPartnerTimeAtSkew() const
Retrieves the time of the partner node when skew was last calculated.
boost::posix_time::time_duration clock_skew_
Clock skew between the active servers.
size_t getAnalyzedMessagesCount() const
Returns the number of analyzed messages while being in the communications interrupted state.
size_t analyzed_messages_count_
Total number of analyzed messages to be responded by partner.
std::function< void()> heartbeat_impl_
Pointer to the function providing heartbeat implementation.
boost::posix_time::ptime poke_time_
Last poke time.
boost::posix_time::time_duration updatePokeTime()
Update the poke time and compute the duration.
bool reportSuccessfulLeaseUpdate(const dhcp::PktPtr &message)
Marks the lease update successful (MT safe).
boost::posix_time::ptime partner_time_at_skew_
Partner reported time when skew was calculated.
CommunicationState(const asiolink::IOServicePtr &io_service, const HAConfigPtr &config)
Constructor.
boost::posix_time::time_duration getDurationSincePartnerStateTime() const
Returns the duration since the partner was first seen in the current state.
int partner_state_
Last known state of the partner server.
boost::posix_time::ptime last_clock_skew_warn_
Holds a time when last warning about too high clock skew was issued.
std::set< std::string > partner_scopes_
Last known set of scopes served by the partner server.
static std::vector< uint8_t > getClientId(const dhcp::PktPtr &message, const uint16_t option_type)
Convenience function attempting to retrieve client identifier from the DHCP message.
uint64_t getUnsentUpdateCount() const
Returns a total number of unsent lease updates.
bool rejectedLeaseUpdatesShouldTerminate()
Indicates whether the HA service should enter "terminated" state due to excessive number of rejected ...
boost::posix_time::ptime getMyTimeAtSkew() const
Retrieves the time of the local node when skew was last calculated.
boost::posix_time::ptime my_time_at_skew_
My time when skew was calculated.
int64_t getDurationInMillisecs() const
Returns duration between the poke time and current time.
bool reportRejectedLeaseUpdate(const dhcp::PktPtr &message, const uint32_t lifetime=86400)
Marks that the lease update failed due to a conflict for the specified DHCP message (MT safe).
size_t getRejectedLeaseUpdatesCount()
Returns the number of lease updates rejected by the partner (MT safe).
asiolink::IntervalTimerPtr timer_
Interval timer triggering heartbeat commands.
void poke()
Pokes the communication state.
This class parses and generates time values used in HTTP.
Definition date_time.h:41
boost::posix_time::ptime getPtime() const
Returns time encapsulated by this class.
Definition date_time.h:59
static HttpDateTime fromRfc1123(const std::string &time_string)
Creates an instance from a string containing time value formatted as specified in RFC 1123.
Definition date_time.cc:54
static MultiThreadingMgr & instance()
Returns a single instance of Multi Threading Manager.
@ D6O_CLIENTID
Definition dhcp6.h:21
@ D6O_ELAPSED_TIME
Definition dhcp6.h:28
#define isc_throw(type, stream)
A shortcut macro to insert known values into exception arguments.
OptionInt< uint16_t > OptionUint16
Definition option_int.h:32
boost::shared_ptr< OptionUint16 > OptionUint16Ptr
Definition option_int.h:33
#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_WARN(LOGGER, MESSAGE)
Macro to conveniently test warn output and log it.
Definition macros.h:26
boost::shared_ptr< const Element > ConstElementPtr
Definition data.h:30
boost::shared_ptr< Element > ElementPtr
Definition data.h:29
boost::shared_ptr< isc::dhcp::Pkt > PktPtr
A pointer to either Pkt4 or Pkt6 packet.
Definition pkt.h:1006
@ DHO_DHCP_CLIENT_IDENTIFIER
Definition dhcp4.h:130
boost::shared_ptr< Pkt4 > Pkt4Ptr
A pointer to Pkt4 object.
Definition pkt4.h:556
boost::shared_ptr< Pkt6 > Pkt6Ptr
A pointer to Pkt6 packet.
Definition pkt6.h:31
boost::shared_ptr< Option > OptionPtr
Definition option.h:38
const isc::log::MessageID HA_COMMUNICATION_INTERRUPTED_CLIENT4_UNACKED
Definition ha_messages.h:22
const isc::log::MessageID HA_COMMUNICATION_INTERRUPTED_CLIENT6
Definition ha_messages.h:23
isc::log::Logger ha_logger("ha-hooks")
Definition ha_log.h:17
const isc::log::MessageID HA_HIGH_CLOCK_SKEW_CAUSED_TERMINATION
Definition ha_messages.h:48
const isc::log::MessageID HA_LEASE_UPDATE_REJECTS_CAUSED_TERMINATION
Definition ha_messages.h:85
boost::shared_ptr< HAConfig > HAConfigPtr
Pointer to the High Availability configuration structure.
Definition ha_config.h:39
const isc::log::MessageID HA_COMMUNICATION_INTERRUPTED_CLIENT6_UNACKED
Definition ha_messages.h:24
std::string stateToString(int state)
Returns state name.
const isc::log::MessageID HA_COMMUNICATION_INTERRUPTED_CLIENT4
Definition ha_messages.h:21
int stringToState(const std::string &state_name)
Returns state for a given name.
const isc::log::MessageID HA_HIGH_CLOCK_SKEW
Definition ha_messages.h:47
std::string ptimeToText(boost::posix_time::ptime t, size_t fsecs_precision=MAX_FSECS_PRECISION)
Converts ptime structure to text.
Defines the logger used by the top-level component of kea-lfc.
Structure holding information about the client which has sent the packet being analyzed.
Structure holding information about the client who has a rejected lease update.
Structure holding information about a client which sent a packet being analyzed.
Structure holding information about the client who has a rejected lease update.