LCOV - code coverage report
Current view: top level - bgp - bgp_xmpp_channel.cc (source / functions) Hit Total Coverage
Test: OpenSDN C/C++ coverage (all TARGET_SET jobs) Lines: 0 1918 0.0 %
Date: 2026-09-28 02:13:17 Functions: 0 164 0.0 %
Legend: Lines: hit not hit

          Line data    Source code
       1             : /*
       2             :  * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved.
       3             :  */
       4             : 
       5             : #include "bgp/bgp_xmpp_channel.h"
       6             : 
       7             : #include <boost/assign/list_of.hpp>
       8             : #include <boost/foreach.hpp>
       9             : 
      10             : #include <limits>
      11             : #include <sstream>
      12             : #include <vector>
      13             : #include <atomic>
      14             : 
      15             : #include "base/regex.h"
      16             : #include "base/task_annotations.h"
      17             : #include "bgp/bgp_config.h"
      18             : #include "bgp/bgp_factory.h"
      19             : #include "bgp/bgp_log.h"
      20             : #include "bgp/bgp_membership.h"
      21             : #include "bgp/bgp_server.h"
      22             : #include "bgp/bgp_update_sender.h"
      23             : #include "bgp/bgp_xmpp_peer_close.h"
      24             : #include "bgp/inet/inet_table.h"
      25             : #include "bgp/inet6/inet6_table.h"
      26             : #include "bgp/extended-community/etree.h"
      27             : #include "bgp/extended-community/load_balance.h"
      28             : #include "bgp/extended-community/mac_mobility.h"
      29             : #include "bgp/extended-community/local_sequence_number.h"
      30             : #include "bgp/extended-community/router_mac.h"
      31             : #include "bgp/extended-community/tag.h"
      32             : #include "bgp/large-community/tag.h"
      33             : #include "bgp/ermvpn/ermvpn_table.h"
      34             : #include "bgp/evpn/evpn_route.h"
      35             : #include "bgp/evpn/evpn_table.h"
      36             : #include "bgp/mvpn/mvpn_table.h"
      37             : #include "bgp/peer_close_manager.h"
      38             : #include "bgp/peer_stats.h"
      39             : #include "bgp/security_group/security_group.h"
      40             : #include "bgp/tunnel_encap/tunnel_encap.h"
      41             : #include "bgp/bgp_xmpp_rtarget_manager.h"
      42             : #include "control-node/sandesh/control_node_types.h"
      43             : #include "net/community_type.h"
      44             : #include "schema/xmpp_multicast_types.h"
      45             : #include "schema/xmpp_enet_types.h"
      46             : #include "schema/xmpp_mvpn_types.h"
      47             : #include "xml/xml_pugi.h"
      48             : #include "xmpp/xmpp_connection.h"
      49             : #include "xmpp/xmpp_init.h"
      50             : #include "xmpp/xmpp_server.h"
      51             : #include "xmpp/sandesh/xmpp_peer_info_types.h"
      52             : 
      53             : using autogen::EnetItemType;
      54             : using autogen::EnetNextHopListType;
      55             : using autogen::EnetSecurityGroupListType;
      56             : using autogen::EnetTunnelEncapsulationListType;
      57             : 
      58             : using autogen::McastItemType;
      59             : using autogen::McastNextHopsType;
      60             : using autogen::McastTunnelEncapsulationListType;
      61             : 
      62             : using autogen::MvpnItemType;
      63             : using autogen::MvpnNextHopType;
      64             : using autogen::MvpnTunnelEncapsulationListType;
      65             : 
      66             : using autogen::ItemType;
      67             : using autogen::NextHopListType;
      68             : using autogen::SecurityGroupListType;
      69             : using autogen::CommunityTagListType;
      70             : using autogen::TunnelEncapsulationListType;
      71             : using autogen::TagListType;
      72             : 
      73             : using boost::assign::list_of;
      74             : using boost::smatch;
      75             : using boost::system::error_code;
      76             : using contrail::regex;
      77             : using contrail::regex_match;
      78             : using contrail::regex_search;
      79             : using pugi::xml_node;
      80             : using std::unique_ptr;
      81             : using std::make_pair;
      82             : using std::numeric_limits;
      83             : using std::ostringstream;
      84             : using std::pair;
      85             : using std::set;
      86             : using std::string;
      87             : using std::vector;
      88             : 
      89             : //
      90             : // Calculate med from local preference.
      91             : // Should move agent definitions to a common location and use those here
      92             : // instead of hard coded values.
      93             : //
      94           0 : static uint32_t GetMedFromLocalPref(uint32_t local_pref) {
      95           0 :     if (local_pref == 0)
      96           0 :         return 0;
      97           0 :     if (local_pref == 100)
      98           0 :         return 200;
      99           0 :     if (local_pref == 200)
     100           0 :         return 100;
     101           0 :     return numeric_limits<uint32_t>::max() - local_pref;
     102             : }
     103             : 
     104           0 : void BgpXmppChannel::ErrorStats::incr_inet6_rx_bad_xml_token_count() {
     105           0 :     ++inet6_rx_bad_xml_token_count;
     106           0 : }
     107             : 
     108           0 : void BgpXmppChannel::ErrorStats::incr_inet6_rx_bad_prefix_count() {
     109           0 :     ++inet6_rx_bad_prefix_count;
     110           0 : }
     111             : 
     112           0 : void BgpXmppChannel::ErrorStats::incr_inet6_rx_bad_nexthop_count() {
     113           0 :     ++inet6_rx_bad_nexthop_count;
     114           0 : }
     115             : 
     116           0 : void BgpXmppChannel::ErrorStats::incr_inet6_rx_bad_afi_safi_count() {
     117           0 :     ++inet6_rx_bad_afi_safi_count;
     118           0 : }
     119             : 
     120           0 : uint64_t BgpXmppChannel::ErrorStats::get_inet6_rx_bad_xml_token_count() const {
     121           0 :     return inet6_rx_bad_xml_token_count;
     122             : }
     123             : 
     124           0 : uint64_t BgpXmppChannel::ErrorStats::get_inet6_rx_bad_prefix_count() const {
     125           0 :     return inet6_rx_bad_prefix_count;
     126             : }
     127             : 
     128           0 : uint64_t BgpXmppChannel::ErrorStats::get_inet6_rx_bad_nexthop_count() const {
     129           0 :     return inet6_rx_bad_nexthop_count;
     130             : }
     131             : 
     132           0 : uint64_t BgpXmppChannel::ErrorStats::get_inet6_rx_bad_afi_safi_count() const {
     133           0 :     return inet6_rx_bad_afi_safi_count;
     134             : }
     135             : 
     136             : class BgpXmppChannel::PeerStats : public IPeerDebugStats {
     137             : public:
     138           0 :     explicit PeerStats(BgpXmppChannel *peer)
     139           0 :         : parent_(peer) {
     140           0 :     }
     141             : 
     142             :     // Used when peer flaps.
     143             :     // Don't need to do anything since the BgpXmppChannel itself gets deleted.
     144           0 :     virtual void Clear() {
     145           0 :     }
     146             : 
     147             :     // Printable name
     148           0 :     virtual string ToString() const {
     149           0 :         return parent_->ToString();
     150             :     }
     151             : 
     152             :     // Previous State of the peer
     153           0 :     virtual string last_state() const {
     154           0 :         return (parent_->channel_->LastStateName());
     155             :     }
     156             :     // Last state change occurred at
     157           0 :     virtual string last_state_change_at() const {
     158           0 :         return (parent_->channel_->LastStateChangeAt());
     159             :     }
     160             : 
     161             :     // Last error on this peer
     162           0 :     virtual string last_error() const {
     163           0 :         return "";
     164             :     }
     165             : 
     166             :     // Last Event on this peer
     167           0 :     virtual string last_event() const {
     168           0 :         return (parent_->channel_->LastEvent());
     169             :     }
     170             : 
     171             :     // When was the Last
     172           0 :     virtual string last_flap() const {
     173           0 :         return (parent_->channel_->LastFlap());
     174             :     }
     175             : 
     176             :     // Total number of flaps
     177           0 :     virtual uint64_t num_flaps() const {
     178           0 :         return (parent_->channel_->FlapCount());
     179             :     }
     180             : 
     181           0 :     virtual void GetRxProtoStats(ProtoStats *stats) const {
     182           0 :         stats->open = parent_->channel_->rx_open();
     183           0 :         stats->close = parent_->channel_->rx_close();
     184           0 :         stats->keepalive = parent_->channel_->rx_keepalive();
     185           0 :         stats->update = parent_->channel_->rx_update();
     186           0 :     }
     187             : 
     188           0 :     virtual void GetTxProtoStats(ProtoStats *stats) const {
     189           0 :         stats->open = parent_->channel_->tx_open();
     190           0 :         stats->close = parent_->channel_->tx_close();
     191           0 :         stats->keepalive = parent_->channel_->tx_keepalive();
     192           0 :         stats->update = parent_->channel_->tx_update();
     193           0 :     }
     194             : 
     195           0 :     virtual void GetRxRouteUpdateStats(UpdateStats *stats)  const {
     196           0 :         stats->reach = parent_->stats_[RX].reach.load();
     197           0 :         stats->unreach = parent_->stats_[RX].unreach.load();
     198           0 :         stats->end_of_rib = parent_->stats_[RX].end_of_rib.load();
     199           0 :         stats->total = stats->reach + stats->unreach + stats->end_of_rib;
     200           0 :     }
     201             : 
     202           0 :     virtual void GetTxRouteUpdateStats(UpdateStats *stats)  const {
     203           0 :         stats->reach = parent_->stats_[TX].reach.load();
     204           0 :         stats->unreach = parent_->stats_[TX].unreach.load();
     205           0 :         stats->end_of_rib = parent_->stats_[TX].end_of_rib.load();
     206           0 :         stats->total = stats->reach + stats->unreach + stats->end_of_rib;
     207           0 :     }
     208             : 
     209           0 :     virtual void GetRxSocketStats(IPeerDebugStats::SocketStats *stats) const {
     210           0 :         const XmppSession *session = parent_->GetSession();
     211           0 :         if (session) {
     212           0 :             const io::SocketStats &socket_stats = session->GetSocketStats();
     213           0 :             stats->calls = socket_stats.read_calls;
     214           0 :             stats->bytes = socket_stats.read_bytes;
     215             :         }
     216           0 :     }
     217             : 
     218           0 :     virtual void GetTxSocketStats(IPeerDebugStats::SocketStats *stats) const {
     219           0 :         const XmppSession *session = parent_->GetSession();
     220           0 :         if (session) {
     221           0 :             const io::SocketStats &socket_stats = session->GetSocketStats();
     222           0 :             stats->calls = socket_stats.write_calls;
     223           0 :             stats->bytes = socket_stats.write_bytes;
     224           0 :             stats->blocked_count = socket_stats.write_blocked;
     225           0 :             stats->blocked_duration_usecs =
     226           0 :                 socket_stats.write_blocked_duration_usecs;
     227             :         }
     228           0 :     }
     229             : 
     230           0 :     virtual void GetRxErrorStats(RxErrorStats *stats) const {
     231           0 :         const BgpXmppChannel::ErrorStats &err_stats = parent_->error_stats();
     232           0 :         stats->inet6_bad_xml_token_count =
     233           0 :             err_stats.get_inet6_rx_bad_xml_token_count();
     234           0 :         stats->inet6_bad_prefix_count =
     235           0 :             err_stats.get_inet6_rx_bad_prefix_count();
     236           0 :         stats->inet6_bad_nexthop_count =
     237           0 :             err_stats.get_inet6_rx_bad_nexthop_count();
     238           0 :         stats->inet6_bad_afi_safi_count =
     239           0 :             err_stats.get_inet6_rx_bad_afi_safi_count();
     240           0 :     }
     241             : 
     242           0 :     virtual void GetRxRouteStats(RxRouteStats *stats) const {
     243           0 :         stats->total_path_count = parent_->Peer()->GetTotalPathCount();
     244           0 :         stats->primary_path_count = parent_->Peer()->GetPrimaryPathCount();
     245           0 :     }
     246             : 
     247           0 :     virtual void UpdateTxUnreachRoute(uint64_t count) {
     248           0 :         parent_->stats_[TX].unreach += count;
     249           0 :     }
     250             : 
     251           0 :     virtual void UpdateTxReachRoute(uint64_t count) {
     252           0 :         parent_->stats_[TX].reach += count;
     253           0 :     }
     254             : 
     255             : private:
     256             :     BgpXmppChannel *parent_;
     257             : };
     258             : 
     259             : class BgpXmppChannel::XmppPeer : public IPeer {
     260             : public:
     261           0 :     XmppPeer(BgpServer *server, BgpXmppChannel *channel)
     262           0 :         : server_(server),
     263           0 :           parent_(channel),
     264           0 :           is_closed_(false),
     265           0 :           send_ready_(true),
     266           0 :           closed_at_(0) {
     267           0 :         total_path_count_ = 0;
     268           0 :         primary_path_count_ = 0;
     269           0 :     }
     270             : 
     271           0 :     virtual ~XmppPeer() {
     272           0 :         assert(GetTotalPathCount() == 0);
     273             : 
     274           0 :         XmppPeerInfoData peer_info;
     275           0 :         peer_info.set_name(ToUVEKey());
     276           0 :         peer_info.set_deleted(true);
     277           0 :         parent_->XMPPPeerInfoSend(peer_info);
     278             : 
     279           0 :         PeerStatsData peer_stats_data;
     280           0 :         peer_stats_data.set_name(ToUVEKey());
     281           0 :         peer_stats_data.set_deleted(true);
     282           0 :         assert(!peer_stats_data.get_name().empty());
     283           0 :         BGP_UVE_SEND2(PeerStatsUve, peer_stats_data, "ObjectXmppPeerInfo");
     284           0 :     }
     285             : 
     286           0 :     virtual bool MembershipPathCallback(DBTablePartBase *tpart, BgpRoute *rt,
     287             :                                         BgpPath *path) {
     288           0 :         if (parent_->close_manager_->IsMembershipInUse())
     289           0 :             return parent_->close_manager_->MembershipPathCallback(tpart, rt,
     290           0 :                                                                    path);
     291             : 
     292           0 :         BgpTable *table = static_cast<BgpTable *>(tpart->parent());
     293           0 :         return table->DeletePath(tpart, rt, path);
     294             :     }
     295             : 
     296             :     virtual bool SendUpdate(const uint8_t *msg, size_t msgsize,
     297             :                             const std::string *msg_str);
     298           0 :     virtual bool SendUpdate(const uint8_t *msg, size_t msgsize) {
     299           0 :         return SendUpdate(msg, msgsize, NULL);
     300             :     }
     301           0 :     virtual const string &ToString() const {
     302           0 :         return parent_->ToString();
     303             :     }
     304             : 
     305           0 :     virtual bool CanUseMembershipManager() const {
     306           0 :         return parent_->GetMembershipRequestQueueSize() == 0;
     307             :     }
     308             : 
     309           0 :     virtual bool IsRegistrationRequired() const { return true; }
     310             : 
     311           0 :     virtual const string &ToUVEKey() const {
     312           0 :         return parent_->ToUVEKey();
     313             :     }
     314             : 
     315           0 :     virtual BgpServer *server() { return server_; }
     316           0 :     virtual BgpServer *server() const { return server_; }
     317           0 :     virtual IPeerClose *peer_close() {
     318           0 :         return parent_->peer_close_.get();
     319             :     }
     320           0 :     virtual IPeerClose *peer_close() const {
     321           0 :         return parent_->peer_close_.get();
     322             :     }
     323             : 
     324           0 :     void UpdateCloseRouteStats(Address::Family family, const BgpPath *old_path,
     325             :                                uint32_t path_flags) const {
     326           0 :         peer_close()->UpdateRouteStats(family, old_path, path_flags);
     327           0 :     }
     328             : 
     329           0 :     virtual IPeerDebugStats *peer_stats() {
     330           0 :         return parent_->peer_stats_.get();
     331             :     }
     332           0 :     virtual const IPeerDebugStats *peer_stats() const {
     333           0 :         return parent_->peer_stats_.get();
     334             :     }
     335             : 
     336           0 :     virtual bool IsReady() const {
     337           0 :         return (parent_->channel_->GetPeerState() == xmps::READY);
     338             :     }
     339           0 :     virtual const string GetStateName() const {
     340           0 :         switch (parent_->channel_->GetPeerState()) {
     341           0 :             case xmps::UNKNOWN: return "UNKNOWN";
     342           0 :             case xmps::READY: return "READY";
     343           0 :             case xmps::NOT_READY: return "NOT_READY";
     344           0 :             case xmps::TIMEDOUT: return "TIMEDOUT";
     345             :         }
     346           0 :         return "UNKNOWN";
     347             :     }
     348           0 :     virtual bool IsXmppPeer() const {
     349           0 :         return true;
     350             :     }
     351             :     virtual void Close(bool graceful);
     352             : 
     353           0 :     const bool IsDeleted() const { return is_closed_; }
     354           0 :     void SetPeerClosed(bool closed) {
     355           0 :         is_closed_ = closed;
     356           0 :         if (is_closed_)
     357           0 :             closed_at_ = UTCTimestampUsec();
     358           0 :     }
     359           0 :     uint64_t closed_at() const { return closed_at_; }
     360             : 
     361           0 :     virtual BgpProto::BgpPeerType PeerType() const {
     362           0 :         return BgpProto::XMPP;
     363             :     }
     364             : 
     365           0 :     virtual uint32_t bgp_identifier() const {
     366           0 :         const TcpSession::Endpoint &remote = parent_->endpoint();
     367           0 :         if (remote.address().is_v4()) {
     368           0 :             return remote.address().to_v4().to_ulong();
     369             :         }
     370           0 :         return 0;
     371             :     }
     372             : 
     373           0 :     virtual void UpdateTotalPathCount(int count) const {
     374           0 :         total_path_count_ += count;
     375           0 :     }
     376           0 :     virtual int GetTotalPathCount() const {
     377           0 :         return total_path_count_;
     378             :     }
     379           0 :     virtual bool IsAs4Supported() const { return true; }
     380           0 :     virtual void UpdatePrimaryPathCount(int count,
     381             :         Address::Family family) const {
     382           0 :         primary_path_count_ += count;
     383           0 :     }
     384           0 :     virtual int GetPrimaryPathCount() const {
     385           0 :          return primary_path_count_;
     386             :     }
     387           0 :     virtual void ProcessPathTunnelEncapsulation(const BgpPath *path,
     388             :         BgpAttr *attr, ExtCommunityDB *extcomm_db, const BgpTable *table)
     389             :         const {
     390           0 :     }
     391           0 :     virtual const std::vector<std::string> GetDefaultTunnelEncap(
     392             :         Address::Family family) const {
     393           0 :         return std::vector<std::string>();
     394             :     }
     395           0 :     virtual bool IsInGRTimerWaitState() const {
     396           0 :         return parent_->close_manager_->IsInGRTimerWaitState();
     397             :     }
     398             : 
     399           0 :     void MembershipRequestCallback(BgpTable *table) {
     400           0 :         parent_->MembershipRequestCallback(table);
     401           0 :     }
     402             : 
     403           0 :     virtual bool send_ready() const { return send_ready_; }
     404           0 :     bool IsRouterTypeBGPaaS() const { return false; }
     405             : 
     406             : private:
     407           0 :     void WriteReadyCb(const boost::system::error_code &ec) {
     408           0 :         if (!server_) return;
     409           0 :         BgpUpdateSender *sender = server_->update_sender();
     410           0 :         BGP_LOG_PEER(Event, this, SandeshLevel::SYS_DEBUG, BGP_LOG_FLAG_ALL,
     411             :                      BGP_PEER_DIR_NA, "Send ready");
     412           0 :         sender->PeerSendReady(this);
     413           0 :         send_ready_ = true;
     414             : 
     415             :         // Restart EndOfRib Send timer if necessary.
     416           0 :         parent_->ResetEndOfRibSendState();
     417             :     }
     418             : 
     419             :     BgpServer *server_;
     420             :     BgpXmppChannel *parent_;
     421             :     mutable std::atomic<int> total_path_count_;
     422             :     mutable std::atomic<int> primary_path_count_;
     423             :     bool is_closed_;
     424             :     bool send_ready_;
     425             :     uint64_t closed_at_;
     426             : };
     427             : 
     428             : // Skip sending updates if the destinatin matches against the pattern.
     429             : // XXX Used in test environments only
     430           0 : bool BgpXmppChannel::SkipUpdateSend() {
     431           0 :     static char *skip_env_ = getenv("XMPP_SKIP_UPDATE_SEND");
     432           0 :     if (!skip_env_)
     433           0 :         return false;
     434             : 
     435             :     // Use XMPP_SKIP_UPDATE_SEND as a regex pattern to match against destination
     436             :     // Cache the result to avoid redundant regex evaluation
     437           0 :     if (!skip_update_send_cached_) {
     438           0 :         smatch matches;
     439           0 :         skip_update_send_ = regex_search(ToString(), matches, regex(skip_env_));
     440           0 :         skip_update_send_cached_ = true;
     441           0 :     }
     442           0 :     return skip_update_send_;
     443             : }
     444             : 
     445           0 : bool BgpXmppChannel::XmppPeer::SendUpdate(const uint8_t *msg, size_t msgsize,
     446             :     const string *msg_str) {
     447           0 :     XmppChannel *channel = parent_->channel_;
     448           0 :     if (channel->GetPeerState() == xmps::READY) {
     449           0 :         parent_->stats_[TX].rt_updates++;
     450           0 :         if (parent_->SkipUpdateSend())
     451           0 :             return true;
     452           0 :         send_ready_ = channel->Send(msg, msgsize, msg_str, xmps::BGP,
     453             :                 boost::bind(&BgpXmppChannel::XmppPeer::WriteReadyCb, this, _1));
     454           0 :         if (!send_ready_) {
     455           0 :             BGP_LOG_PEER(Event, this, SandeshLevel::SYS_DEBUG, BGP_LOG_FLAG_ALL,
     456             :                          BGP_PEER_DIR_NA, "Send blocked");
     457             : 
     458             :             // If EndOfRib Send timer is running, cancel it and reschedule it
     459             :             // after socket gets unblocked.
     460           0 :             if (parent_->eor_send_timer_ && parent_->eor_send_timer_->running())
     461           0 :                 parent_->eor_send_timer_->Cancel();
     462             :         }
     463           0 :         return send_ready_;
     464             :     } else {
     465           0 :         return false;
     466             :     }
     467             : }
     468             : 
     469           0 : void BgpXmppChannel::XmppPeer::Close(bool graceful) {
     470           0 :     send_ready_ = true;
     471           0 :     parent_->set_peer_closed(true);
     472           0 :     if (server_ == NULL)
     473           0 :         return;
     474             : 
     475             :     XmppConnection *connection =
     476           0 :         const_cast<XmppConnection *>(parent_->channel_->connection());
     477             : 
     478           0 :     if (connection && !connection->IsActiveChannel()) {
     479             : 
     480             :         // Clear EOR state.
     481           0 :         parent_->ClearEndOfRibState();
     482             : 
     483           0 :         parent_->peer_close_->Close(graceful);
     484             :     }
     485             : }
     486             : 
     487           0 : BgpXmppChannel::BgpXmppChannel(XmppChannel *channel,
     488             :                                BgpServer *bgp_server,
     489           0 :                                BgpXmppChannelManager *manager)
     490           0 :     : channel_(channel),
     491           0 :       peer_id_(xmps::BGP),
     492           0 :       rtarget_manager_(new BgpXmppRTargetManager(this)),
     493           0 :       bgp_server_(bgp_server),
     494           0 :       peer_(new XmppPeer(bgp_server, this)),
     495           0 :       peer_close_(new BgpXmppPeerClose(this)),
     496           0 :       peer_stats_(new PeerStats(this)),
     497           0 :       bgp_policy_(BgpProto::XMPP, RibExportPolicy::XMPP, -1, 0),
     498           0 :       manager_(manager),
     499           0 :       delete_in_progress_(false),
     500           0 :       deleted_(false),
     501           0 :       defer_peer_close_(false),
     502           0 :       skip_update_send_(false),
     503           0 :       skip_update_send_cached_(false),
     504           0 :       eor_sent_(false),
     505           0 :       eor_receive_timer_(NULL),
     506           0 :       eor_send_timer_(NULL),
     507           0 :       eor_receive_timer_start_time_(0),
     508           0 :       eor_send_timer_start_time_(0),
     509           0 :       membership_response_worker_(
     510             :             TaskScheduler::GetInstance()->GetTaskId("xmpp::StateMachine"),
     511           0 :             channel->GetTaskInstance(),
     512             :             boost::bind(&BgpXmppChannel::MembershipResponseHandler, this, _1)),
     513           0 :       lb_mgr_(new LabelBlockManager()) {
     514           0 :     close_manager_.reset(
     515           0 :         BgpStaticObjectFactory::Create<PeerCloseManager>(static_cast<IPeerClose*>(peer_close_.get())));
     516           0 :     if (bgp_server) {
     517           0 :         eor_receive_timer_ =
     518           0 :             TimerManager::CreateTimer(*bgp_server->ioservice(),
     519             :                 "EndOfRib receive timer",
     520             :                 TaskScheduler::GetInstance()->GetTaskId("xmpp::StateMachine"),
     521           0 :                 channel->GetTaskInstance());
     522           0 :         eor_send_timer_ =
     523           0 :             TimerManager::CreateTimer(*bgp_server->ioservice(),
     524             :                 "EndOfRib send timer",
     525             :                 TaskScheduler::GetInstance()->GetTaskId("xmpp::StateMachine"),
     526           0 :                 channel->GetTaskInstance());
     527             :     }
     528           0 :     channel_->RegisterReferer(peer_id_);
     529           0 :     channel_->RegisterReceive(peer_id_,
     530             :          boost::bind(&BgpXmppChannel::ReceiveUpdate, this, _1));
     531           0 :     BGP_LOG_PEER(Event, peer_.get(), SandeshLevel::SYS_INFO, BGP_LOG_FLAG_ALL,
     532             :         BGP_PEER_DIR_NA, "Created");
     533           0 : }
     534             : 
     535           0 : BgpXmppChannel::~BgpXmppChannel() {
     536           0 :     if (channel_->connection() && !channel_->connection()->IsActiveChannel()) {
     537           0 :         CHECK_CONCURRENCY("bgp::Config");
     538             :     }
     539             : 
     540           0 :     if (manager_)
     541           0 :         manager_->RemoveChannel(channel_);
     542           0 :     if (manager_ && delete_in_progress_)
     543           0 :         manager_->decrement_deleting_count();
     544           0 :     STLDeleteElements(&defer_q_);
     545           0 :     assert(peer_deleted());
     546           0 :     assert(!close_manager_->IsMembershipInUse());
     547           0 :     assert(table_membership_request_map_.empty());
     548           0 :     TimerManager::DeleteTimer(eor_receive_timer_);
     549           0 :     TimerManager::DeleteTimer(eor_send_timer_);
     550           0 :     BGP_LOG_PEER(Event, peer_.get(), SandeshLevel::SYS_INFO, BGP_LOG_FLAG_ALL,
     551             :         BGP_PEER_DIR_NA, "Deleted");
     552           0 :     channel_->UnRegisterWriteReady(peer_id_);
     553           0 :     channel_->UnRegisterReferer(peer_id_);
     554           0 :     channel_->UnRegisterReceive(peer_id_);
     555           0 : }
     556             : 
     557           0 : void BgpXmppChannel::XMPPPeerInfoSend(const XmppPeerInfoData &peer_info) const {
     558           0 :     assert(!peer_info.get_name().empty());
     559           0 :     BGP_UVE_SEND(XMPPPeerInfo, peer_info);
     560           0 : }
     561             : 
     562           0 : const XmppSession *BgpXmppChannel::GetSession() const {
     563           0 :     if (channel_ && channel_->connection()) {
     564           0 :         return channel_->connection()->session();
     565             :     }
     566           0 :     return NULL;
     567             : }
     568             : 
     569           0 : const string &BgpXmppChannel::ToString() const {
     570           0 :     return channel_->ToString();
     571             : }
     572             : 
     573           0 : const string &BgpXmppChannel::ToUVEKey() const {
     574           0 :     if (channel_->connection()) {
     575           0 :         return channel_->connection()->ToUVEKey();
     576             :     } else {
     577           0 :         return channel_->ToString();
     578             :     }
     579             : }
     580             : 
     581           0 : string BgpXmppChannel::StateName() const {
     582           0 :     return channel_->StateName();
     583             : }
     584             : 
     585             : 
     586           0 : size_t BgpXmppChannel::GetMembershipRequestQueueSize() const {
     587           0 :     return table_membership_request_map_.size();
     588             : }
     589             : 
     590           0 : void BgpXmppChannel::RoutingInstanceCallback(string vrf_name, int op) {
     591           0 :     if (delete_in_progress_)
     592           0 :         return;
     593           0 :     if (vrf_name == BgpConfigManager::kMasterInstance)
     594           0 :         return;
     595           0 :     if (op == RoutingInstanceMgr::INSTANCE_DELETE)
     596           0 :         return;
     597             : 
     598           0 :     RoutingInstanceMgr *instance_mgr = bgp_server_->routing_instance_mgr();
     599           0 :     assert(instance_mgr);
     600           0 :     RoutingInstance *rt_instance = instance_mgr->GetRoutingInstance(vrf_name);
     601           0 :     assert(rt_instance);
     602             : 
     603           0 :     if (op == RoutingInstanceMgr::INSTANCE_ADD) {
     604             :         const InstanceMembershipRequestState *imr_state =
     605           0 :             GetInstanceMembershipState(vrf_name);
     606           0 :         if (!imr_state)
     607           0 :             return;
     608           0 :         ProcessDeferredSubscribeRequest(rt_instance, *imr_state);
     609           0 :         DeleteInstanceMembershipState(vrf_name);
     610             :     } else {
     611           0 :         SubscriptionState *sub_state = GetSubscriptionState(rt_instance);
     612           0 :         if (!sub_state)
     613           0 :             return;
     614           0 :         rtarget_manager_->RoutingInstanceCallback(
     615             :             rt_instance, &sub_state->targets);
     616             :     }
     617             : }
     618             : 
     619           0 : IPeer *BgpXmppChannel::Peer() {
     620           0 :     return peer_.get();
     621             : }
     622             : 
     623           0 : const IPeer *BgpXmppChannel::Peer() const {
     624           0 :     return peer_.get();
     625             : }
     626             : 
     627           0 : TcpSession::Endpoint BgpXmppChannel::endpoint() const {
     628           0 :     return channel_->connection()->endpoint();
     629             : }
     630             : 
     631           0 : bool BgpXmppChannel::XmppDecodeAddress(int af, const string &address,
     632             :                                        IpAddress *addrp, bool zero_ok) {
     633           0 :     if (af != BgpAf::IPv4 && af != BgpAf::IPv6 && af != BgpAf::L2Vpn)
     634           0 :         return false;
     635             : 
     636           0 :     error_code error;
     637           0 :     *addrp = IpAddress::from_string(address, error);
     638           0 :     if (error)
     639           0 :         return false;
     640             : 
     641           0 :     return (zero_ok ? true : !addrp->is_unspecified());
     642             : }
     643             : 
     644             : //
     645             : // Return true if there's a pending request, false otherwise.
     646             : //
     647           0 : bool BgpXmppChannel::GetMembershipInfo(BgpTable *table,
     648             :     int *instance_id, uint64_t *subscription_gen_id, RequestType *req_type) {
     649           0 :     *instance_id = -1;
     650           0 :     *subscription_gen_id = 0;
     651             :     TableMembershipRequestState *tmr_state =
     652           0 :         GetTableMembershipState(table->name());
     653           0 :     if (tmr_state) {
     654           0 :         *req_type = tmr_state->pending_req;
     655           0 :         *instance_id = tmr_state->instance_id;
     656           0 :         return true;
     657             :     } else {
     658           0 :         *req_type = NONE;
     659           0 :         BgpMembershipManager *mgr = bgp_server_->membership_mgr();
     660           0 :         mgr->GetRegistrationInfo(peer_.get(), table,
     661             :                                  instance_id, subscription_gen_id);
     662           0 :         return false;
     663             :     }
     664             : }
     665             : 
     666             : //
     667             : // Add entry to the pending table request map.
     668             : //
     669           0 : void BgpXmppChannel::AddTableMembershipState(const string &table_name,
     670             :     TableMembershipRequestState tmr_state) {
     671           0 :     table_membership_request_map_.insert(make_pair(table_name, tmr_state));
     672           0 : }
     673             : 
     674             : //
     675             : // Delete entry from the pending table request map.
     676             : // Return true if the entry was found and deleted.
     677             : //
     678           0 : bool BgpXmppChannel::DeleteTableMembershipState(const string &table_name) {
     679           0 :     return (table_membership_request_map_.erase(table_name) > 0);
     680             : }
     681             : 
     682             : //
     683             : // Find entry in the pending table request map.
     684             : //
     685             : BgpXmppChannel::TableMembershipRequestState *
     686           0 : BgpXmppChannel::GetTableMembershipState(
     687             :     const string &table_name) {
     688             :     TableMembershipRequestMap::iterator loc =
     689           0 :         table_membership_request_map_.find(table_name);
     690           0 :     return (loc == table_membership_request_map_.end() ? NULL : &loc->second);
     691             : }
     692             : 
     693             : //
     694             : // Find entry in the pending table request map.
     695             : // Const version.
     696             : //
     697             : const BgpXmppChannel::TableMembershipRequestState *
     698           0 : BgpXmppChannel::GetTableMembershipState(
     699             :     const string &table_name) const {
     700             :     TableMembershipRequestMap::const_iterator loc =
     701           0 :         table_membership_request_map_.find(table_name);
     702           0 :     return (loc == table_membership_request_map_.end() ? NULL : &loc->second);
     703             : }
     704             : 
     705             : //
     706             : // Add entry to the pending instance request map.
     707             : //
     708           0 : void BgpXmppChannel::AddInstanceMembershipState(const string &instance,
     709             :     InstanceMembershipRequestState imr_state) {
     710           0 :     instance_membership_request_map_.insert(make_pair(instance, imr_state));
     711           0 : }
     712             : 
     713             : //
     714             : // Delete entry from the pending instance request map.
     715             : // Return true if the entry was found and deleted.
     716             : //
     717           0 : bool BgpXmppChannel::DeleteInstanceMembershipState(const string &instance) {
     718           0 :     return (instance_membership_request_map_.erase(instance) > 0);
     719             : }
     720             : 
     721             : //
     722             : // Find the entry in the pending instance request map.
     723             : //
     724             : const BgpXmppChannel::InstanceMembershipRequestState *
     725           0 : BgpXmppChannel::GetInstanceMembershipState(const string &instance) const {
     726             :     InstanceMembershipRequestMap::const_iterator loc =
     727           0 :         instance_membership_request_map_.find(instance);
     728           0 :     return loc != instance_membership_request_map_.end() ? &loc->second : NULL;
     729             : }
     730             : 
     731             : //
     732             : // Verify that there's a subscribe or pending subscribe for the table
     733             : // corresponding to the vrf and family.
     734             : // If there's a subscribe, populate the table and instance_id.
     735             : // If there's a pending subscribe, populate the instance_id.
     736             : // The subscribe_pending parameter is set appropriately.
     737             : //
     738             : // Return true if there's a subscribe or pending subscribe, false otherwise.
     739             : //
     740           0 : bool BgpXmppChannel::VerifyMembership(const string &vrf_name,
     741             :     Address::Family family, BgpTable **table,
     742             :     int *instance_id, uint64_t *subscription_gen_id, bool *subscribe_pending,
     743             :     bool add_change) {
     744           0 :     *table = NULL;
     745           0 :     *subscribe_pending = false;
     746             : 
     747           0 :     RoutingInstanceMgr *instance_mgr = bgp_server_->routing_instance_mgr();
     748           0 :     RoutingInstance *rt_instance = instance_mgr->GetRoutingInstance(vrf_name);
     749           0 :     if (rt_instance)
     750           0 :         *table = rt_instance->GetTable(family);
     751           0 :     if (rt_instance != NULL && !rt_instance->deleted()) {
     752             :         RequestType req_type;
     753           0 :         if (GetMembershipInfo(*table, instance_id,
     754             :                               subscription_gen_id, &req_type)) {
     755             :             // Bail if there's a pending unsubscribe.
     756           0 :             if (req_type != SUBSCRIBE) {
     757           0 :                 BGP_LOG_PEER_INSTANCE_CRITICAL(Peer(), vrf_name,
     758             :                     BGP_PEER_DIR_IN, BGP_LOG_FLAG_ALL,
     759             :                     "Received route after unsubscribe");
     760           0 :                 return false;
     761             :             }
     762           0 :             *subscribe_pending = true;
     763             :         } else {
     764             :             // Bail if we are not subscribed to the table.
     765           0 :             if (*instance_id < 0) {
     766           0 :                 BGP_LOG_PEER_INSTANCE_CRITICAL(Peer(), vrf_name,
     767             :                     BGP_PEER_DIR_IN, BGP_LOG_FLAG_ALL,
     768             :                     "Received route without subscribe");
     769           0 :                 return false;
     770             :             }
     771             :         }
     772             :     } else {
     773             :         // Bail if there's no pending subscribe for the instance.
     774             :         // Note that route retract can be received while the instance is
     775             :         // marked for deletion.
     776             :         const InstanceMembershipRequestState *imr_state =
     777           0 :             GetInstanceMembershipState(vrf_name);
     778           0 :         if (imr_state) {
     779           0 :             *instance_id = imr_state->instance_id;
     780           0 :             *subscribe_pending = true;
     781           0 :         } else if (add_change || !rt_instance) {
     782           0 :             BGP_LOG_PEER_INSTANCE_CRITICAL(Peer(), vrf_name, BGP_PEER_DIR_IN,
     783             :                BGP_LOG_FLAG_ALL, "Received route without pending subscribe");
     784           0 :             return false;
     785             :         }
     786             :     }
     787             : 
     788           0 :     return true;
     789             : }
     790             : 
     791           0 : bool BgpXmppChannel::ProcessMcastItem(string vrf_name,
     792             :     const pugi::xml_node &node, bool add_change) {
     793           0 :     McastItemType item;
     794           0 :     item.Clear();
     795             : 
     796           0 :     if (!item.XmlParse(node)) {
     797           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
     798             :             BGP_LOG_FLAG_ALL, "Invalid multicast route message received");
     799           0 :         return false;
     800             :     }
     801             : 
     802           0 :     if (item.entry.nlri.af != BgpAf::IPv4) {
     803           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
     804             :             "Unsupported address family " << item.entry.nlri.af <<
     805             :             " for multicast route");
     806           0 :         return false;
     807             :     }
     808             : 
     809           0 :     if (item.entry.nlri.safi != BgpAf::Mcast) {
     810           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
     811             :             BGP_LOG_FLAG_ALL, "Unsupported subsequent address family " <<
     812             :             item.entry.nlri.safi << " for multicast route");
     813           0 :         return false;
     814             :     }
     815             : 
     816           0 :     error_code error;
     817           0 :     IpAddress grp_address = IpAddress::from_string("0.0.0.0", error);
     818           0 :     if (!item.entry.nlri.group.empty()) {
     819           0 :         if (!XmppDecodeAddress(item.entry.nlri.af,
     820             :             item.entry.nlri.group, &grp_address, false)) {
     821           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
     822             :                 "Bad group address " << item.entry.nlri.group);
     823           0 :             return false;
     824             :         }
     825             :     }
     826             : 
     827           0 :     IpAddress src_address = IpAddress::from_string("0.0.0.0", error);
     828           0 :     if (!item.entry.nlri.source.empty()) {
     829           0 :         if (!XmppDecodeAddress(item.entry.nlri.af,
     830             :             item.entry.nlri.source, &src_address, true)) {
     831           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
     832             :                 "Bad source address " << item.entry.nlri.source);
     833           0 :             return false;
     834             :         }
     835             :     }
     836             : 
     837             :     bool subscribe_pending;
     838             :     int instance_id;
     839             :     uint64_t subscription_gen_id;
     840             :     BgpTable *table;
     841           0 :     if (!VerifyMembership(vrf_name, Address::ERMVPN, &table, &instance_id,
     842             :         &subscription_gen_id, &subscribe_pending, add_change)) {
     843           0 :         channel_->Close();
     844           0 :         return false;
     845             :     }
     846             : 
     847             :     // Build the key to the Multicast DBTable
     848             :     uint16_t cluster_seed =
     849           0 :         bgp_server_->global_config()->rd_cluster_seed();
     850           0 :     RouteDistinguisher mc_rd;
     851           0 :     if (cluster_seed) {
     852           0 :         mc_rd = RouteDistinguisher(cluster_seed, peer_->bgp_identifier(),
     853           0 :                                    instance_id);
     854             :     } else {
     855           0 :         mc_rd = RouteDistinguisher(peer_->bgp_identifier(), instance_id);
     856             :     }
     857             : 
     858             :     ErmVpnPrefix mc_prefix(ErmVpnPrefix::NativeRoute, mc_rd,
     859           0 :         grp_address.to_v4(), src_address.to_v4());
     860             : 
     861             :     // Build and enqueue a DB request for route-addition
     862           0 :     DBRequest req;
     863           0 :     req.key.reset(new ErmVpnTable::RequestKey(mc_prefix, peer_.get()));
     864             : 
     865           0 :     uint32_t flags = 0;
     866           0 :     ExtCommunitySpec ext;
     867           0 :     string label_range("none");
     868             : 
     869           0 :     if (add_change) {
     870           0 :         req.oper = DBRequest::DB_ENTRY_ADD_CHANGE;
     871           0 :         vector<uint32_t> labels;
     872           0 :         const McastNextHopsType &inh_list = item.entry.next_hops;
     873             : 
     874           0 :         if (inh_list.next_hop.empty()) {
     875           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
     876             :                 BGP_LOG_FLAG_ALL, "Missing next-hop for multicast route " <<
     877             :                 mc_prefix.ToString());
     878           0 :             return false;
     879             :         }
     880             : 
     881             :         // Agents should send only one next-hop in the item
     882           0 :         if (inh_list.next_hop.size() != 1) {
     883           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
     884             :                 "More than one nexthop received for multicast route " <<
     885             :                 mc_prefix.ToString());
     886           0 :             return false;
     887             :         }
     888             : 
     889           0 :         McastNextHopsType::const_iterator nit = inh_list.begin();
     890             : 
     891             :         // Label Allocation item.entry.label by parsing the range
     892           0 :         label_range = nit->label;
     893           0 :         if (!stringToIntegerList(label_range, "-", labels) ||
     894           0 :             labels.size() != 2) {
     895           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
     896             :                 "Bad label range " << label_range <<
     897             :                 " for multicast route " << mc_prefix.ToString());
     898           0 :             return false;
     899             :         }
     900             : 
     901           0 :         if (!labels[0] || !labels[1] || labels[1] < labels[0]) {
     902           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
     903             :                 BGP_LOG_FLAG_ALL, "Bad label range " << label_range <<
     904             :                 " for multicast route " << mc_prefix.ToString());
     905           0 :             return false;
     906             :         }
     907             : 
     908           0 :         BgpAttrSpec attrs;
     909           0 :         LabelBlockPtr lbptr = lb_mgr_->LocateBlock(labels[0], labels[1]);
     910             : 
     911           0 :         BgpAttrLabelBlock attr_label(lbptr);
     912           0 :         attrs.push_back(&attr_label);
     913             : 
     914             :         // Next-hop ip address
     915           0 :         IpAddress nh_address;
     916           0 :         if (!XmppDecodeAddress(nit->af, nit->address, &nh_address)) {
     917           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
     918             :                 BGP_LOG_FLAG_ALL, "Bad nexthop address " << nit->address <<
     919             :                 " for multicast route " << mc_prefix.ToString());
     920           0 :             return false;
     921             :         }
     922           0 :         BgpAttrNextHop nexthop(nh_address.to_v4().to_ulong());
     923           0 :         attrs.push_back(&nexthop);
     924             : 
     925             :         // Process tunnel encapsulation list.
     926           0 :         bool no_tunnel_encap = true;
     927           0 :         bool no_valid_tunnel_encap = true;
     928           0 :         for (McastTunnelEncapsulationListType::const_iterator eit =
     929           0 :              nit->tunnel_encapsulation_list.begin();
     930           0 :              eit != nit->tunnel_encapsulation_list.end(); ++eit) {
     931           0 :             no_tunnel_encap = false;
     932           0 :             TunnelEncap tun_encap(*eit);
     933           0 :             if (tun_encap.tunnel_encap() == TunnelEncapType::UNSPEC)
     934           0 :                 continue;
     935           0 :             no_valid_tunnel_encap = false;
     936           0 :             ext.communities.push_back(tun_encap.GetExtCommunityValue());
     937             :         }
     938             : 
     939             :         // Mark the path as infeasible if all tunnel encaps published
     940             :         // by agent are invalid.
     941           0 :         if (!no_tunnel_encap && no_valid_tunnel_encap) {
     942           0 :             flags = BgpPath::NoTunnelEncap;
     943             :         }
     944             : 
     945           0 :         if (!ext.communities.empty())
     946           0 :             attrs.push_back(&ext);
     947             : 
     948           0 :         BgpAttrPtr attr = bgp_server_->attr_db()->Locate(attrs);
     949           0 :         req.data.reset(new ErmVpnTable::RequestData(
     950           0 :             attr, flags, 0, 0, subscription_gen_id));
     951           0 :         stats_[RX].reach++;
     952           0 :     } else {
     953           0 :         req.oper = DBRequest::DB_ENTRY_DELETE;
     954           0 :         stats_[RX].unreach++;
     955             :     }
     956             : 
     957             :     // Defer all requests till subscribe is processed.
     958           0 :     if (subscribe_pending) {
     959           0 :         DBRequest *request_entry = new DBRequest();
     960           0 :         request_entry->Swap(&req);
     961             :         string table_name =
     962           0 :             RoutingInstance::GetTableName(vrf_name, Address::ERMVPN);
     963           0 :         defer_q_.insert(make_pair(
     964           0 :             make_pair(vrf_name, table_name), request_entry));
     965           0 :         return true;
     966           0 :     }
     967             : 
     968           0 :     assert(table);
     969           0 :     BGP_LOG_PEER_INSTANCE(Peer(), vrf_name,
     970             :         SandeshLevel::SYS_DEBUG, BGP_LOG_FLAG_TRACE,
     971             :         "Multicast group " << item.entry.nlri.group <<
     972             :         " source " << item.entry.nlri.source <<
     973             :         " and label range " << label_range <<
     974             :         " enqueued for " << (add_change ? "add/change" : "delete"));
     975           0 :     table->Enqueue(&req);
     976           0 :     return true;
     977           0 : }
     978             : 
     979           0 : void BgpXmppChannel::CreateType5MvpnRouteRequest(IpAddress grp_address,
     980             :         IpAddress src_address, bool add_change, uint64_t subscription_gen_id,
     981             :         int instance_id, DBRequest& req, const MvpnNextHopType &nexthop) {
     982           0 :     RouteDistinguisher mc_rd =  RouteDistinguisher::kZeroRd;
     983             :     MvpnPrefix mc_prefix(MvpnPrefix::SourceActiveADRoute, mc_rd,
     984           0 :             grp_address.to_v4(), src_address.to_v4());
     985           0 :     uint32_t flags = 0;
     986             : 
     987             :     // Build and enqueue a DB request for route-addition
     988           0 :     req.key.reset(new MvpnTable::RequestKey(mc_prefix, peer_.get()));
     989             : 
     990           0 :     if (add_change) {
     991           0 :         req.oper = DBRequest::DB_ENTRY_ADD_CHANGE;
     992             : 
     993           0 :         BgpAttrSpec attrs;
     994             :         // Next-hop ip address
     995           0 :         IpAddress nh_address;
     996           0 :         if (!XmppDecodeAddress(nexthop.af, nexthop.address, &nh_address)) {
     997           0 :             return;
     998             :         }
     999             : 
    1000             :         BgpAttrSourceRd source_rd(
    1001           0 :                 RouteDistinguisher(nh_address.to_v4().to_ulong(), instance_id));
    1002           0 :         attrs.push_back(&source_rd);
    1003           0 :         BgpAttrPtr attr = bgp_server_->attr_db()->Locate(attrs);
    1004           0 :         req.data.reset(new MvpnTable::RequestData(
    1005           0 :             attr, flags, 0, 0, subscription_gen_id));
    1006           0 :         stats_[RX].reach++;
    1007           0 :     } else {
    1008           0 :         req.oper = DBRequest::DB_ENTRY_DELETE;
    1009           0 :         stats_[RX].unreach++;
    1010             :     }
    1011           0 : }
    1012             : 
    1013           0 : void BgpXmppChannel::CreateType7MvpnRouteRequest(IpAddress grp_address,
    1014             :         IpAddress src_address, bool add_change, uint64_t subscription_gen_id,
    1015             :         DBRequest& req) {
    1016           0 :     RouteDistinguisher mc_rd =  RouteDistinguisher::kZeroRd;
    1017             :     MvpnPrefix mc_prefix(MvpnPrefix::SourceTreeJoinRoute, mc_rd, 0,
    1018           0 :             grp_address.to_v4(), src_address.to_v4());
    1019           0 :     uint32_t flags = BgpPath::ResolveNexthop;
    1020             : 
    1021             :     // Build and enqueue a DB request for route-addition
    1022           0 :     req.key.reset(new MvpnTable::RequestKey(mc_prefix, peer_.get()));
    1023             : 
    1024           0 :     if (add_change) {
    1025           0 :         req.oper = DBRequest::DB_ENTRY_ADD_CHANGE;
    1026             : 
    1027           0 :         BgpAttrSpec attrs;
    1028             : 
    1029             :         // Next-hop ip address
    1030           0 :         BgpAttrNextHop nexthop(src_address);
    1031           0 :         attrs.push_back(&nexthop);
    1032             : 
    1033           0 :         BgpAttrPtr attr = bgp_server_->attr_db()->Locate(attrs);
    1034           0 :         req.data.reset(new MvpnTable::RequestData(
    1035           0 :             attr, flags, 0, 0, subscription_gen_id));
    1036           0 :         stats_[RX].reach++;
    1037           0 :     } else {
    1038           0 :         req.oper = DBRequest::DB_ENTRY_DELETE;
    1039           0 :         stats_[RX].unreach++;
    1040             :     }
    1041           0 : }
    1042             : 
    1043           0 : bool BgpXmppChannel::ProcessMvpnItem(string vrf_name,
    1044             :     const pugi::xml_node &node, bool add_change) {
    1045           0 :     MvpnItemType item;
    1046           0 :     item.Clear();
    1047             : 
    1048           0 :     if (!item.XmlParse(node)) {
    1049           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1050             :             BGP_LOG_FLAG_ALL, "Invalid multicast route message received");
    1051           0 :         return false;
    1052             :     }
    1053             : 
    1054           0 :     if (item.entry.nlri.af != BgpAf::IPv4) {
    1055           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1056             :             "Unsupported address family " << item.entry.nlri.af <<
    1057             :             " for multicast route");
    1058           0 :         return false;
    1059             :     }
    1060             : 
    1061           0 :     if (item.entry.nlri.safi != BgpAf::MVpn) {
    1062           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1063             :             BGP_LOG_FLAG_ALL, "Unsupported subsequent address family " <<
    1064             :             item.entry.nlri.safi << " for multicast route");
    1065           0 :         return false;
    1066             :     }
    1067             : 
    1068           0 :     if (item.entry.nlri.group.empty()) {
    1069           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1070             :             "Mandatory group address not specified");
    1071           0 :         return false;
    1072             :     }
    1073             : 
    1074           0 :     error_code error;
    1075           0 :     IpAddress grp_address = IpAddress::from_string("0.0.0.0", error);
    1076           0 :     if (!XmppDecodeAddress(item.entry.nlri.af,
    1077             :         item.entry.nlri.group, &grp_address, false)) {
    1078           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1079             :             "Bad group address " << item.entry.nlri.group);
    1080           0 :         return false;
    1081             :     }
    1082             : 
    1083           0 :     if (item.entry.nlri.source.empty()) {
    1084           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1085             :             "Mandatory source address not specified");
    1086           0 :         return false;
    1087             :     }
    1088             : 
    1089           0 :     IpAddress src_address = IpAddress::from_string("0.0.0.0", error);
    1090           0 :     if (!XmppDecodeAddress(item.entry.nlri.af,
    1091             :         item.entry.nlri.source, &src_address, true)) {
    1092           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1093             :             "Bad source address " << item.entry.nlri.source);
    1094           0 :         return false;
    1095             :     }
    1096             : 
    1097             :     bool subscribe_pending;
    1098             :     int instance_id;
    1099             :     uint64_t subscription_gen_id;
    1100             :     BgpTable *table;
    1101           0 :     if (!VerifyMembership(vrf_name, Address::MVPN, &table, &instance_id,
    1102             :         &subscription_gen_id, &subscribe_pending, add_change)) {
    1103           0 :         channel_->Close();
    1104           0 :         return false;
    1105             :     }
    1106             : 
    1107           0 :     int rt_type = item.entry.nlri.route_type;
    1108           0 :     DBRequest req;
    1109             :     // Build the key to the Multicast DBTable
    1110           0 :     if (rt_type == MvpnPrefix::SourceTreeJoinRoute) {
    1111           0 :         CreateType7MvpnRouteRequest(grp_address, src_address, add_change,
    1112             :                 subscription_gen_id, req);
    1113           0 :     } else if (rt_type == MvpnPrefix::SourceActiveADRoute) {
    1114           0 :         CreateType5MvpnRouteRequest(grp_address, src_address, add_change,
    1115             :                 subscription_gen_id, instance_id, req, item.entry.next_hop);
    1116             :     } else {
    1117           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1118             :             "Unsupported route type " << item.entry.nlri.route_type);
    1119           0 :         return false;
    1120             :     }
    1121             : 
    1122             :     // Need to locate path resolver if not done already
    1123           0 :     assert(table);
    1124           0 :     table->LocatePathResolver();
    1125             : 
    1126             :     // Defer all requests till subscribe is processed.
    1127           0 :     if (subscribe_pending) {
    1128           0 :         DBRequest *request_entry = new DBRequest();
    1129           0 :         request_entry->Swap(&req);
    1130             :         string table_name =
    1131           0 :             RoutingInstance::GetTableName(vrf_name, Address::MVPN);
    1132           0 :         defer_q_.insert(make_pair(
    1133           0 :             make_pair(vrf_name, table_name), request_entry));
    1134           0 :         return true;
    1135           0 :     }
    1136             : 
    1137           0 :     BGP_LOG_PEER_INSTANCE(Peer(), vrf_name,
    1138             :         SandeshLevel::SYS_DEBUG, BGP_LOG_FLAG_TRACE,
    1139             :         "Multicast group " << item.entry.nlri.group <<
    1140             :         " source " << item.entry.nlri.source <<
    1141             :         " enqueued for " << (add_change ? "add/change" : "delete"));
    1142           0 :     table->Enqueue(&req);
    1143           0 :     return true;
    1144           0 : }
    1145             : 
    1146           0 : bool BgpXmppChannel::ProcessItem(string vrf_name,
    1147             :     const pugi::xml_node &node, bool add_change, int primary_instance_id) {
    1148           0 :     ItemType item;
    1149           0 :     item.Clear();
    1150             : 
    1151           0 :     if (!item.XmlParse(node)) {
    1152           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1153             :             "Invalid inet route message received");
    1154           0 :         return false;
    1155             :     }
    1156             : 
    1157           0 :     if (item.entry.nlri.af != BgpAf::IPv4) {
    1158           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1159             :             "Unsupported address family " << item.entry.nlri.af <<
    1160             :             " for inet route " << item.entry.nlri.address);
    1161           0 :         return false;
    1162             :     }
    1163             : 
    1164           0 :     if ((item.entry.nlri.safi != BgpAf::Unicast) &&
    1165           0 :           (item.entry.nlri.safi != BgpAf::Mpls)) {
    1166           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1167             :             "Unsupported subsequent address family " << item.entry.nlri.safi <<
    1168             :             " for inet route " << item.entry.nlri.address);
    1169           0 :         return false;
    1170             :     }
    1171           0 :     error_code error;
    1172             :     Ip4Prefix inet_prefix =
    1173           0 :         Ip4Prefix::FromString(item.entry.nlri.address, &error);
    1174           0 :     if (error) {
    1175           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1176             :             "Bad inet route " << item.entry.nlri.address);
    1177           0 :         return false;
    1178             :     }
    1179             : 
    1180           0 :     if (add_change && item.entry.next_hops.next_hop.empty()) {
    1181           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1182             :             "Missing next-hops for inet route " << inet_prefix.ToString());
    1183           0 :         return false;
    1184             :     }
    1185             : 
    1186             :     // Rules for routes in master instance:
    1187             :     // - Label must be 0 unless it is INETMPLS
    1188             :     // - Tunnel encapsulation is not required
    1189             :     // - Do not add SourceRd and ExtCommunitySpec
    1190           0 :     bool master = (vrf_name == BgpConfigManager::kMasterInstance);
    1191             :     bool subscribe_pending;
    1192             :     int instance_id;
    1193             :     uint64_t subscription_gen_id;
    1194             :     BgpTable *table;
    1195           0 :     Address::Family family = BgpAf::AfiSafiToFamily(item.entry.nlri.af,
    1196           0 :                                                     item.entry.nlri.safi);
    1197           0 :     if (!VerifyMembership(vrf_name, family, &table, &instance_id,
    1198             :         &subscription_gen_id, &subscribe_pending, add_change)) {
    1199           0 :         channel_->Close();
    1200           0 :         return false;
    1201             :     }
    1202             : 
    1203           0 :     DBRequest req;
    1204           0 :     req.key.reset(new InetTable::RequestKey(inet_prefix, peer_.get()));
    1205             : 
    1206           0 :     IpAddress nh_address(Ip4Address(0));
    1207           0 :     uint32_t label = 0;
    1208           0 :     uint32_t flags = 0;
    1209           0 :     ExtCommunitySpec ext;
    1210           0 :     LargeCommunitySpec largecomm;
    1211           0 :     CommunitySpec comm;
    1212             : 
    1213           0 :     if (add_change) {
    1214           0 :         req.oper = DBRequest::DB_ENTRY_ADD_CHANGE;
    1215           0 :         BgpAttrSpec attrs;
    1216             : 
    1217           0 :         const NextHopListType &inh_list = item.entry.next_hops;
    1218             : 
    1219             :         // Agents should send only one next-hop in the item.
    1220           0 :         if (inh_list.next_hop.size() != 1) {
    1221           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1222             :                 BGP_LOG_FLAG_ALL,
    1223             :                 "More than one nexthop received for inet route " <<
    1224             :                 inet_prefix.ToString());
    1225           0 :             return false;
    1226             :         }
    1227             : 
    1228           0 :         NextHopListType::const_iterator nit = inh_list.begin();
    1229             : 
    1230           0 :         IpAddress nhop_address(Ip4Address(0));
    1231           0 :         if (!XmppDecodeAddress(nit->af, nit->address, &nhop_address)) {
    1232           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1233             :                 BGP_LOG_FLAG_ALL,
    1234             :                 "Bad nexthop address " << nit->address <<
    1235             :                 " for inet route " << inet_prefix.ToString());
    1236           0 :             return false;
    1237             :         }
    1238             : 
    1239           0 :         if (nit->label > EvpnPrefix::kMaxVniSigned ||
    1240           0 :             ((master && nit->label) &&
    1241             :              !(family == Address::INETMPLS))) {
    1242           0 :              BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1243             :                 BGP_LOG_FLAG_ALL,
    1244             :                 "Bad label " << nit->label <<
    1245             :                 " for inet route " << inet_prefix.ToString());
    1246           0 :             return false;
    1247             :         }
    1248           0 :         if ((!master || (master && (family == Address::INETMPLS))) &&
    1249           0 :             !nit->label) {
    1250           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1251             :                BGP_LOG_FLAG_ALL,
    1252             :                "Bad label " << nit->label <<
    1253             :                " for inet route in master instance(mpls)/non master instance" <<
    1254             :                inet_prefix.ToString());
    1255           0 :            return false;
    1256             :         }
    1257             : 
    1258           0 :         nh_address = nhop_address;
    1259           0 :         label = nit->label;
    1260             : 
    1261             :         // Process tunnel encapsulation list.
    1262           0 :         bool no_tunnel_encap = true;
    1263           0 :         bool no_valid_tunnel_encap = true;
    1264           0 :         for (TunnelEncapsulationListType::const_iterator eit =
    1265           0 :             nit->tunnel_encapsulation_list.begin();
    1266           0 :             eit != nit->tunnel_encapsulation_list.end(); ++eit) {
    1267           0 :             no_tunnel_encap = false;
    1268           0 :             TunnelEncap tun_encap(*eit);
    1269           0 :             if (tun_encap.tunnel_encap() == TunnelEncapType::UNSPEC)
    1270           0 :                 continue;
    1271           0 :             no_valid_tunnel_encap = false;
    1272           0 :             ext.communities.push_back(tun_encap.GetExtCommunityValue());
    1273             :         }
    1274             : 
    1275             :         // Mark the path as infeasible if all tunnel encaps published
    1276             :         // by agent are invalid.
    1277           0 :         if (!no_tunnel_encap && no_valid_tunnel_encap && !master) {
    1278           0 :             flags = BgpPath::NoTunnelEncap;
    1279             :         }
    1280             : 
    1281             :         // Process router-mac as ext-community.
    1282           0 :         if (!nit->mac.empty()) {
    1283             :             MacAddress mac_addr =
    1284           0 :                         MacAddress::FromString(nit->mac, &error);
    1285           0 :             if (error) {
    1286           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1287             :                     BGP_LOG_FLAG_ALL,
    1288             :                     "Bad next-hop mac address " << nit->mac);
    1289           0 :                 return false;
    1290             :             }
    1291           0 :             if (!mac_addr.IsZero()) {
    1292           0 :                 RouterMac router_mac(mac_addr);
    1293           0 :                 ext.communities.push_back(router_mac.GetExtCommunityValue());
    1294             :             }
    1295             :         }
    1296             : 
    1297             :         // Process tags list.
    1298           0 :         uint64_t tag_val = 0;
    1299           0 :         for (TagListType::const_iterator tit = nit->tag_list.begin();
    1300           0 :             tit != nit->tag_list.end(); ++tit) {
    1301           0 :             tag_val = nit->is_new_tags_list ? *tit :
    1302           0 :                 ((0xFFFF0000 & *tit) << 16) | (0x0000FFFF & *tit);
    1303           0 :             TagLC tag_lc(bgp_server_->autonomous_system(), tag_val);
    1304           0 :             for (const auto &value_data : tag_lc.GetLargeCommunityValue()) {
    1305           0 :                 largecomm.communities.push_back(value_data);
    1306           0 :             }
    1307             :         }
    1308             : 
    1309             :         // Process local sequence_number
    1310           0 :         if (nit->local_sequence_number) {
    1311           0 :             LocalSequenceNumber lsn (nit->local_sequence_number);
    1312           0 :             ext.communities.push_back(lsn.GetExtCommunityValue());
    1313             :         }
    1314             : 
    1315           0 :         BgpAttrLocalPref local_pref(item.entry.local_preference);
    1316           0 :         if (local_pref.local_pref != 0)
    1317           0 :             attrs.push_back(&local_pref);
    1318             : 
    1319             :         // If there's no explicit med, calculate it automatically from the
    1320             :         // local pref.
    1321           0 :         uint32_t med_value = item.entry.med;
    1322           0 :         if (!med_value)
    1323           0 :             med_value = GetMedFromLocalPref(local_pref.local_pref);
    1324           0 :         BgpAttrMultiExitDisc med(med_value);
    1325           0 :         if (med.med != 0)
    1326           0 :             attrs.push_back(&med);
    1327             : 
    1328             :         // Process community tags.
    1329           0 :         const CommunityTagListType &ict_list = item.entry.community_tag_list;
    1330           0 :         for (CommunityTagListType::const_iterator cit = ict_list.begin();
    1331           0 :             cit != ict_list.end(); ++cit) {
    1332           0 :             error_code error;
    1333             :             uint32_t rt_community =
    1334           0 :                 CommunityType::CommunityFromString(*cit, &error);
    1335           0 :             if (error)
    1336           0 :                 continue;
    1337           0 :             comm.communities.push_back(rt_community);
    1338             :         }
    1339             : 
    1340           0 :         uint32_t addr = nh_address.to_v4().to_ulong();
    1341           0 :         BgpAttrNextHop nexthop(addr);
    1342           0 :         attrs.push_back(&nexthop);
    1343             :         uint16_t cluster_seed =
    1344           0 :             bgp_server_->global_config()->rd_cluster_seed();
    1345           0 :         BgpAttrSourceRd source_rd;
    1346           0 :         if (!master || primary_instance_id) {
    1347           0 :             if (master)
    1348           0 :                 instance_id = primary_instance_id;
    1349           0 :             if (cluster_seed) {
    1350           0 :                 source_rd = BgpAttrSourceRd(
    1351           0 :                     RouteDistinguisher(cluster_seed, addr, instance_id));
    1352             :             } else {
    1353           0 :                 source_rd = BgpAttrSourceRd(
    1354           0 :                     RouteDistinguisher(addr, instance_id));
    1355             :             }
    1356           0 :             attrs.push_back(&source_rd);
    1357             :         }
    1358             : 
    1359             :         // Process security group list.
    1360           0 :         uint16_t sg_index = 0;
    1361           0 :         const SecurityGroupListType &isg_list = item.entry.security_group_list;
    1362           0 :         for (SecurityGroupListType::const_iterator sit = isg_list.begin();
    1363           0 :             sit != isg_list.end(); ++sit) {
    1364           0 :             if (bgp_server_->autonomous_system() <= AS2_MAX) {
    1365           0 :                 SecurityGroup sg(bgp_server_->autonomous_system(), *sit);
    1366           0 :                 ext.communities.push_back(sg.GetExtCommunityValue());
    1367             :             } else {
    1368           0 :                 SecurityGroup sg(sg_index, *sit);
    1369           0 :                 SecurityGroup4ByteAs sg4(bgp_server_->autonomous_system(),
    1370           0 :                                          sg_index++);
    1371           0 :                 ext.communities.push_back(sg4.GetExtCommunityValue());
    1372           0 :                 ext.communities.push_back(sg.GetExtCommunityValue());
    1373             :             }
    1374             :         }
    1375             : 
    1376           0 :         if (item.entry.mobility.seqno) {
    1377           0 :             MacMobility mm(item.entry.mobility.seqno,
    1378           0 :                            item.entry.mobility.sticky);
    1379           0 :             ext.communities.push_back(mm.GetExtCommunityValue());
    1380           0 :         } else if (item.entry.sequence_number) {
    1381           0 :             MacMobility mm(item.entry.sequence_number);
    1382           0 :             ext.communities.push_back(mm.GetExtCommunityValue());
    1383             :         }
    1384             : 
    1385             :         // Process load-balance extended community.
    1386           0 :         LoadBalance load_balance(item.entry.load_balance);
    1387           0 :         if (!load_balance.IsDefault())
    1388           0 :             ext.communities.push_back(load_balance.GetExtCommunityValue());
    1389             : 
    1390           0 :         if (!comm.communities.empty())
    1391           0 :             attrs.push_back(&comm);
    1392           0 :         if (!master && !ext.communities.empty())
    1393           0 :             attrs.push_back(&ext);
    1394           0 :         if (!master && !largecomm.communities.empty())
    1395           0 :             attrs.push_back(&largecomm);
    1396             : 
    1397             :         // Process sub-protocol(route types)
    1398           0 :         BgpAttrSubProtocol sbp(item.entry.sub_protocol);
    1399           0 :         attrs.push_back(&sbp);
    1400             : 
    1401           0 :         BgpAttrPtr attr = bgp_server_->attr_db()->Locate(attrs);
    1402           0 :         req.data.reset(new BgpTable::RequestData(
    1403           0 :             attr, flags, label, 0, subscription_gen_id));
    1404           0 :     } else {
    1405           0 :         req.oper = DBRequest::DB_ENTRY_DELETE;
    1406             :     }
    1407             : 
    1408             :     // Defer all requests till subscribe is processed.
    1409           0 :     if (subscribe_pending) {
    1410           0 :         DBRequest *request_entry = new DBRequest();
    1411           0 :         request_entry->Swap(&req);
    1412             :         string table_name =
    1413           0 :             RoutingInstance::GetTableName(vrf_name, family);
    1414           0 :         defer_q_.insert(make_pair(
    1415           0 :             make_pair(vrf_name, table_name), request_entry));
    1416           0 :        return true;
    1417           0 :     }
    1418             : 
    1419           0 :     assert(table);
    1420           0 :     BGP_LOG_PEER_INSTANCE(Peer(), vrf_name,
    1421             :         SandeshLevel::SYS_DEBUG, BGP_LOG_FLAG_TRACE,
    1422             :         "Inet route " << item.entry.nlri.address <<
    1423             :         " with next-hop " << nh_address << " and label " << label <<
    1424             :         " enqueued for " << (add_change ? "add/change" : "delete") <<
    1425             :         " to table " << table->name());
    1426           0 :     table->Enqueue(&req);
    1427             : 
    1428           0 :     if (add_change) {
    1429           0 :         stats_[RX].reach++;
    1430             :     } else {
    1431           0 :         stats_[RX].unreach++;
    1432             :     }
    1433             : 
    1434           0 :     return true;
    1435           0 : }
    1436             : 
    1437           0 : bool BgpXmppChannel::ProcessInet6Item(string vrf_name,
    1438             :     const pugi::xml_node &node, bool add_change) {
    1439           0 :     ItemType item;
    1440           0 :     item.Clear();
    1441             : 
    1442           0 :     if (!item.XmlParse(node)) {
    1443           0 :         error_stats().incr_inet6_rx_bad_xml_token_count();
    1444           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1445             :             "Invalid inet6 route message received");
    1446           0 :         return false;
    1447             :     }
    1448             : 
    1449           0 :     if (item.entry.nlri.af != BgpAf::IPv6) {
    1450           0 :         error_stats().incr_inet6_rx_bad_afi_safi_count();
    1451           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1452             :             "Unsupported address family " << item.entry.nlri.af <<
    1453             :             " for inet6 route " << item.entry.nlri.address);
    1454           0 :         return false;
    1455             :     }
    1456             : 
    1457           0 :     if (item.entry.nlri.safi != BgpAf::Unicast) {
    1458           0 :         error_stats().incr_inet6_rx_bad_afi_safi_count();
    1459           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1460             :             "Unsupported subsequent address family " << item.entry.nlri.safi <<
    1461             :             " for inet6 route " << item.entry.nlri.address);
    1462           0 :         return false;
    1463             :     }
    1464             : 
    1465           0 :     error_code error;
    1466             :     Inet6Prefix inet6_prefix =
    1467           0 :         Inet6Prefix::FromString(item.entry.nlri.address, &error);
    1468           0 :     if (error) {
    1469           0 :         error_stats().incr_inet6_rx_bad_prefix_count();
    1470           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1471             :             "Bad inet6 route " << item.entry.nlri.address);
    1472           0 :         return false;
    1473             :     }
    1474             : 
    1475           0 :     if (add_change && item.entry.next_hops.next_hop.empty()) {
    1476           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1477             :             "Missing next-hops for inet6 route " << inet6_prefix.ToString());
    1478           0 :         return false;
    1479             :     }
    1480             : 
    1481             :     // Rules for routes in master instance:
    1482             :     // - Label must be 0
    1483             :     // - Tunnel encapsulation is not required
    1484             :     // - Do not add SourceRd and ExtCommunitySpec
    1485           0 :     bool master = (vrf_name == BgpConfigManager::kMasterInstance);
    1486             : 
    1487             :     // vector<Address::Family> family_list = list_of(Address::INET6)(Address::EVPN);
    1488           0 :     vector<Address::Family> family_list = list_of(Address::INET6);
    1489           0 :     BOOST_FOREACH(Address::Family family, family_list) {
    1490             :         bool subscribe_pending;
    1491             :         int instance_id;
    1492             :         uint64_t subscription_gen_id;
    1493             :         BgpTable *table;
    1494           0 :         if (!VerifyMembership(vrf_name, family, &table, &instance_id,
    1495             :             &subscription_gen_id, &subscribe_pending, add_change)) {
    1496           0 :             channel_->Close();
    1497           0 :             return false;
    1498             :         }
    1499             : 
    1500           0 :         DBRequest req;
    1501           0 :         if (family == Address::INET6) {
    1502           0 :             req.key.reset(new Inet6Table::RequestKey(inet6_prefix, peer_.get()));
    1503             :         } else {
    1504             :             EvpnPrefix evpn_prefix(RouteDistinguisher::kZeroRd,
    1505           0 :                 inet6_prefix.addr(), inet6_prefix.prefixlen());
    1506           0 :             req.key.reset(new EvpnTable::RequestKey(evpn_prefix, peer_.get()));
    1507             :         }
    1508             : 
    1509           0 :         IpAddress nh_address(Ip4Address(0));
    1510           0 :         uint32_t label = 0;
    1511           0 :         uint32_t flags = 0;
    1512           0 :         ExtCommunitySpec ext;
    1513           0 :         LargeCommunitySpec largecomm;
    1514           0 :         CommunitySpec comm;
    1515             : 
    1516           0 :         if (add_change) {
    1517           0 :             req.oper = DBRequest::DB_ENTRY_ADD_CHANGE;
    1518           0 :             BgpAttrSpec attrs;
    1519             : 
    1520           0 :             const NextHopListType &inh_list = item.entry.next_hops;
    1521             : 
    1522             :             // Agents should send only one next-hop in the item.
    1523           0 :             if (inh_list.next_hop.size() != 1) {
    1524           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1525             :                     BGP_LOG_FLAG_ALL,
    1526             :                     "More than one nexthop received for inet6 route " <<
    1527             :                     inet6_prefix.ToString());
    1528           0 :                 return false;
    1529             :             }
    1530             : 
    1531           0 :             NextHopListType::const_iterator nit = inh_list.begin();
    1532             : 
    1533           0 :             IpAddress nhop_address(Ip4Address(0));
    1534           0 :             if (!XmppDecodeAddress(nit->af, nit->address, &nhop_address)) {
    1535           0 :                 error_stats().incr_inet6_rx_bad_nexthop_count();
    1536           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1537             :                     BGP_LOG_FLAG_ALL,
    1538             :                     "Bad nexthop address " << nit->address <<
    1539             :                     " for inet6 route " << inet6_prefix.ToString());
    1540           0 :                 return false;
    1541             :             }
    1542             : 
    1543           0 :             if (family == Address::EVPN) {
    1544           0 :                 if (nit->vni > EvpnPrefix::kMaxVniSigned) {
    1545           0 :                     BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1546             :                         BGP_LOG_FLAG_ALL,
    1547             :                         "Bad label " << nit->vni <<
    1548             :                         " for inet6 route " << inet6_prefix.ToString());
    1549           0 :                     return false;
    1550             :                 }
    1551           0 :                 if (!nit->vni)
    1552           0 :                     continue;
    1553           0 :                 if (nit->mac.empty())
    1554           0 :                     continue;
    1555             : 
    1556             :                 MacAddress mac_addr =
    1557           0 :                     MacAddress::FromString(nit->mac, &error);
    1558           0 :                 if (error) {
    1559           0 :                     BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1560             :                         BGP_LOG_FLAG_ALL,
    1561             :                         "Bad next-hop mac address " << nit->mac);
    1562           0 :                     return false;
    1563             :                 }
    1564           0 :                 RouterMac router_mac(mac_addr);
    1565           0 :                 ext.communities.push_back(router_mac.GetExtCommunityValue());
    1566             :             } else {
    1567           0 :                 if (nit->label > EvpnPrefix::kMaxVniSigned || (master && nit->label)) {
    1568           0 :                     BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1569             :                         BGP_LOG_FLAG_ALL,
    1570             :                         "Bad label " << nit->label <<
    1571             :                         " for inet6 route " << inet6_prefix.ToString());
    1572           0 :                     return false;
    1573             :                 }
    1574           0 :                 if (!master && !nit->label)
    1575           0 :                     continue;
    1576             :             }
    1577             : 
    1578           0 :             nh_address = nhop_address;
    1579           0 :             if (family == Address::INET6) {
    1580           0 :                 label = nit->label;
    1581             :             } else {
    1582           0 :                 label = nit->vni;
    1583             :             }
    1584             : 
    1585             :             // Process tunnel encapsulation list.
    1586           0 :             bool no_tunnel_encap = true;
    1587           0 :             bool no_valid_tunnel_encap = true;
    1588           0 :             for (TunnelEncapsulationListType::const_iterator eit =
    1589           0 :                 nit->tunnel_encapsulation_list.begin();
    1590           0 :                 eit != nit->tunnel_encapsulation_list.end(); ++eit) {
    1591           0 :                 no_tunnel_encap = false;
    1592           0 :                 TunnelEncap tun_encap(*eit);
    1593           0 :                 if (tun_encap.tunnel_encap() == TunnelEncapType::UNSPEC)
    1594           0 :                     continue;
    1595           0 :                 no_valid_tunnel_encap = false;
    1596           0 :                 ext.communities.push_back(tun_encap.GetExtCommunityValue());
    1597             :             }
    1598             : 
    1599             :             // Mark the path as infeasible if all tunnel encaps published
    1600             :             // by agent are invalid.
    1601           0 :             if (!no_tunnel_encap && no_valid_tunnel_encap && !master) {
    1602           0 :                 flags = BgpPath::NoTunnelEncap;
    1603             :             }
    1604             : 
    1605             :             // Process router-mac as ext-community.
    1606           0 :             if (!nit->mac.empty()) {
    1607             :                 MacAddress mac_addr =
    1608           0 :                         MacAddress::FromString(nit->mac, &error);
    1609           0 :                 if (error) {
    1610           0 :                     BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1611             :                         BGP_LOG_FLAG_ALL,
    1612             :                         "Bad next-hop mac address " << nit->mac);
    1613           0 :                     return false;
    1614             :                 }
    1615           0 :                 if (!mac_addr.IsZero()) {
    1616           0 :                     RouterMac router_mac(mac_addr);
    1617           0 :                     ext.communities.push_back(router_mac.GetExtCommunityValue());
    1618             :                 }
    1619             :             }
    1620             : 
    1621             :             // Process tags list.
    1622           0 :             uint64_t tag_val = 0;
    1623           0 :             for (TagListType::const_iterator tit = nit->tag_list.begin();
    1624           0 :                 tit != nit->tag_list.end(); ++tit) {
    1625           0 :                 tag_val = nit->is_new_tags_list ? *tit :
    1626           0 :                     ((0xFFFF0000 & *tit) << 16) | (0x0000FFFF & *tit);
    1627           0 :                 TagLC tag_lc(bgp_server_->autonomous_system(), tag_val);
    1628           0 :                 for (const auto &value_data :
    1629           0 :                     tag_lc.GetLargeCommunityValue()) {
    1630           0 :                     largecomm.communities.push_back(value_data);
    1631           0 :                 }
    1632             :             }
    1633             : 
    1634             :             // Process local sequence_number
    1635           0 :             if (nit->local_sequence_number) {
    1636           0 :                 LocalSequenceNumber lsn (nit->local_sequence_number);
    1637           0 :                 ext.communities.push_back(lsn.GetExtCommunityValue());
    1638             :             }
    1639             : 
    1640           0 :             BgpAttrLocalPref local_pref(item.entry.local_preference);
    1641           0 :             if (local_pref.local_pref != 0)
    1642           0 :                 attrs.push_back(&local_pref);
    1643             : 
    1644             :             // If there's no explicit med, calculate it automatically from the
    1645             :             // local pref.
    1646           0 :             uint32_t med_value = item.entry.med;
    1647           0 :             if (!med_value)
    1648           0 :                 med_value = GetMedFromLocalPref(local_pref.local_pref);
    1649           0 :             BgpAttrMultiExitDisc med(med_value);
    1650           0 :             if (med.med != 0)
    1651           0 :                 attrs.push_back(&med);
    1652             : 
    1653             :             // Process community tags.
    1654           0 :             const CommunityTagListType &ict_list =
    1655             :                 item.entry.community_tag_list;
    1656           0 :             for (CommunityTagListType::const_iterator cit = ict_list.begin();
    1657           0 :                 cit != ict_list.end(); ++cit) {
    1658           0 :                 error_code error;
    1659             :                 uint32_t rt_community =
    1660           0 :                     CommunityType::CommunityFromString(*cit, &error);
    1661           0 :                 if (error)
    1662           0 :                     continue;
    1663           0 :                 comm.communities.push_back(rt_community);
    1664             :             }
    1665             : 
    1666           0 :             BgpAttrNextHop nexthop(nh_address);
    1667           0 :             attrs.push_back(&nexthop);
    1668             : 
    1669           0 :             BgpAttrSourceRd source_rd;
    1670           0 :             if (!master) {
    1671           0 :                 uint32_t addr = nh_address.to_v4().to_ulong();
    1672             :                 uint16_t cluster_seed =
    1673           0 :                   bgp_server_->global_config()->rd_cluster_seed();
    1674           0 :                 if (cluster_seed) {
    1675           0 :                     source_rd = BgpAttrSourceRd(
    1676           0 :                         RouteDistinguisher(cluster_seed, addr, instance_id));
    1677             :                 } else {
    1678           0 :                     source_rd = BgpAttrSourceRd(
    1679           0 :                         RouteDistinguisher(addr, instance_id));
    1680             :                 }
    1681           0 :                 attrs.push_back(&source_rd);
    1682             :             }
    1683             : 
    1684             :             // Process security group list.
    1685           0 :             const SecurityGroupListType &isg_list =
    1686             :                 item.entry.security_group_list;
    1687           0 :             uint16_t sg_index = 0;
    1688           0 :             for (SecurityGroupListType::const_iterator sit = isg_list.begin();
    1689           0 :                 sit != isg_list.end(); ++sit) {
    1690           0 :                 if (bgp_server_->autonomous_system() <= AS2_MAX) {
    1691           0 :                     SecurityGroup sg(bgp_server_->autonomous_system(), *sit);
    1692           0 :                     ext.communities.push_back(sg.GetExtCommunityValue());
    1693             :                 } else {
    1694           0 :                     SecurityGroup sg(sg_index, *sit);
    1695           0 :                     SecurityGroup4ByteAs sg4(bgp_server_->autonomous_system(),
    1696           0 :                                              sg_index++);
    1697           0 :                     ext.communities.push_back(sg4.GetExtCommunityValue());
    1698           0 :                     ext.communities.push_back(sg.GetExtCommunityValue());
    1699             :                 }
    1700             :             }
    1701             : 
    1702           0 :             if (item.entry.mobility.seqno) {
    1703           0 :                 MacMobility mm(item.entry.mobility.seqno,
    1704           0 :                     item.entry.mobility.sticky);
    1705           0 :                 ext.communities.push_back(mm.GetExtCommunityValue());
    1706           0 :             } else if (item.entry.sequence_number) {
    1707           0 :                 MacMobility mm(item.entry.sequence_number);
    1708           0 :                 ext.communities.push_back(mm.GetExtCommunityValue());
    1709             :             }
    1710             : 
    1711             :             // Process load-balance extended community.
    1712           0 :             LoadBalance load_balance(item.entry.load_balance);
    1713           0 :             if (!load_balance.IsDefault())
    1714           0 :                 ext.communities.push_back(load_balance.GetExtCommunityValue());
    1715             : 
    1716             :             // Process sub-protocol(route types)
    1717           0 :             BgpAttrSubProtocol sbp(item.entry.sub_protocol);
    1718           0 :             attrs.push_back(&sbp);
    1719             : 
    1720           0 :             if (!comm.communities.empty())
    1721           0 :                 attrs.push_back(&comm);
    1722           0 :             if (!master && !ext.communities.empty())
    1723           0 :                 attrs.push_back(&ext);
    1724           0 :             if (!master && !largecomm.communities.empty())
    1725           0 :                 attrs.push_back(&largecomm);
    1726             : 
    1727           0 :             BgpAttrPtr attr = bgp_server_->attr_db()->Locate(attrs);
    1728           0 :             req.data.reset(new BgpTable::RequestData(
    1729           0 :                 attr, flags, label, 0, subscription_gen_id));
    1730           0 :         } else {
    1731           0 :             req.oper = DBRequest::DB_ENTRY_DELETE;
    1732             :         }
    1733             : 
    1734             :         // Defer all requests till subscribe is processed.
    1735           0 :         if (subscribe_pending) {
    1736           0 :             DBRequest *request_entry = new DBRequest();
    1737           0 :             request_entry->Swap(&req);
    1738             :             string table_name =
    1739           0 :                 RoutingInstance::GetTableName(vrf_name, family);
    1740           0 :             defer_q_.insert(make_pair(
    1741           0 :                 make_pair(vrf_name, table_name), request_entry));
    1742           0 :             continue;
    1743           0 :         }
    1744             : 
    1745           0 :         assert(table);
    1746           0 :         BGP_LOG_PEER_INSTANCE(Peer(), vrf_name,
    1747             :             SandeshLevel::SYS_DEBUG, BGP_LOG_FLAG_TRACE,
    1748             :             "Inet6 route " << item.entry.nlri.address <<
    1749             :             " with next-hop " << nh_address << " and label " << label <<
    1750             :             " enqueued for " << (add_change ? "add/change" : "delete") <<
    1751             :             " to table " << table->name());
    1752           0 :         table->Enqueue(&req);
    1753           0 :     }
    1754             : 
    1755           0 :     if (add_change) {
    1756           0 :         stats_[RX].reach++;
    1757             :     } else {
    1758           0 :         stats_[RX].unreach++;
    1759             :     }
    1760             : 
    1761           0 :     return true;
    1762           0 : }
    1763             : 
    1764           0 : bool BgpXmppChannel::ProcessEnetItem(string vrf_name,
    1765             :     const pugi::xml_node &node, bool add_change) {
    1766           0 :     EnetItemType item;
    1767           0 :     item.Clear();
    1768             : 
    1769           0 :     if (!item.XmlParse(node)) {
    1770           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1771             :             "Invalid enet route message received");
    1772           0 :         return false;
    1773             :     }
    1774             : 
    1775           0 :     if (item.entry.nlri.af != BgpAf::L2Vpn) {
    1776           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1777             :             "Unsupported address family " << item.entry.nlri.af <<
    1778             :             " for enet route " << item.entry.nlri.address);
    1779           0 :         return false;
    1780             :     }
    1781             : 
    1782           0 :     if (item.entry.nlri.safi != BgpAf::Enet) {
    1783           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1784             :             "Unsupported subsequent address family " << item.entry.nlri.safi <<
    1785             :             " for enet route " << item.entry.nlri.mac);
    1786           0 :         return false;
    1787             :     }
    1788             : 
    1789           0 :     bool type6 = false;
    1790           0 :     error_code error;
    1791           0 :     IpAddress group= IpAddress::from_string("0.0.0.0", error);
    1792           0 :     if (!item.entry.nlri.group.empty()) {
    1793           0 :         type6 = true;
    1794           0 :         if (!XmppDecodeAddress(item.entry.nlri.af,
    1795             :                     item.entry.nlri.group, &group, false)) {
    1796           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1797             :                 "Bad group address " << item.entry.nlri.group);
    1798             :         }
    1799             :     }
    1800             : 
    1801           0 :     IpAddress source = IpAddress::from_string("0.0.0.0", error);
    1802           0 :     if (!item.entry.nlri.source.empty() && !XmppDecodeAddress(
    1803             :                 item.entry.nlri.af, item.entry.nlri.source, &source, true)) {
    1804           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1805             :             "Bad source address " << item.entry.nlri.source);
    1806             :     }
    1807             : 
    1808             :     //error_code error;
    1809           0 :     MacAddress mac_addr = MacAddress::FromString(item.entry.nlri.mac, &error);
    1810             : 
    1811           0 :     if (error) {
    1812           0 :         BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1813             :             "Bad mac address " << item.entry.nlri.mac);
    1814           0 :         return false;
    1815             :     }
    1816             : 
    1817           0 :     bool type2 = type6 ? false : !mac_addr.IsZero();
    1818           0 :     Ip4Prefix inet_prefix;
    1819           0 :     Inet6Prefix inet6_prefix;
    1820           0 :     IpAddress ip_addr;
    1821           0 :     int prefix_len = 0;
    1822           0 :     if (!item.entry.nlri.address.empty()) {
    1823           0 :         size_t pos = item.entry.nlri.address.find('/');
    1824           0 :         if (pos == string::npos) {
    1825           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name, BGP_LOG_FLAG_ALL,
    1826             :                 "Missing / in address " << item.entry.nlri.address);
    1827           0 :             return false;
    1828             :         }
    1829             : 
    1830           0 :         bool ipv6 = item.entry.nlri.address.find(':') != string::npos;
    1831           0 :         if (!ipv6) {
    1832             :             inet_prefix =
    1833           0 :                 Ip4Prefix::FromString(item.entry.nlri.address, &error);
    1834           0 :             if (error) {
    1835           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1836             :                     BGP_LOG_FLAG_ALL,
    1837             :                     "Cannot parse inet prefix string " <<
    1838             :                         item.entry.nlri.address);
    1839           0 :                 return false;
    1840             :             }
    1841             : 
    1842           0 :             if (type2 && inet_prefix.prefixlen() != 32 &&
    1843           0 :                 item.entry.nlri.address != "0.0.0.0/0") {
    1844           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1845             :                     BGP_LOG_FLAG_ALL,
    1846             :                     "Bad inet address " << item.entry.nlri.address);
    1847           0 :                 return false;
    1848             :             }
    1849             : 
    1850           0 :             ip_addr = inet_prefix.ip4_addr();
    1851           0 :             prefix_len = inet_prefix.prefixlen();
    1852             :         } else {
    1853             :             inet6_prefix =
    1854           0 :                 Inet6Prefix::FromString(item.entry.nlri.address, &error);
    1855             : 
    1856           0 :             if (error) {
    1857           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1858             :                     BGP_LOG_FLAG_ALL,
    1859             :                     "Cannot parse inet6 prefix string " <<
    1860             :                         item.entry.nlri.address);
    1861           0 :                 return false;
    1862             :             }
    1863             : 
    1864           0 :             if (type2 && inet6_prefix.prefixlen() != 128 &&
    1865           0 :                 item.entry.nlri.address != "::/0") {
    1866           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1867             :                     BGP_LOG_FLAG_ALL,
    1868             :                     "Bad inet6 address " << item.entry.nlri.address);
    1869           0 :                 return false;
    1870             :             }
    1871           0 :             ip_addr = inet6_prefix.ip6_addr();
    1872           0 :             prefix_len = inet6_prefix.prefixlen();
    1873             :         }
    1874             :     }
    1875             : 
    1876             :     bool subscribe_pending;
    1877             :     int instance_id;
    1878             :     uint64_t subscription_gen_id;
    1879             :     BgpTable *table;
    1880           0 :     if (!VerifyMembership(vrf_name, Address::EVPN, &table, &instance_id,
    1881             :         &subscription_gen_id, &subscribe_pending, add_change)) {
    1882           0 :         channel_->Close();
    1883           0 :         return false;
    1884             :     }
    1885             : 
    1886           0 :     RouteDistinguisher rd;
    1887           0 :     if (mac_addr.IsBroadcast()) {
    1888           0 :         rd = RouteDistinguisher(peer_->bgp_identifier(), instance_id);
    1889           0 :     } else if (type6) {
    1890           0 :         rd = RouteDistinguisher(bgp_server_->bgp_identifier(),
    1891           0 :                                 table->routing_instance()->index());
    1892             :     } else {
    1893           0 :         rd = RouteDistinguisher::kZeroRd;
    1894             :     }
    1895             : 
    1896           0 :     uint32_t ethernet_tag = item.entry.nlri.ethernet_tag;
    1897             :     EvpnPrefix evpn_prefix = type6 ?
    1898             :         EvpnPrefix(rd, ethernet_tag, source, group,
    1899           0 :                    Ip4Address(bgp_server_->bgp_identifier())) :
    1900             :         type2 ? EvpnPrefix(rd, ethernet_tag, mac_addr, ip_addr) :
    1901           0 :         EvpnPrefix(rd, ip_addr, prefix_len);
    1902             : 
    1903           0 :     DBRequest req;
    1904           0 :     ExtCommunitySpec ext;
    1905           0 :     LargeCommunitySpec largecomm;
    1906           0 :     req.key.reset(new EvpnTable::RequestKey(evpn_prefix, peer_.get()));
    1907             : 
    1908           0 :     IpAddress nh_address(Ip4Address(0));
    1909           0 :     uint32_t label = 0;
    1910           0 :     uint32_t l3_label = 0;
    1911           0 :     uint32_t flags = 0;
    1912             : 
    1913           0 :     if (add_change) {
    1914           0 :         req.oper = DBRequest::DB_ENTRY_ADD_CHANGE;
    1915           0 :         BgpAttrSpec attrs;
    1916           0 :         const EnetNextHopListType &inh_list = item.entry.next_hops;
    1917             : 
    1918           0 :         if (inh_list.next_hop.empty()) {
    1919           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1920             :                 BGP_LOG_FLAG_ALL, "Missing next-hops for enet route " <<
    1921             :                                   evpn_prefix.ToXmppIdString());
    1922           0 :             return false;
    1923             :         }
    1924             : 
    1925             :         // Agents should send only one next-hop in the item.
    1926           0 :         if (inh_list.next_hop.size() != 1) {
    1927           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1928             :                 BGP_LOG_FLAG_ALL,
    1929             :                 "More than one nexthop received for enet route " <<
    1930             :                 evpn_prefix.ToXmppIdString());
    1931           0 :             return false;
    1932             :         }
    1933             : 
    1934           0 :         EnetNextHopListType::const_iterator nit = inh_list.begin();
    1935             : 
    1936           0 :         IpAddress nhop_address(Ip4Address(0));
    1937             : 
    1938           0 :         if (!XmppDecodeAddress(nit->af, nit->address, &nhop_address)) {
    1939           0 :             BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1940             :                 BGP_LOG_FLAG_ALL, "Bad nexthop address " << nit->address <<
    1941             :                 " for enet route " << evpn_prefix.ToXmppIdString());
    1942           0 :             return false;
    1943             :         }
    1944             : 
    1945           0 :         nh_address = nhop_address;
    1946           0 :         label = nit->label;
    1947           0 :         l3_label = nit->l3_label;
    1948           0 :         if (!nit->mac.empty()) {
    1949             :             MacAddress rmac_addr =
    1950           0 :                 MacAddress::FromString(nit->mac, &error);
    1951           0 :             if (error) {
    1952           0 :                 BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    1953             :                     BGP_LOG_FLAG_ALL,
    1954             :                     "Bad next-hop mac address " << nit->mac <<
    1955             :                     " for enet route " << evpn_prefix.ToXmppIdString());
    1956           0 :                 return false;
    1957             :             }
    1958           0 :             RouterMac router_mac(rmac_addr);
    1959           0 :             ext.communities.push_back(router_mac.GetExtCommunityValue());
    1960             :         }
    1961             : 
    1962             :         // Process tunnel encapsulation list.
    1963           0 :         bool no_tunnel_encap = true;
    1964           0 :         bool no_valid_tunnel_encap = true;
    1965           0 :         for (EnetTunnelEncapsulationListType::const_iterator eit =
    1966           0 :             nit->tunnel_encapsulation_list.begin();
    1967           0 :             eit != nit->tunnel_encapsulation_list.end(); ++eit) {
    1968           0 :             no_tunnel_encap = false;
    1969           0 :             TunnelEncap tun_encap(*eit);
    1970           0 :             if (tun_encap.tunnel_encap() == TunnelEncapType::UNSPEC)
    1971           0 :                 continue;
    1972           0 :             no_valid_tunnel_encap = false;
    1973           0 :             ext.communities.push_back(tun_encap.GetExtCommunityValue());
    1974           0 :             if (tun_encap.tunnel_encap() == TunnelEncapType::GRE) {
    1975           0 :                 TunnelEncap alt_tun_encap(TunnelEncapType::MPLS_O_GRE);
    1976           0 :                 ext.communities.push_back(alt_tun_encap.GetExtCommunityValue());
    1977             :             }
    1978             :         }
    1979             : 
    1980             :         // Mark the path as infeasible if all tunnel encaps published
    1981             :         // by agent are invalid.
    1982           0 :         if (!no_tunnel_encap && no_valid_tunnel_encap) {
    1983           0 :             flags = BgpPath::NoTunnelEncap;
    1984             :         }
    1985             : 
    1986             :         // Process tags list.
    1987           0 :         uint64_t tag_val = 0;
    1988           0 :         for (TagListType::const_iterator tit = nit->tag_list.begin();
    1989           0 :             tit != nit->tag_list.end(); ++tit) {
    1990           0 :             tag_val = nit->is_new_tags_list ? *tit :
    1991           0 :                 ((0x0000FFFF0000 & *tit) << 16) |
    1992           0 :                 (0x00000000FFFF & *tit);
    1993           0 :             TagLC tag_lc(bgp_server_->autonomous_system(), tag_val);
    1994           0 :             for (const auto &value_data : tag_lc.GetLargeCommunityValue()) {
    1995           0 :                 largecomm.communities.push_back(value_data);
    1996           0 :             }
    1997             :         }
    1998             : 
    1999             :         // Process local sequence_number
    2000           0 :         if (nit->local_sequence_number) {
    2001           0 :             LocalSequenceNumber lsn (nit->local_sequence_number);
    2002           0 :             ext.communities.push_back(lsn.GetExtCommunityValue());
    2003             :         }
    2004             : 
    2005           0 :         BgpAttrLocalPref local_pref(item.entry.local_preference);
    2006           0 :         if (local_pref.local_pref != 0) {
    2007           0 :             attrs.push_back(&local_pref);
    2008             :         }
    2009             : 
    2010             :         // If there's no explicit med, calculate it automatically from the
    2011             :         // local pref.
    2012           0 :         uint32_t med_value = item.entry.med;
    2013           0 :         if (!med_value)
    2014           0 :             med_value = GetMedFromLocalPref(local_pref.local_pref);
    2015           0 :         BgpAttrMultiExitDisc med(med_value);
    2016           0 :         if (med.med != 0)
    2017           0 :             attrs.push_back(&med);
    2018             : 
    2019           0 :         BgpAttrNextHop nexthop(nh_address.to_v4().to_ulong());
    2020           0 :         if (type6) {
    2021           0 :             flags |= BgpPath::CheckGlobalErmVpnRoute;
    2022           0 :             if (item.entry.replicator_address.empty() &&
    2023           0 :                     item.entry.edge_replication_not_supported) {
    2024             :                 // Only for test to inject remote smet routes
    2025           0 :                 flags &= ~BgpPath::CheckGlobalErmVpnRoute;
    2026           0 :                 attrs.push_back(&nexthop);
    2027             :             }
    2028             :         } else {
    2029           0 :             attrs.push_back(&nexthop);
    2030             :         }
    2031             : 
    2032           0 :         uint16_t cluster_seed = bgp_server_->global_config()->rd_cluster_seed();
    2033           0 :         BgpAttrSourceRd source_rd;
    2034           0 :         if (cluster_seed) {
    2035           0 :             source_rd = BgpAttrSourceRd(RouteDistinguisher(cluster_seed,
    2036           0 :                 nh_address.to_v4().to_ulong(), instance_id));
    2037             :         } else {
    2038           0 :             source_rd = BgpAttrSourceRd(RouteDistinguisher(
    2039           0 :                 nh_address.to_v4().to_ulong(), instance_id));
    2040             :         }
    2041           0 :         attrs.push_back(&source_rd);
    2042             : 
    2043             :         // Process security group list.
    2044           0 :         const EnetSecurityGroupListType &isg_list =
    2045             :             item.entry.security_group_list;
    2046           0 :         uint16_t sg_index = 0;
    2047           0 :         for (EnetSecurityGroupListType::const_iterator sit = isg_list.begin();
    2048           0 :             sit != isg_list.end(); ++sit) {
    2049           0 :             if (bgp_server_->autonomous_system() <= AS2_MAX) {
    2050           0 :                 SecurityGroup sg(bgp_server_->autonomous_system(), *sit);
    2051           0 :                 ext.communities.push_back(sg.GetExtCommunityValue());
    2052             :             } else {
    2053           0 :                 SecurityGroup sg(sg_index, *sit);
    2054           0 :                 SecurityGroup4ByteAs sg4(bgp_server_->autonomous_system(),
    2055           0 :                                              sg_index++);
    2056           0 :                 ext.communities.push_back(sg4.GetExtCommunityValue());
    2057           0 :                 ext.communities.push_back(sg.GetExtCommunityValue());
    2058             :             }
    2059             :         }
    2060             : 
    2061           0 :         if (item.entry.mobility.seqno) {
    2062           0 :             MacMobility mm(item.entry.mobility.seqno,
    2063           0 :                 item.entry.mobility.sticky);
    2064           0 :             ext.communities.push_back(mm.GetExtCommunityValue());
    2065           0 :         } else if (item.entry.sequence_number) {
    2066           0 :             MacMobility mm(item.entry.sequence_number);
    2067           0 :             ext.communities.push_back(mm.GetExtCommunityValue());
    2068             :         }
    2069             : 
    2070           0 :         ETree etree(item.entry.etree_leaf);
    2071           0 :         ext.communities.push_back(etree.GetExtCommunityValue());
    2072             : 
    2073           0 :         if (!ext.communities.empty())
    2074           0 :             attrs.push_back(&ext);
    2075           0 :         if (!largecomm.communities.empty())
    2076           0 :             attrs.push_back(&largecomm);
    2077             : 
    2078           0 :         PmsiTunnelSpec pmsi_spec;
    2079           0 :         if (mac_addr.IsBroadcast()) {
    2080           0 :             if (!item.entry.replicator_address.empty()) {
    2081           0 :                 IpAddress replicator_address;
    2082           0 :                 if (!XmppDecodeAddress(BgpAf::IPv4,
    2083             :                     item.entry.replicator_address, &replicator_address)) {
    2084           0 :                     BGP_LOG_PEER_INSTANCE_WARNING(Peer(), vrf_name,
    2085             :                         BGP_LOG_FLAG_ALL,
    2086             :                         "Bad replicator address " <<
    2087             :                         item.entry.replicator_address <<
    2088             :                         " for enet route " << evpn_prefix.ToXmppIdString());
    2089           0 :                     return false;
    2090             :                 }
    2091           0 :                 pmsi_spec.tunnel_type =
    2092             :                     PmsiTunnelSpec::AssistedReplicationContrail;
    2093           0 :                 pmsi_spec.tunnel_flags = PmsiTunnelSpec::ARLeaf;
    2094           0 :                 pmsi_spec.SetIdentifier(replicator_address.to_v4());
    2095             :             } else {
    2096           0 :                 pmsi_spec.tunnel_type = PmsiTunnelSpec::IngressReplication;
    2097           0 :                 if (item.entry.assisted_replication_supported) {
    2098           0 :                     pmsi_spec.tunnel_flags |= PmsiTunnelSpec::ARReplicator;
    2099           0 :                     pmsi_spec.tunnel_flags |= PmsiTunnelSpec::LeafInfoRequired;
    2100             :                 }
    2101           0 :                 if (!item.entry.edge_replication_not_supported) {
    2102           0 :                     pmsi_spec.tunnel_flags |=
    2103             :                         PmsiTunnelSpec::EdgeReplicationSupported;
    2104             :                 }
    2105           0 :                 pmsi_spec.SetIdentifier(nh_address.to_v4());
    2106             :             }
    2107           0 :             ExtCommunity ext_comm(bgp_server_->extcomm_db(), ext);
    2108           0 :             pmsi_spec.SetLabel(label, &ext_comm);
    2109           0 :             attrs.push_back(&pmsi_spec);
    2110           0 :         }
    2111             : 
    2112           0 :         BgpAttrPtr attr = bgp_server_->attr_db()->Locate(attrs);
    2113             : 
    2114           0 :         req.data.reset(new EvpnTable::RequestData(
    2115           0 :             attr, flags, label, l3_label, subscription_gen_id));
    2116           0 :         stats_[RX].reach++;
    2117           0 :     } else {
    2118           0 :         req.oper = DBRequest::DB_ENTRY_DELETE;
    2119           0 :         stats_[RX].unreach++;
    2120             :     }
    2121             : 
    2122             :     // Defer all requests till subscribe is processed.
    2123           0 :     if (subscribe_pending) {
    2124           0 :         DBRequest *request_entry = new DBRequest();
    2125           0 :         request_entry->Swap(&req);
    2126             :         string table_name =
    2127           0 :             RoutingInstance::GetTableName(vrf_name, Address::EVPN);
    2128           0 :         defer_q_.insert(make_pair(
    2129           0 :             make_pair(vrf_name, table_name), request_entry));
    2130           0 :         return true;
    2131           0 :     }
    2132             : 
    2133           0 :     assert(table);
    2134           0 :     BGP_LOG_PEER_INSTANCE(Peer(), vrf_name,
    2135             :         SandeshLevel::SYS_DEBUG, BGP_LOG_FLAG_TRACE,
    2136             :         "Enet route " << evpn_prefix.ToXmppIdString() <<
    2137             :         " with next-hop " << nh_address <<
    2138             :         " label " << label << " l3-label " << l3_label <<
    2139             :         " enqueued for " << (add_change ? "add/change" : "delete"));
    2140           0 :     table->Enqueue(&req);
    2141           0 :     return true;
    2142           0 : }
    2143             : 
    2144           0 : void BgpXmppChannel::DequeueRequest(const string &table_name,
    2145             :                                     DBRequest *request) {
    2146           0 :     unique_ptr<DBRequest> ptr(request);
    2147             : 
    2148             :     BgpTable *table = static_cast<BgpTable *>
    2149           0 :         (bgp_server_->database()->FindTable(table_name));
    2150           0 :     if (table == NULL || table->IsDeleted()) {
    2151           0 :         return;
    2152             :     }
    2153             : 
    2154           0 :     BgpMembershipManager *mgr = bgp_server_->membership_mgr();
    2155           0 :     if (mgr) {
    2156           0 :         int instance_id = -1;
    2157           0 :         uint64_t subscription_gen_id = 0;
    2158           0 :         bool is_registered = mgr->GetRegistrationInfo(peer_.get(), table,
    2159             :                                             &instance_id, &subscription_gen_id);
    2160           0 :         if (!is_registered) {
    2161           0 :             BGP_LOG_PEER_WARNING(Membership, Peer(),
    2162             :                 BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2163             :                 "Not subscribed to table " << table->name());
    2164           0 :             return;
    2165             :         }
    2166           0 :         if (ptr->oper == DBRequest::DB_ENTRY_ADD_CHANGE) {
    2167           0 :             ((BgpTable::RequestData *)ptr->data.get())
    2168           0 :                 ->set_subscription_gen_id(subscription_gen_id);
    2169             :         }
    2170             :     }
    2171             : 
    2172           0 :     table->Enqueue(ptr.get());
    2173           0 : }
    2174             : 
    2175           0 : bool BgpXmppChannel::ResumeClose() {
    2176           0 :     peer_->Close(true);
    2177           0 :     return true;
    2178             : }
    2179             : 
    2180           0 : void BgpXmppChannel::RegisterTable(int line, BgpTable *table,
    2181             :     const TableMembershipRequestState *tmr_state) {
    2182             :     // Defer if Membership manager is in use (by close manager).
    2183           0 :     if (close_manager_->IsMembershipInUse()) {
    2184           0 :         BGP_LOG_PEER_TABLE(Peer(), SandeshLevel::SYS_DEBUG,
    2185             :                            BGP_LOG_FLAG_ALL, table, "RegisterTable deferred "
    2186             :                            "from :" << line);
    2187           0 :         return;
    2188             :     }
    2189             : 
    2190           0 :     BgpMembershipManager *mgr = bgp_server_->membership_mgr();
    2191           0 :     BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2192             :                  BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2193             :                  "Subscribe to table " << table->name() <<
    2194             :                  (tmr_state->no_ribout ? " (no ribout)" : "") <<
    2195             :                  " with id " << tmr_state->instance_id);
    2196           0 :     if (tmr_state->no_ribout) {
    2197           0 :         mgr->RegisterRibIn(peer_.get(), table);
    2198           0 :         mgr->SetRegistrationInfo(peer_.get(), table, tmr_state->instance_id,
    2199           0 :             manager_->get_subscription_gen_id());
    2200           0 :         channel_stats_.table_subscribe++;
    2201           0 :         MembershipRequestCallback(table);
    2202             :     } else {
    2203           0 :         mgr->Register(peer_.get(), table, bgp_policy_, tmr_state->instance_id);
    2204           0 :         channel_stats_.table_subscribe++;
    2205             :     }
    2206             : 
    2207             :     // If EndOfRib Send timer is running, cancel it and reschedule it after all
    2208             :     // outstanding membership registrations are complete.
    2209           0 :     if (eor_send_timer_->running())
    2210           0 :         eor_send_timer_->Cancel();
    2211             : }
    2212             : 
    2213           0 : void BgpXmppChannel::UnregisterTable(int line, BgpTable *table) {
    2214             :     // Defer if Membership manager is in use (by close manager).
    2215           0 :     if (close_manager_->IsMembershipInUse()) {
    2216           0 :         BGP_LOG_PEER_TABLE(Peer(), SandeshLevel::SYS_DEBUG,
    2217             :                            BGP_LOG_FLAG_ALL, table, "UnregisterTable deferred "
    2218             :                            "from :" << line);
    2219           0 :         return;
    2220             :     }
    2221             : 
    2222           0 :     BgpMembershipManager *mgr = bgp_server_->membership_mgr();
    2223           0 :     BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2224             :                  BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2225             :                  "Unsubscribe to table " << table->name());
    2226           0 :     mgr->Unregister(peer_.get(), table);
    2227           0 :     channel_stats_.table_unsubscribe++;
    2228             : }
    2229             : 
    2230             : #define RegisterTable(table, tmr_state) \
    2231             :     RegisterTable(__LINE__, table, tmr_state)
    2232             : #define UnregisterTable(table) UnregisterTable(__LINE__, table)
    2233             : 
    2234             : // Process all pending membership requests of various tables.
    2235           0 : void BgpXmppChannel::ProcessPendingSubscriptions() {
    2236           0 :     assert(!close_manager_->IsMembershipInUse());
    2237           0 :     BOOST_FOREACH(TableMembershipRequestMap::value_type &entry,
    2238             :                   table_membership_request_map_) {
    2239             :         BgpTable *table = static_cast<BgpTable *>(
    2240           0 :             bgp_server_->database()->FindTable(entry.first));
    2241           0 :         const TableMembershipRequestState &tmr_state = entry.second;
    2242           0 :         if (tmr_state.current_req == SUBSCRIBE) {
    2243           0 :             RegisterTable(table, &tmr_state);
    2244             :         } else {
    2245           0 :             assert(tmr_state.current_req == UNSUBSCRIBE);
    2246           0 :             UnregisterTable(table);
    2247             :         }
    2248             :     }
    2249           0 : }
    2250             : 
    2251           0 : size_t BgpXmppChannel::table_membership_requests() const {
    2252           0 :     return table_membership_request_map_.size();
    2253             : }
    2254             : 
    2255           0 : bool BgpXmppChannel::MembershipResponseHandler(string table_name) {
    2256           0 :     if (close_manager_->IsMembershipInUse()) {
    2257           0 :         close_manager_->MembershipRequestCallback();
    2258           0 :         return true;
    2259             :     }
    2260             : 
    2261             :     TableMembershipRequestState *tmr_state =
    2262           0 :         GetTableMembershipState(table_name);
    2263           0 :     if (!tmr_state) {
    2264           0 :         BGP_LOG_PEER_INSTANCE_CRITICAL(Peer(), table_name,
    2265             :                      BGP_PEER_DIR_IN, BGP_LOG_FLAG_ALL,
    2266             :                      "Table not in subscribe/unsubscribe request queue");
    2267           0 :         assert(false);
    2268             :     }
    2269             : 
    2270           0 :     if (tmr_state->current_req == SUBSCRIBE) {
    2271           0 :         BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2272             :                      BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2273             :                      "Subscribe to table " << table_name << " completed");
    2274           0 :         channel_stats_.table_subscribe_complete++;
    2275             :     } else {
    2276           0 :         BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2277             :                      BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2278             :                      "Unsubscribe to table " << table_name << " completed");
    2279           0 :         channel_stats_.table_unsubscribe_complete++;
    2280             :     }
    2281             : 
    2282           0 :     if (defer_peer_close_) {
    2283           0 :         DeleteTableMembershipState(table_name);
    2284           0 :         if (table_membership_requests())
    2285           0 :             return true;
    2286           0 :         defer_peer_close_ = false;
    2287           0 :         ResumeClose();
    2288             :     } else {
    2289           0 :         ProcessMembershipResponse(table_name, tmr_state);
    2290             :     }
    2291             : 
    2292           0 :     assert(channel_stats_.table_subscribe_complete <=
    2293             :                channel_stats_.table_subscribe);
    2294           0 :     assert(channel_stats_.table_unsubscribe_complete <=
    2295             :                channel_stats_.table_unsubscribe);
    2296             : 
    2297             :     // Restart EndOfRib send if necessary.
    2298           0 :     ResetEndOfRibSendState();
    2299             : 
    2300             :     // If Close manager is waiting to use membership, try now.
    2301           0 :     if (close_manager_->IsMembershipInWait())
    2302           0 :         close_manager_->MembershipRequest();
    2303             : 
    2304           0 :     return true;
    2305             : }
    2306             : 
    2307           0 : bool BgpXmppChannel::ProcessMembershipResponse(string table_name,
    2308             :     TableMembershipRequestState *tmr_state) {
    2309             :     BgpTable *table = static_cast<BgpTable *>
    2310           0 :         (bgp_server_->database()->FindTable(table_name));
    2311           0 :     if (!table) {
    2312           0 :         DeleteTableMembershipState(table_name);
    2313           0 :         return true;
    2314             :     }
    2315           0 :     BgpMembershipManager *mgr = bgp_server_->membership_mgr();
    2316             : 
    2317           0 :     if ((tmr_state->current_req == UNSUBSCRIBE) &&
    2318           0 :         (tmr_state->pending_req == SUBSCRIBE)) {
    2319             :         // Process pending subscribe now that unsubscribe has completed.
    2320           0 :         tmr_state->current_req = SUBSCRIBE;
    2321           0 :         RegisterTable(table, tmr_state);
    2322           0 :         return true;
    2323           0 :     } else if ((tmr_state->current_req == SUBSCRIBE) &&
    2324           0 :                (tmr_state->pending_req == UNSUBSCRIBE)) {
    2325             :         // Process pending unsubscribe now that subscribe has completed.
    2326           0 :         tmr_state->current_req = UNSUBSCRIBE;
    2327           0 :         UnregisterTable(table);
    2328           0 :         return true;
    2329           0 :     } else if ((tmr_state->current_req == SUBSCRIBE) &&
    2330           0 :         (tmr_state->pending_req == SUBSCRIBE) &&
    2331           0 :         (mgr->IsRibOutRegistered(peer_.get(), table) == tmr_state->no_ribout)) {
    2332             :         // Trigger an unsubscribe so that we can subsequently subscribe with
    2333             :         // the updated value of no_ribout.
    2334           0 :         tmr_state->current_req = UNSUBSCRIBE;
    2335           0 :         UnregisterTable(table);
    2336           0 :         return true;
    2337             :     }
    2338             : 
    2339           0 :     string vrf_name = table->routing_instance()->name();
    2340           0 :     VrfTableName vrf_n_table = make_pair(vrf_name, table->name());
    2341             : 
    2342           0 :     if (tmr_state->pending_req == UNSUBSCRIBE) {
    2343           0 :         if (!GetInstanceMembershipState(vrf_name))
    2344           0 :             assert(defer_q_.count(vrf_n_table) == 0);
    2345           0 :         DeleteTableMembershipState(table_name);
    2346           0 :         return true;
    2347           0 :     } else if (tmr_state->pending_req == SUBSCRIBE) {
    2348           0 :         mgr->SetRegistrationInfo(peer_.get(), table, tmr_state->instance_id,
    2349           0 :             manager_->get_subscription_gen_id());
    2350           0 :         DeleteTableMembershipState(table_name);
    2351             :     }
    2352             : 
    2353           0 :     for (DeferQ::iterator it = defer_q_.find(vrf_n_table);
    2354           0 :          it != defer_q_.end() && it->first.second == table->name(); ++it) {
    2355           0 :         DequeueRequest(table->name(), it->second);
    2356             :     }
    2357             : 
    2358             :     // Erase all elements for the table
    2359           0 :     defer_q_.erase(vrf_n_table);
    2360             : 
    2361           0 :     return true;
    2362           0 : }
    2363             : 
    2364           0 : void BgpXmppChannel::MembershipRequestCallback(BgpTable *table) {
    2365           0 :     membership_response_worker_.Enqueue(table->name());
    2366           0 : }
    2367             : 
    2368           0 : void BgpXmppChannel::FillCloseInfo(BgpNeighborResp *resp) const {
    2369           0 :     close_manager_->FillCloseInfo(resp);
    2370           0 : }
    2371             : 
    2372           0 : void BgpXmppChannel::FillInstanceMembershipInfo(BgpNeighborResp *resp) const {
    2373           0 :     vector<BgpNeighborRoutingInstance> instance_list;
    2374           0 :     BOOST_FOREACH(const SubscribedRoutingInstanceList::value_type &entry,
    2375             :         routing_instances_) {
    2376           0 :         BgpNeighborRoutingInstance instance;
    2377           0 :         instance.set_name(entry.first->name());
    2378           0 :         if (entry.second.IsLlgrStale()) {
    2379           0 :             instance.set_state("subscribed-llgr-stale");
    2380           0 :         } else if (entry.second.IsGrStale()) {
    2381           0 :             instance.set_state("subscribed-gr-stale");
    2382             :         } else {
    2383           0 :             instance.set_state("subscribed");
    2384             :         }
    2385           0 :         instance.set_index(entry.second.index);
    2386           0 :         rtarget_manager_->FillInfo(&instance, entry.second.targets);
    2387           0 :         instance_list.push_back(instance);
    2388           0 :     }
    2389           0 :     BOOST_FOREACH(const InstanceMembershipRequestMap::value_type &entry,
    2390             :         instance_membership_request_map_) {
    2391           0 :         const InstanceMembershipRequestState &imr_state = entry.second;
    2392           0 :         BgpNeighborRoutingInstance instance;
    2393           0 :         instance.set_name(entry.first);
    2394           0 :         instance.set_state("pending");
    2395           0 :         instance.set_index(imr_state.instance_id);
    2396           0 :         instance_list.push_back(instance);
    2397           0 :     }
    2398           0 :     resp->set_routing_instances(instance_list);
    2399           0 : }
    2400             : 
    2401           0 : void BgpXmppChannel::FillTableMembershipInfo(BgpNeighborResp *resp) const {
    2402           0 :     vector<BgpNeighborRoutingTable> old_table_list = resp->get_routing_tables();
    2403           0 :     set<string> old_table_set;
    2404           0 :     vector<BgpNeighborRoutingTable> new_table_list;
    2405             : 
    2406           0 :     BOOST_FOREACH(const BgpNeighborRoutingTable &table, old_table_list) {
    2407           0 :         old_table_set.insert(table.get_name());
    2408           0 :         if (!GetTableMembershipState(table.get_name()))
    2409           0 :             new_table_list.push_back(table);
    2410             :     }
    2411             : 
    2412           0 :     BOOST_FOREACH(const TableMembershipRequestMap::value_type &entry,
    2413             :         table_membership_request_map_) {
    2414           0 :         BgpNeighborRoutingTable table;
    2415           0 :         table.set_name(entry.first);
    2416           0 :         if (old_table_set.find(entry.first) != old_table_set.end())
    2417           0 :             table.set_current_state("subscribed");
    2418           0 :         const TableMembershipRequestState &tmr_state = entry.second;
    2419           0 :         if (tmr_state.current_req == SUBSCRIBE) {
    2420           0 :             table.set_current_request("subscribe");
    2421             :         } else {
    2422           0 :             table.set_current_request("unsubscribe");
    2423             :         }
    2424           0 :         if (tmr_state.pending_req == SUBSCRIBE) {
    2425           0 :             table.set_pending_request("subscribe");
    2426             :         } else {
    2427           0 :             table.set_pending_request("unsubscribe");
    2428             :         }
    2429           0 :         new_table_list.push_back(table);
    2430           0 :     }
    2431           0 :     resp->set_routing_tables(new_table_list);
    2432           0 : }
    2433             : 
    2434             : //
    2435             : // Erase all defer_q_ elements with the given (vrf, table).
    2436             : //
    2437           0 : void BgpXmppChannel::FlushDeferQ(string vrf_name, string table_name) {
    2438           0 :     for (DeferQ::iterator it =
    2439           0 :         defer_q_.find(make_pair(vrf_name, table_name)), itnext;
    2440           0 :         (it != defer_q_.end() && it->first.second == table_name);
    2441           0 :         it = itnext) {
    2442           0 :         itnext = it;
    2443           0 :         itnext++;
    2444           0 :         delete it->second;
    2445           0 :         defer_q_.erase(it);
    2446             :     }
    2447           0 : }
    2448             : 
    2449             : //
    2450             : // Erase all defer_q_ elements for all tables for the given vrf.
    2451             : //
    2452           0 : void BgpXmppChannel::FlushDeferQ(string vrf_name) {
    2453           0 :     for (DeferQ::iterator it =
    2454           0 :         defer_q_.lower_bound(make_pair(vrf_name, string())), itnext;
    2455           0 :         (it != defer_q_.end() && it->first.first == vrf_name);
    2456           0 :         it = itnext) {
    2457           0 :         itnext = it;
    2458           0 :         itnext++;
    2459           0 :         delete it->second;
    2460           0 :         defer_q_.erase(it);
    2461             :     }
    2462           0 : }
    2463             : 
    2464             : // Mark all current subscriptions as 'stale'. This is called when peer close
    2465             : // process is initiated by BgpXmppChannel via PeerCloseManager.
    2466           0 : void BgpXmppChannel::StaleCurrentSubscriptions() {
    2467           0 :     CHECK_CONCURRENCY(peer_close_->GetTaskName());
    2468           0 :     BOOST_FOREACH(SubscribedRoutingInstanceList::value_type &entry,
    2469             :                   routing_instances_) {
    2470           0 :         entry.second.SetGrStale();
    2471           0 :         rtarget_manager_->UpdateRouteTargetRouteFlag(entry.first,
    2472           0 :                 entry.second.targets, BgpPath::Stale);
    2473             :     }
    2474           0 : }
    2475             : 
    2476             : // Mark all current subscriptions as 'llgr_stale'.
    2477           0 : void BgpXmppChannel::LlgrStaleCurrentSubscriptions() {
    2478           0 :     CHECK_CONCURRENCY(peer_close_->GetTaskName());
    2479           0 :     BOOST_FOREACH(SubscribedRoutingInstanceList::value_type &entry,
    2480             :                   routing_instances_) {
    2481           0 :         assert(entry.second.IsGrStale());
    2482           0 :         entry.second.SetLlgrStale();
    2483           0 :         rtarget_manager_->UpdateRouteTargetRouteFlag(entry.first,
    2484           0 :                 entry.second.targets, BgpPath::Stale | BgpPath::LlgrStale);
    2485             :     }
    2486           0 : }
    2487             : 
    2488             : // Sweep all current subscriptions which are still marked as 'stale'.
    2489           0 : void BgpXmppChannel::SweepCurrentSubscriptions() {
    2490           0 :     CHECK_CONCURRENCY(peer_close_->GetTaskName());
    2491           0 :     for (SubscribedRoutingInstanceList::iterator i = routing_instances_.begin();
    2492           0 :             i != routing_instances_.end();) {
    2493           0 :         if (i->second.IsGrStale()) {
    2494           0 :             string name = i->first->name();
    2495             : 
    2496             :             // Increment the iterator first as we expect the entry to be
    2497             :             // soon removed.
    2498           0 :             i++;
    2499           0 :             BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2500             :                          BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2501             :                          "Instance subscription " << name <<
    2502             :                          " is still stale and hence unsubscribed");
    2503           0 :             ProcessSubscriptionRequest(name, NULL, false);
    2504           0 :         } else {
    2505           0 :             i++;
    2506             :         }
    2507             :     }
    2508           0 : }
    2509             : 
    2510             : // Clear staled subscription state as new subscription has been received.
    2511           0 : void BgpXmppChannel::ClearStaledSubscription(RoutingInstance *rt_instance,
    2512             :         SubscriptionState *sub_state) {
    2513           0 :     if (!sub_state->IsGrStale())
    2514           0 :         return;
    2515             : 
    2516           0 :     BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2517             :                  BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2518             :                  "Instance subscription " << rt_instance->name() <<
    2519             :                  " stale flag is cleared");
    2520           0 :     sub_state->ClearStale();
    2521           0 :     rtarget_manager_->Stale(sub_state->targets);
    2522             : }
    2523             : 
    2524           0 : void BgpXmppChannel::AddSubscriptionState(RoutingInstance *rt_instance,
    2525             :         int index) {
    2526           0 :     SubscriptionState state(rt_instance->GetImportList(), index);
    2527             :     pair<SubscribedRoutingInstanceList::iterator, bool> ret =
    2528           0 :         routing_instances_.insert(pair<RoutingInstance *, SubscriptionState> (
    2529             :                                       rt_instance, state));
    2530             : 
    2531             :     // During GR, we expect duplicate subscription requests. Clear stale
    2532             :     // state, as agent did re-subscribe after restart.
    2533           0 :     if (!ret.second) {
    2534           0 :         ClearStaledSubscription(rt_instance, &ret.first->second);
    2535             :     } else {
    2536           0 :         rtarget_manager_->PublishRTargetRoute(rt_instance, true);
    2537             :     }
    2538           0 : }
    2539             : 
    2540           0 : void BgpXmppChannel::DeleteSubscriptionState(RoutingInstance *rt_instance) {
    2541           0 :     routing_instances_.erase(rt_instance);
    2542           0 : }
    2543             : 
    2544           0 : BgpXmppChannel::SubscriptionState *BgpXmppChannel::GetSubscriptionState(
    2545             :     RoutingInstance *rt_instance) {
    2546             :     SubscribedRoutingInstanceList::iterator loc =
    2547           0 :         routing_instances_.find(rt_instance);
    2548           0 :     return (loc != routing_instances_.end() ? &loc->second : NULL);
    2549             : }
    2550             : 
    2551           0 : const BgpXmppChannel::SubscriptionState *BgpXmppChannel::GetSubscriptionState(
    2552             :     RoutingInstance *rt_instance) const {
    2553             :     SubscribedRoutingInstanceList::const_iterator loc =
    2554           0 :         routing_instances_.find(rt_instance);
    2555           0 :     return (loc != routing_instances_.end() ? &loc->second : NULL);
    2556             : }
    2557             : 
    2558           0 : void BgpXmppChannel::ProcessDeferredSubscribeRequest(RoutingInstance *instance,
    2559             :     const InstanceMembershipRequestState &imr_state) {
    2560           0 :     int instance_id = imr_state.instance_id;
    2561           0 :     bool no_ribout = imr_state.no_ribout;
    2562           0 :     AddSubscriptionState(instance, instance_id);
    2563           0 :     RoutingInstance::RouteTableList const rt_list = instance->GetTables();
    2564           0 :     for (RoutingInstance::RouteTableList::const_iterator it = rt_list.begin();
    2565           0 :          it != rt_list.end(); ++it) {
    2566           0 :         BgpTable *table = it->second;
    2567           0 :         if (table->IsVpnTable() || table->family() == Address::RTARGET)
    2568           0 :             continue;
    2569             : 
    2570             :         TableMembershipRequestState tmr_state(
    2571           0 :             SUBSCRIBE, instance_id, no_ribout);
    2572           0 :         AddTableMembershipState(table->name(), tmr_state);
    2573           0 :         RegisterTable(table, &tmr_state);
    2574             :     }
    2575           0 : }
    2576             : 
    2577           0 : void BgpXmppChannel::ProcessSubscriptionRequest(
    2578             :         string vrf_name, const XmppStanza::XmppMessageIq *iq,
    2579             :         bool add_change) {
    2580           0 :     int instance_id = -1;
    2581           0 :     bool no_ribout = false;
    2582             : 
    2583           0 :     if (add_change) {
    2584           0 :         XmlPugi *pugi = reinterpret_cast<XmlPugi *>(iq->dom.get());
    2585           0 :         xml_node options = pugi->FindNode("options");
    2586           0 :         for (xml_node node = options.first_child(); node;
    2587           0 :              node = node.next_sibling()) {
    2588           0 :             if (strcmp(node.name(), "instance-id") == 0) {
    2589           0 :                 instance_id = node.text().as_int();
    2590             :             }
    2591           0 :             if (strcmp(node.name(), "no-ribout") == 0) {
    2592           0 :                 no_ribout = node.text().as_bool();
    2593             :             }
    2594             :         }
    2595             :     }
    2596             : 
    2597           0 :     RoutingInstanceMgr *instance_mgr = bgp_server_->routing_instance_mgr();
    2598           0 :     assert(instance_mgr);
    2599           0 :     RoutingInstance *rt_instance = instance_mgr->GetRoutingInstance(vrf_name);
    2600           0 :     if (rt_instance == NULL) {
    2601           0 :         BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_INFO,
    2602             :                      BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2603             :                      "Routing instance " << vrf_name <<
    2604             :                      " not found when processing " <<
    2605             :                      (add_change ? "subscribe" : "unsubscribe"));
    2606           0 :         if (add_change) {
    2607           0 :             if (GetInstanceMembershipState(vrf_name)) {
    2608           0 :                 BGP_LOG_PEER_WARNING(Membership, Peer(),
    2609             :                              BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2610             :                              "Duplicate subscribe for routing instance " <<
    2611             :                              vrf_name << ", triggering close");
    2612           0 :                 channel_->Close();
    2613             :             } else {
    2614           0 :                 AddInstanceMembershipState(vrf_name,
    2615             :                     InstanceMembershipRequestState(instance_id, no_ribout));
    2616           0 :                 channel_stats_.instance_subscribe++;
    2617             :             }
    2618             :         } else {
    2619           0 :             if (DeleteInstanceMembershipState(vrf_name)) {
    2620           0 :                 FlushDeferQ(vrf_name);
    2621           0 :                 channel_stats_.instance_unsubscribe++;
    2622             :             } else {
    2623           0 :                 BGP_LOG_PEER_WARNING(Membership, Peer(),
    2624             :                              BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2625             :                              "Spurious unsubscribe for routing instance " <<
    2626             :                              vrf_name << ", triggering close");
    2627           0 :                 channel_->Close();
    2628             :             }
    2629             :         }
    2630           0 :         return;
    2631           0 :     } else if (rt_instance->deleted()) {
    2632           0 :         BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2633             :                      BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2634             :                      "Routing instance " << vrf_name <<
    2635             :                      " is being deleted when processing " <<
    2636             :                      (add_change ? "subscribe" : "unsubscribe"));
    2637           0 :         if (add_change) {
    2638           0 :             if (GetInstanceMembershipState(vrf_name)) {
    2639           0 :                 BGP_LOG_PEER_WARNING(Membership, Peer(),
    2640             :                              BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2641             :                              "Duplicate subscribe for routing instance " <<
    2642             :                              vrf_name << ", triggering close");
    2643           0 :                 channel_->Close();
    2644           0 :             } else if (GetSubscriptionState(rt_instance)) {
    2645           0 :                 BGP_LOG_PEER_WARNING(Membership, Peer(),
    2646             :                              BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2647             :                              "Duplicate subscribe for routing instance " <<
    2648             :                              vrf_name << ", triggering close");
    2649           0 :                 channel_->Close();
    2650             :             } else {
    2651           0 :                 AddInstanceMembershipState(vrf_name,
    2652             :                     InstanceMembershipRequestState(instance_id, no_ribout));
    2653           0 :                 channel_stats_.instance_subscribe++;
    2654             :             }
    2655           0 :             return;
    2656             :         } else {
    2657             :             // If instance is being deleted and agent is trying to unsubscribe
    2658             :             // we need to process the unsubscribe if vrf is not in the request
    2659             :             // map.  This would be the normal case where we wait for agent to
    2660             :             // unsubscribe in order to remove routes added by it.
    2661           0 :             if (DeleteInstanceMembershipState(vrf_name)) {
    2662           0 :                 FlushDeferQ(vrf_name);
    2663           0 :                 channel_stats_.instance_unsubscribe++;
    2664           0 :                 return;
    2665           0 :             } else if (!GetSubscriptionState(rt_instance)) {
    2666           0 :                 BGP_LOG_PEER_WARNING(Membership, Peer(),
    2667             :                              BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2668             :                              "Spurious unsubscribe for routing instance " <<
    2669             :                              vrf_name << ", triggering close");
    2670           0 :                 channel_->Close();
    2671           0 :                 return;
    2672             :             }
    2673           0 :             channel_stats_.instance_unsubscribe++;
    2674             :         }
    2675             :     } else {
    2676           0 :         if (add_change) {
    2677             :             const SubscriptionState *sub_state =
    2678           0 :                 GetSubscriptionState(rt_instance);
    2679           0 :             if (sub_state) {
    2680           0 :                 if (!close_manager_->IsCloseInProgress()) {
    2681           0 :                     BGP_LOG_PEER_WARNING(Membership, Peer(),
    2682             :                                  BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2683             :                                  "Duplicate subscribe for routing instance " <<
    2684             :                                  vrf_name << ", triggering close");
    2685           0 :                     channel_->Close();
    2686           0 :                     return;
    2687             :                 }
    2688           0 :                 if (!sub_state->IsGrStale()) {
    2689           0 :                     BGP_LOG_PEER_WARNING(Membership, Peer(),
    2690             :                                  BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2691             :                                  "Duplicate subscribe for routing instance " <<
    2692             :                                  vrf_name << " under GR, triggering close");
    2693           0 :                     channel_->Close();
    2694           0 :                     return;
    2695             :                 }
    2696             :             }
    2697           0 :             channel_stats_.instance_subscribe++;
    2698             :         } else {
    2699           0 :             if (!GetSubscriptionState(rt_instance)) {
    2700           0 :                 BGP_LOG_PEER_WARNING(Membership, Peer(),
    2701             :                              BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2702             :                              "Spurious unsubscribe for routing instance " <<
    2703             :                              vrf_name << ", triggering close");
    2704           0 :                 channel_->Close();
    2705           0 :                 return;
    2706             :             }
    2707           0 :             channel_stats_.instance_unsubscribe++;
    2708             :         }
    2709             :     }
    2710             : 
    2711           0 :     if (add_change) {
    2712           0 :         AddSubscriptionState(rt_instance, instance_id);
    2713             :     } else  {
    2714           0 :         rtarget_manager_->PublishRTargetRoute(rt_instance, false);
    2715           0 :         DeleteSubscriptionState(rt_instance);
    2716             :     }
    2717             : 
    2718           0 :     RoutingInstance::RouteTableList const rt_list = rt_instance->GetTables();
    2719           0 :     for (RoutingInstance::RouteTableList::const_iterator it = rt_list.begin();
    2720           0 :          it != rt_list.end(); ++it) {
    2721           0 :         BgpTable *table = it->second;
    2722           0 :         if (table->IsVpnTable() || table->family() == Address::RTARGET)
    2723           0 :             continue;
    2724             : 
    2725           0 :         if (add_change) {
    2726             :             TableMembershipRequestState *tmr_state =
    2727           0 :                 GetTableMembershipState(table->name());
    2728           0 :             if (!tmr_state) {
    2729             :                 TableMembershipRequestState tmp_tmr_state(
    2730           0 :                     SUBSCRIBE, instance_id, no_ribout);
    2731           0 :                 AddTableMembershipState(table->name(), tmp_tmr_state);
    2732           0 :                 RegisterTable(table, &tmp_tmr_state);
    2733             :             } else {
    2734           0 :                 tmr_state->instance_id = instance_id;
    2735           0 :                 tmr_state->pending_req = SUBSCRIBE;
    2736           0 :                 tmr_state->no_ribout = no_ribout;
    2737             :             }
    2738             :         } else {
    2739           0 :             if (defer_q_.count(make_pair(vrf_name, table->name()))) {
    2740           0 :                 BGP_LOG_PEER(Membership, Peer(), SandeshLevel::SYS_DEBUG,
    2741             :                              BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2742             :                              "Flush deferred route requests for table " <<
    2743             :                              table->name() << " on unsubscribe");
    2744           0 :                 FlushDeferQ(vrf_name, table->name());
    2745             :             }
    2746             : 
    2747             :             // Erase all elements for the table.
    2748             : 
    2749             :             TableMembershipRequestState *tmr_state =
    2750           0 :                 GetTableMembershipState(table->name());
    2751           0 :             if (!tmr_state) {
    2752           0 :                 AddTableMembershipState(table->name(),
    2753             :                     TableMembershipRequestState(
    2754             :                         UNSUBSCRIBE, instance_id, no_ribout));
    2755           0 :                 UnregisterTable(table);
    2756             :             } else {
    2757           0 :                 tmr_state->instance_id = -1;
    2758           0 :                 tmr_state->pending_req = UNSUBSCRIBE;
    2759           0 :                 tmr_state->no_ribout = false;
    2760             :             }
    2761             :         }
    2762             :     }
    2763           0 : }
    2764             : 
    2765           0 : void BgpXmppChannel::ClearEndOfRibState() {
    2766           0 :     eor_receive_timer_->Cancel();
    2767           0 :     eor_send_timer_->Cancel();
    2768           0 :     eor_sent_ = false;
    2769           0 : }
    2770             : 
    2771           0 : void BgpXmppChannel::ReceiveEndOfRIB(Address::Family family) {
    2772           0 :     eor_receive_timer_->Cancel();
    2773           0 :     close_manager_->ProcessEORMarkerReceived(family);
    2774           0 : }
    2775             : 
    2776           0 : void BgpXmppChannel::EndOfRibTimerErrorHandler(string error_name,
    2777             :                                                string error_message) {
    2778           0 :     BGP_LOG_PEER_CRITICAL(Timer, Peer(), BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA,
    2779             :                  "Timer error: " << error_name << " " << error_message);
    2780           0 : }
    2781             : 
    2782           0 : bool BgpXmppChannel::EndOfRibReceiveTimerExpired() {
    2783           0 :     if (!peer_->IsReady())
    2784           0 :         return false;
    2785             : 
    2786           0 :     uint32_t timeout = manager() && manager()->xmpp_server() ?
    2787           0 :         manager()->xmpp_server()->GetEndOfRibReceiveTime() :
    2788           0 :         BgpGlobalSystemConfig::kEndOfRibTime;
    2789             : 
    2790             :     // If max timeout has not reached yet, check if we can exit GR sooner by
    2791             :     // looking at the activity in the channel.
    2792           0 :     if (UTCTimestamp() - eor_receive_timer_start_time_ < timeout) {
    2793             : 
    2794             :         // If there is some send or receive activity in the channel in last few
    2795             :         // seconds, delay EoR receive event.
    2796           0 :         if (channel_->LastReceived(kEndOfRibSendRetryTime * 6) ||
    2797           0 :                 channel_->LastSent(kEndOfRibSendRetryTime * 6)) {
    2798           0 :             eor_receive_timer_->Reschedule(kEndOfRibSendRetryTime * 1000);
    2799           0 :             BGP_LOG_PEER(Message, Peer(), SandeshLevel::SYS_INFO,
    2800             :                          BGP_LOG_FLAG_ALL, BGP_PEER_DIR_IN,
    2801             :                          "EndOfRib Receive timer rescheduled to fire after " <<
    2802             :                          kEndOfRibSendRetryTime << " second(s)");
    2803           0 :             return true;
    2804             :         }
    2805             :     }
    2806             : 
    2807           0 :     BGP_LOG_PEER(Message, Peer(), SandeshLevel::SYS_INFO, BGP_LOG_FLAG_ALL,
    2808             :                  BGP_PEER_DIR_IN, "EndOfRib Receive timer expired");
    2809           0 :     ReceiveEndOfRIB(Address::UNSPEC);
    2810           0 :     return false;
    2811             : }
    2812             : 
    2813           0 : time_t BgpXmppChannel::GetEndOfRibSendTime() const {
    2814           0 :     return manager() && manager()->xmpp_server() ?
    2815           0 :         manager()->xmpp_server()->GetEndOfRibSendTime() :
    2816           0 :         BgpGlobalSystemConfig::kEndOfRibTime;
    2817             : }
    2818             : 
    2819           0 : bool BgpXmppChannel::EndOfRibSendTimerExpired() {
    2820           0 :     if (!peer_->IsReady())
    2821           0 :         return false;
    2822             : 
    2823             :     // If max timeout has not reached yet, check if we can exit GR sooner by
    2824             :     // looking at the activity in the channel.
    2825           0 :     if (UTCTimestamp() - eor_send_timer_start_time_ < GetEndOfRibSendTime()) {
    2826             : 
    2827             :         // If there is some send or receive activity in the channel in last few
    2828             :         // seconds, delay EoR send event.
    2829           0 :         if (channel_->LastReceived(kEndOfRibSendRetryTime * 6) ||
    2830           0 :                 channel_->LastSent(kEndOfRibSendRetryTime * 6) ||
    2831           0 :                 manager()->bgp_server()->IsServerStartingUp()) {
    2832           0 :             eor_send_timer_->Reschedule(kEndOfRibSendRetryTime * 1000);
    2833           0 :             BGP_LOG_PEER(Message, Peer(), SandeshLevel::SYS_INFO,
    2834             :                          BGP_LOG_FLAG_ALL, BGP_PEER_DIR_OUT,
    2835             :                          "EndOfRib Send timer rescheduled to fire after " <<
    2836             :                          kEndOfRibSendRetryTime << " second(s)");
    2837           0 :             return true;
    2838             :         }
    2839             :     }
    2840             : 
    2841           0 :     SendEndOfRIB();
    2842           0 :     return false;
    2843             : }
    2844             : 
    2845           0 : void BgpXmppChannel::StartEndOfRibReceiveTimer() {
    2846           0 :     uint32_t timeout = manager() && manager()->xmpp_server() ?
    2847           0 :                            manager()->xmpp_server()->GetEndOfRibReceiveTime() :
    2848           0 :                            BgpGlobalSystemConfig::kEndOfRibTime;
    2849           0 :     eor_receive_timer_start_time_ = UTCTimestamp();
    2850           0 :     eor_receive_timer_->Cancel();
    2851             : 
    2852           0 :     BGP_LOG_PEER(Message, Peer(), SandeshLevel::SYS_INFO, BGP_LOG_FLAG_ALL,
    2853             :         BGP_PEER_DIR_IN, "EndOfRib Receive timer scheduled to fire after " <<
    2854             :         timeout << " second(s)");
    2855           0 :     eor_receive_timer_->Start(timeout * 1000,
    2856             :         boost::bind(&BgpXmppChannel::EndOfRibReceiveTimerExpired, this),
    2857             :         boost::bind(&BgpXmppChannel::EndOfRibTimerErrorHandler, this, _1, _2));
    2858           0 : }
    2859             : 
    2860           0 : void BgpXmppChannel::ResetEndOfRibSendState() {
    2861           0 :     if (eor_sent_)
    2862           0 :         return;
    2863             : 
    2864             :     // If socket is blocked, then wait for it to get unblocked first.
    2865           0 :     if (!peer_->send_ready())
    2866           0 :         return;
    2867             : 
    2868             :     // If there is any outstanding subscribe pending, wait for its completion.
    2869           0 :     if (channel_stats_.table_subscribe_complete !=
    2870           0 :             channel_stats_.table_subscribe)
    2871           0 :         return;
    2872             : 
    2873           0 :     eor_send_timer_start_time_ = UTCTimestamp();
    2874           0 :     eor_send_timer_->Cancel();
    2875             : 
    2876           0 :     BGP_LOG_PEER(Message, Peer(), SandeshLevel::SYS_INFO, BGP_LOG_FLAG_ALL,
    2877             :         BGP_PEER_DIR_OUT, "EndOfRib Send timer scheduled to fire after " <<
    2878             :         kEndOfRibSendRetryTime << " second(s)");
    2879           0 :     eor_send_timer_->Start(kEndOfRibSendRetryTime * 1000,
    2880             :         boost::bind(&BgpXmppChannel::EndOfRibSendTimerExpired, this),
    2881             :         boost::bind(&BgpXmppChannel::EndOfRibTimerErrorHandler, this, _1, _2));
    2882             : }
    2883             : 
    2884             : /*
    2885             :  * Empty items list constitute eor marker.
    2886             :  */
    2887           0 : void BgpXmppChannel::SendEndOfRIB() {
    2888           0 :     eor_send_timer_->Cancel();
    2889           0 :     eor_sent_ = true;
    2890             : 
    2891           0 :     string msg;
    2892           0 :     msg += "\n<message from=\"";
    2893           0 :     msg += XmppInit::kControlNodeJID;
    2894           0 :     msg += "\" to=\"";
    2895           0 :     msg += peer_->ToString();
    2896           0 :     msg += "/";
    2897           0 :     msg += XmppInit::kBgpPeer;
    2898           0 :     msg += "\">";
    2899           0 :     msg += "\n\t<event xmlns=\"http://jabber.org/protocol/pubsub\">";
    2900           0 :     msg = (msg + "\n<items node=\"") + XmppInit::kEndOfRibMarker +
    2901           0 :           "\"></items>";
    2902           0 :     msg += "\n\t</event>\n</message>\n";
    2903             : 
    2904           0 :     if (channel_->connection())
    2905           0 :         channel_->connection()->Send((const uint8_t *) msg.data(), msg.size());
    2906             : 
    2907           0 :     stats_[TX].end_of_rib++;
    2908           0 :     BGP_LOG_PEER(Message, Peer(), SandeshLevel::SYS_INFO, BGP_LOG_FLAG_ALL,
    2909             :                  BGP_PEER_DIR_OUT, "EndOfRib marker sent");
    2910           0 : }
    2911             : 
    2912             : // Process any associated primary instance-id.
    2913           0 : int BgpXmppChannel::GetPrimaryInstanceID(const string &s,
    2914             :                                          bool expect_prefix_len) const {
    2915           0 :     if (s.empty())
    2916           0 :         return 0;
    2917           0 :     char *str = const_cast<char *>(s.c_str());
    2918             :     char *saveptr, *token;
    2919           0 :     token = strtok_r(str, "/", &saveptr); // Get afi
    2920           0 :     if (!token || !saveptr)
    2921           0 :         return 0;
    2922           0 :     token = strtok_r(NULL, "/", &saveptr); // Get safi
    2923           0 :     if (!token || !saveptr)
    2924           0 :         return 0;
    2925           0 :     token = strtok_r(NULL, "/", &saveptr); // vrf name
    2926           0 :     if (!token || !saveptr)
    2927           0 :         return 0;
    2928           0 :     token = strtok_r(NULL, "/", &saveptr); // address
    2929           0 :     if (!token || !saveptr)
    2930           0 :         return 0;
    2931           0 :     if (expect_prefix_len) {
    2932           0 :         token = strtok_r(NULL, "/", &saveptr); // prefix-length
    2933           0 :         if (!token || !saveptr)
    2934           0 :             return 0;
    2935             :     }
    2936           0 :     token = strtok_r(NULL, "/", &saveptr); // primary instance-id
    2937           0 :     if (!token)
    2938           0 :         return 0;
    2939           0 :     return strtoul(token, NULL, 0);
    2940             : }
    2941             : 
    2942           0 : void BgpXmppChannel::ReceiveUpdate(const XmppStanza::XmppMessage *msg) {
    2943           0 :     CHECK_CONCURRENCY("xmpp::StateMachine");
    2944             : 
    2945             :     // Bail if the connection is being deleted. It's not safe to assert
    2946             :     // because the Delete method can be called from the main thread.
    2947           0 :     if (channel_->connection() && channel_->connection()->IsDeleted())
    2948           0 :         return;
    2949             : 
    2950             :     // Make sure that peer is not set for closure already.
    2951           0 :     assert(!defer_peer_close_);
    2952           0 :     assert(!peer_deleted());
    2953             : 
    2954           0 :     if (msg->type == XmppStanza::IQ_STANZA) {
    2955           0 :         const XmppStanza::XmppMessageIq *iq =
    2956             :                    static_cast<const XmppStanza::XmppMessageIq *>(msg);
    2957           0 :         if (iq->iq_type.compare("set") == 0) {
    2958           0 :             if (iq->action.compare("subscribe") == 0) {
    2959           0 :                 ProcessSubscriptionRequest(iq->node, iq, true);
    2960           0 :             } else if (iq->action.compare("unsubscribe") == 0) {
    2961           0 :                 ProcessSubscriptionRequest(iq->node, iq, false);
    2962           0 :             } else if (iq->action.compare("publish") == 0) {
    2963           0 :                 XmlBase *impl = msg->dom.get();
    2964           0 :                 stats_[RX].rt_updates++;
    2965           0 :                 XmlPugi *pugi = reinterpret_cast<XmlPugi *>(impl);
    2966           0 :                 xml_node item = pugi->FindNode("item");
    2967             : 
    2968             :                 // Empty items-list can be considered as EOR Marker for all afis
    2969           0 :                 if (item == 0) {
    2970           0 :                     BGP_LOG_PEER(Message, Peer(), SandeshLevel::SYS_INFO,
    2971             :                                  BGP_LOG_FLAG_ALL, BGP_PEER_DIR_IN,
    2972             :                                  "EndOfRib marker received");
    2973           0 :                     stats_[RX].end_of_rib++;
    2974           0 :                     ReceiveEndOfRIB(Address::UNSPEC);
    2975           0 :                     return;
    2976             :                 }
    2977           0 :                 for (; item; item = item.next_sibling()) {
    2978           0 :                     if (strcmp(item.name(), "item") != 0) continue;
    2979             : 
    2980           0 :                     string id(iq->as_node.c_str());
    2981           0 :                     char *str = const_cast<char *>(id.c_str());
    2982             :                     char *saveptr;
    2983           0 :                     char *af = strtok_r(str, "/", &saveptr);
    2984           0 :                     char *safi = strtok_r(NULL, "/", &saveptr);
    2985             : 
    2986           0 :                     if (atoi(af) == BgpAf::IPv4 &&
    2987           0 :                         ((atoi(safi) == BgpAf::Unicast) ||
    2988           0 :                          (atoi(safi) == BgpAf::Mpls))) {
    2989           0 :                         ProcessItem(iq->node, item, iq->is_as_node,
    2990           0 :                             GetPrimaryInstanceID(iq->as_node, true));
    2991           0 :                     } else if (atoi(af) == BgpAf::IPv6 &&
    2992           0 :                                atoi(safi) == BgpAf::Unicast) {
    2993           0 :                         ProcessInet6Item(iq->node, item, iq->is_as_node);
    2994           0 :                     } else if (atoi(af) == BgpAf::IPv4 &&
    2995           0 :                         atoi(safi) == BgpAf::Mcast) {
    2996           0 :                         ProcessMcastItem(iq->node, item, iq->is_as_node);
    2997           0 :                     } else if (atoi(af) == BgpAf::IPv4 &&
    2998           0 :                         atoi(safi) == BgpAf::MVpn) {
    2999           0 :                         ProcessMvpnItem(iq->node, item, iq->is_as_node);
    3000           0 :                     } else if (atoi(af) == BgpAf::L2Vpn &&
    3001           0 :                                atoi(safi) == BgpAf::Enet) {
    3002           0 :                         ProcessEnetItem(iq->node, item, iq->is_as_node);
    3003             :                     }
    3004           0 :                 }
    3005             :             }
    3006             :         }
    3007             :     }
    3008             : }
    3009             : 
    3010           0 : bool BgpXmppChannelManager::DeleteChannel(BgpXmppChannel *channel) {
    3011           0 :     if (!channel->deleted()) {
    3012           0 :         channel->set_deleted(true);
    3013           0 :         delete channel;
    3014             :     }
    3015           0 :     return true;
    3016             : }
    3017             : 
    3018             : // BgpXmppChannelManager routines.
    3019           0 : BgpXmppChannelManager::BgpXmppChannelManager(XmppServer *xmpp_server,
    3020           0 :                                              BgpServer *server)
    3021           0 :     : xmpp_server_(xmpp_server),
    3022           0 :       bgp_server_(server),
    3023           0 :       queue_(TaskScheduler::GetInstance()->GetTaskId("bgp::Config"), 0,
    3024             :           boost::bind(&BgpXmppChannelManager::DeleteChannel, this, _1)),
    3025           0 :       id_(-1),
    3026           0 :       asn_listener_id_(-1),
    3027           0 :       identifier_listener_id_(-1),
    3028           0 :       dscp_listener_id_(-1) {
    3029             :     // Initialize the gen id counter
    3030           0 :     subscription_gen_id_ = 1;
    3031           0 :     deleting_count_ = 0;
    3032             : 
    3033           0 :     if (xmpp_server)
    3034           0 :         xmpp_server->CreateConfigUpdater(server->config_manager());
    3035           0 :     queue_.SetEntryCallback(
    3036             :             boost::bind(&BgpXmppChannelManager::IsReadyForDeletion, this));
    3037           0 :     if (xmpp_server) {
    3038           0 :         xmpp_server->RegisterConnectionEvent(xmps::BGP,
    3039             :                boost::bind(&BgpXmppChannelManager::XmppHandleChannelEvent,
    3040             :                            this, _1, _2));
    3041             :     }
    3042           0 :     admin_down_listener_id_ =
    3043           0 :         server->RegisterAdminDownCallback(boost::bind(
    3044             :             &BgpXmppChannelManager::AdminDownCallback, this));
    3045           0 :     asn_listener_id_ =
    3046           0 :         server->RegisterASNUpdateCallback(boost::bind(
    3047             :             &BgpXmppChannelManager::ASNUpdateCallback, this, _1, _2));
    3048           0 :     identifier_listener_id_ =
    3049           0 :         server->RegisterIdentifierUpdateCallback(boost::bind(
    3050             :             &BgpXmppChannelManager::IdentifierUpdateCallback, this, _1));
    3051           0 :     dscp_listener_id_ =
    3052           0 :         server->RegisterDSCPUpdateCallback(boost::bind(
    3053             :             &BgpXmppChannelManager::DSCPUpdateCallback, this, _1));
    3054             : 
    3055           0 :     id_ = server->routing_instance_mgr()->RegisterInstanceOpCallback(
    3056             :         boost::bind(&BgpXmppChannelManager::RoutingInstanceCallback,
    3057             :                     this, _1, _2));
    3058           0 : }
    3059             : 
    3060           0 : BgpXmppChannelManager::~BgpXmppChannelManager() {
    3061           0 :     assert(channel_map_.empty());
    3062           0 :     assert(channel_name_map_.empty());
    3063           0 :     assert(deleting_count_ == 0);
    3064           0 :     if (xmpp_server_) {
    3065           0 :         xmpp_server_->UnRegisterConnectionEvent(xmps::BGP);
    3066             :     }
    3067             : 
    3068           0 :     queue_.Shutdown();
    3069           0 :     bgp_server_->UnregisterAdminDownCallback(admin_down_listener_id_);
    3070           0 :     bgp_server_->UnregisterASNUpdateCallback(asn_listener_id_);
    3071           0 :     bgp_server_->routing_instance_mgr()->UnregisterInstanceOpCallback(id_);
    3072           0 :     bgp_server_->UnregisterDSCPUpdateCallback(dscp_listener_id_);
    3073           0 : }
    3074             : 
    3075           0 : bool BgpXmppChannelManager::IsReadyForDeletion() {
    3076           0 :     return bgp_server_->IsReadyForDeletion();
    3077             : }
    3078             : 
    3079           0 : void BgpXmppChannelManager::SetQueueDisable(bool disabled) {
    3080           0 :     queue_.set_disable(disabled);
    3081           0 : }
    3082             : 
    3083           0 : size_t BgpXmppChannelManager::GetQueueSize() const {
    3084           0 :     return queue_.Length();
    3085             : }
    3086             : 
    3087           0 : void BgpXmppChannelManager::AdminDownCallback() {
    3088           0 :     xmpp_server_->ClearAllConnections();
    3089           0 : }
    3090             : 
    3091           0 : void BgpXmppChannelManager::DSCPUpdateCallback(uint8_t dscp_value) {
    3092           0 :     xmpp_server_->SetDscpValue(dscp_value);
    3093           0 : }
    3094             : 
    3095           0 : void BgpXmppChannelManager::ASNUpdateCallback(as_t old_asn,
    3096             :     as_t old_local_asn) {
    3097           0 :     CHECK_CONCURRENCY("bgp::Config");
    3098           0 :     BOOST_FOREACH(XmppChannelMap::value_type &i, channel_map_) {
    3099           0 :         i.second->rtarget_manager()->ASNUpdateCallback(old_asn, old_local_asn);
    3100             :     }
    3101           0 :     if (bgp_server_->autonomous_system() != old_asn) {
    3102           0 :         xmpp_server_->ClearAllConnections();
    3103             :     }
    3104           0 : }
    3105             : 
    3106           0 : void BgpXmppChannelManager::IdentifierUpdateCallback(
    3107             :         Ip4Address old_identifier) {
    3108           0 :     CHECK_CONCURRENCY("bgp::Config");
    3109           0 :     xmpp_server_->ClearAllConnections();
    3110           0 : }
    3111             : 
    3112           0 : void BgpXmppChannelManager::RoutingInstanceCallback(string vrf_name, int op) {
    3113           0 :     CHECK_CONCURRENCY("bgp::Config", "bgp::ConfigHelper");
    3114           0 :     BOOST_FOREACH(XmppChannelMap::value_type &i, channel_map_) {
    3115           0 :         i.second->RoutingInstanceCallback(vrf_name, op);
    3116             :     }
    3117           0 : }
    3118             : 
    3119           0 : void BgpXmppChannelManager::VisitChannels(BgpXmppChannelManager::VisitorFn fn) {
    3120           0 :     std::scoped_lock lock(mutex_);
    3121           0 :     BOOST_FOREACH(XmppChannelMap::value_type &i, channel_map_) {
    3122           0 :         fn(i.second);
    3123             :     }
    3124           0 : }
    3125             : 
    3126           0 : void BgpXmppChannelManager::VisitChannels(BgpXmppChannelManager::VisitorFn fn)
    3127             :         const {
    3128           0 :     std::scoped_lock lock(mutex_);
    3129           0 :     BOOST_FOREACH(const XmppChannelMap::value_type &i, channel_map_) {
    3130           0 :         fn(i.second);
    3131             :     }
    3132           0 : }
    3133             : 
    3134           0 : BgpXmppChannel *BgpXmppChannelManager::FindChannel(string client) {
    3135           0 :     BOOST_FOREACH(XmppChannelMap::value_type &i, channel_map_) {
    3136           0 :         if (i.second->ToString() == client) {
    3137           0 :             return i.second;
    3138             :         }
    3139             :     }
    3140           0 :     return NULL;
    3141             : }
    3142             : 
    3143           0 : BgpXmppChannel *BgpXmppChannelManager::FindChannel(
    3144             :         const XmppChannel *ch) {
    3145           0 :     XmppChannelMap::iterator it = channel_map_.find(ch);
    3146           0 :     if (it == channel_map_.end())
    3147           0 :         return NULL;
    3148           0 :     return it->second;
    3149             : }
    3150             : 
    3151           0 : void BgpXmppChannelManager::RemoveChannel(XmppChannel *channel) {
    3152           0 :     if (channel->connection() && !channel->connection()->IsActiveChannel()) {
    3153           0 :         CHECK_CONCURRENCY("bgp::Config");
    3154             :     }
    3155           0 :     channel_map_.erase(channel);
    3156           0 :     channel_name_map_.erase(channel->ToString());
    3157           0 : }
    3158             : 
    3159           0 : BgpXmppChannel *BgpXmppChannelManager::CreateChannel(XmppChannel *channel) {
    3160           0 :     CHECK_CONCURRENCY("xmpp::StateMachine");
    3161           0 :     BgpXmppChannel *ch = new BgpXmppChannel(channel, bgp_server_, this);
    3162             : 
    3163           0 :     return ch;
    3164             : }
    3165             : 
    3166           0 : void BgpXmppChannelManager::XmppHandleChannelEvent(XmppChannel *channel,
    3167             :                                                    xmps::PeerState state) {
    3168           0 :     std::scoped_lock lock(mutex_);
    3169             : 
    3170           0 :     XmppChannelMap::iterator it = channel_map_.find(channel);
    3171           0 :     BgpXmppChannel *bgp_xmpp_channel = NULL;
    3172           0 :     if (state == xmps::READY) {
    3173           0 :         if (it == channel_map_.end()) {
    3174           0 :             bgp_xmpp_channel = CreateChannel(channel);
    3175           0 :             channel_map_.insert(make_pair(channel, bgp_xmpp_channel));
    3176           0 :             channel_name_map_.insert(
    3177           0 :                 make_pair(channel->ToString(), bgp_xmpp_channel));
    3178           0 :             BGP_LOG_PEER(Message, bgp_xmpp_channel->Peer(),
    3179             :                          Sandesh::LoggingUtLevel(), BGP_LOG_FLAG_SYSLOG,
    3180             :                          BGP_PEER_DIR_IN,
    3181             :                          "Received XmppChannel up event");
    3182           0 :             if (!bgp_server_->HasSelfConfiguration()) {
    3183           0 :                 BGP_LOG_PEER(Message, bgp_xmpp_channel->Peer(),
    3184             :                              SandeshLevel::SYS_INFO, BGP_LOG_FLAG_SYSLOG,
    3185             :                              BGP_PEER_DIR_IN,
    3186             :                              "No BGP configuration for self - closing channel");
    3187           0 :                 if (!getenv("CONTRAIL_CAT_FRAMEWORK"))
    3188           0 :                     channel->Close();
    3189             :             }
    3190           0 :             if (bgp_server_->admin_down()) {
    3191           0 :                 BGP_LOG_PEER(Message, bgp_xmpp_channel->Peer(),
    3192             :                              SandeshLevel::SYS_INFO, BGP_LOG_FLAG_SYSLOG,
    3193             :                              BGP_PEER_DIR_IN,
    3194             :                              "BGP is administratively down - closing channel");
    3195           0 :                 channel->Close();
    3196             :             }
    3197             :         } else {
    3198           0 :             bgp_xmpp_channel = (*it).second;
    3199           0 :             if (bgp_xmpp_channel->peer_deleted())
    3200           0 :                 return;
    3201             : 
    3202             :             // Gracefully close the channel if GR closure is in progress.
    3203             :             // This can happen if GR timers fire just after session comes
    3204             :             // back up.
    3205           0 :             if (bgp_xmpp_channel->close_manager()->IsCloseInProgress() &&
    3206           0 :                 !bgp_xmpp_channel->close_manager()->IsInGRTimerWaitState()) {
    3207           0 :                 BGP_LOG_PEER(Message, bgp_xmpp_channel->Peer(),
    3208             :                              SandeshLevel::SYS_INFO, BGP_LOG_FLAG_SYSLOG,
    3209             :                              BGP_PEER_DIR_IN,
    3210             :                              "Graceful Closure in progress - Closing channel");
    3211           0 :                 channel->Close();
    3212             :             }
    3213           0 :             channel->RegisterReceive(xmps::BGP,
    3214             :                 boost::bind(&BgpXmppChannel::ReceiveUpdate, bgp_xmpp_channel,
    3215             :                             _1));
    3216             :         }
    3217             : 
    3218           0 :         bgp_xmpp_channel->eor_sent_ = false;
    3219           0 :         bgp_xmpp_channel->StartEndOfRibReceiveTimer();
    3220           0 :         bgp_xmpp_channel->ResetEndOfRibSendState();
    3221           0 :     } else if (state == xmps::NOT_READY) {
    3222           0 :         if (it != channel_map_.end()) {
    3223           0 :             bgp_xmpp_channel = (*it).second;
    3224           0 :             BGP_LOG_PEER(Message, bgp_xmpp_channel->Peer(),
    3225             :                          Sandesh::LoggingUtLevel(), BGP_LOG_FLAG_SYSLOG,
    3226             :                          BGP_PEER_DIR_IN,
    3227             :                          "Received XmppChannel down event");
    3228             : 
    3229             :             // Trigger closure of this channel
    3230           0 :             bgp_xmpp_channel->Close();
    3231             :         } else {
    3232           0 :             ostringstream os;
    3233           0 :             os << "Peer not found for " << channel->ToString() <<
    3234           0 :                   " on channel down event";
    3235           0 :             BGP_LOG_NOTICE(BgpMessage, BGP_LOG_FLAG_ALL, os.str());
    3236           0 :         }
    3237             :     }
    3238           0 : }
    3239             : 
    3240           0 : void BgpXmppChannelManager::FillPeerInfo(const BgpXmppChannel *channel) const {
    3241           0 :     PeerStatsInfo stats;
    3242           0 :     PeerStats::FillPeerDebugStats(channel->Peer()->peer_stats(), &stats);
    3243             : 
    3244           0 :     XmppPeerInfoData peer_info;
    3245           0 :     peer_info.set_name(channel->Peer()->ToUVEKey());
    3246           0 :     peer_info.set_peer_stats_info(stats);
    3247           0 :     assert(!peer_info.get_name().empty());
    3248           0 :     BGP_UVE_SEND(XMPPPeerInfo, peer_info);
    3249             : 
    3250           0 :     PeerStatsData peer_stats_data;
    3251           0 :     peer_stats_data.set_name(channel->Peer()->ToUVEKey());
    3252           0 :     peer_stats_data.set_encoding("XMPP");
    3253           0 :     PeerStats::FillPeerUpdateStats(channel->Peer()->peer_stats(),
    3254             :                                    &peer_stats_data);
    3255           0 :     assert(!peer_stats_data.get_name().empty());
    3256           0 :     BGP_UVE_SEND2(PeerStatsUve, peer_stats_data, "ObjectXmppPeerInfo");
    3257           0 : }
    3258             : 
    3259           0 : bool BgpXmppChannelManager::CollectStats(BgpRouterState *state, bool first)
    3260             :          const {
    3261           0 :     CHECK_CONCURRENCY("bgp::ShowCommand");
    3262             : 
    3263           0 :     VisitChannels(boost::bind(&BgpXmppChannelManager::FillPeerInfo, this, _1));
    3264           0 :     bool change = false;
    3265           0 :     uint32_t num_xmpp = count();
    3266           0 :     if (first || num_xmpp != state->get_num_xmpp_peer()) {
    3267           0 :         state->set_num_xmpp_peer(num_xmpp);
    3268           0 :         change = true;
    3269             :     }
    3270             : 
    3271           0 :     uint32_t num_up_xmpp = NumUpPeer();
    3272           0 :     if (first || num_up_xmpp != state->get_num_up_xmpp_peer()) {
    3273           0 :         state->set_num_up_xmpp_peer(num_up_xmpp);
    3274           0 :         change = true;
    3275             :     }
    3276             : 
    3277           0 :     uint32_t num_deleting_xmpp = deleting_count();
    3278           0 :     if (first || num_deleting_xmpp != state->get_num_deleting_xmpp_peer()) {
    3279           0 :         state->set_num_deleting_xmpp_peer(num_deleting_xmpp);
    3280           0 :         change = true;
    3281             :     }
    3282             : 
    3283           0 :     return change;
    3284             : }
    3285             : 
    3286           0 : void BgpXmppChannel::Close() {
    3287           0 :     instance_membership_request_map_.clear();
    3288           0 :     STLDeleteElements(&defer_q_);
    3289             : 
    3290           0 :     if (table_membership_requests()) {
    3291           0 :         BGP_LOG_PEER(Event, peer_.get(), SandeshLevel::SYS_INFO,
    3292             :             BGP_LOG_FLAG_ALL, BGP_PEER_DIR_NA, "Close procedure deferred");
    3293           0 :         defer_peer_close_ = true;
    3294           0 :         return;
    3295             :     }
    3296           0 :     peer_->Close(true);
    3297             : }
    3298             : 
    3299             : //
    3300             : // Return connection's remote tcp endpoint if available
    3301             : //
    3302           0 : TcpSession::Endpoint BgpXmppChannel::remote_endpoint() const {
    3303           0 :     const XmppSession *session = GetSession();
    3304           0 :     if (session) {
    3305           0 :         return session->remote_endpoint();
    3306             :     }
    3307           0 :     return TcpSession::Endpoint();
    3308             : }
    3309             : 
    3310             : //
    3311             : // Return connection's local tcp endpoint if available
    3312             : //
    3313           0 : TcpSession::Endpoint BgpXmppChannel::local_endpoint() const {
    3314           0 :     const XmppSession *session = GetSession();
    3315           0 :     if (session) {
    3316           0 :         return session->local_endpoint();
    3317             :     }
    3318           0 :     return TcpSession::Endpoint();
    3319             : }
    3320             : 
    3321             : //
    3322             : // Return connection's remote tcp endpoint string.
    3323             : //
    3324           0 : string BgpXmppChannel::transport_address_string() const {
    3325           0 :     TcpSession::Endpoint endpoint = remote_endpoint();
    3326           0 :     ostringstream oss;
    3327           0 :     oss << endpoint;
    3328           0 :     return oss.str();
    3329           0 : }
    3330             : 
    3331             : //
    3332             : // Mark the XmppPeer as deleted.
    3333             : //
    3334           0 : void BgpXmppChannel::set_peer_closed(bool flag) {
    3335           0 :     peer_->SetPeerClosed(flag);
    3336           0 : }
    3337             : 
    3338             : //
    3339             : // Return true if the XmppPeer is deleted.
    3340             : //
    3341           0 : bool BgpXmppChannel::peer_deleted() const {
    3342           0 :     return peer_->IsDeleted();
    3343             : }
    3344             : 
    3345             : //
    3346             : // Return time stamp of when the XmppPeer delete was initiated.
    3347             : //
    3348           0 : uint64_t BgpXmppChannel::peer_closed_at() const {
    3349           0 :     return peer_->closed_at();
    3350             : }
    3351             : 
    3352           0 : bool BgpXmppChannel::IsSubscriptionGrStale(RoutingInstance *instance) const {
    3353             :     SubscribedRoutingInstanceList::const_iterator it =
    3354           0 :         routing_instances_.find(instance);
    3355           0 :     assert(it != routing_instances_.end());
    3356           0 :     return it->second.IsGrStale();
    3357             : }
    3358             : 
    3359           0 : bool BgpXmppChannel::IsSubscriptionLlgrStale(RoutingInstance *instance) const {
    3360             :     SubscribedRoutingInstanceList::const_iterator it =
    3361           0 :         routing_instances_.find(instance);
    3362           0 :     assert(it != routing_instances_.end());
    3363           0 :     return it->second.IsLlgrStale();
    3364             : }
    3365             : 
    3366           0 : bool BgpXmppChannel::IsSubscriptionEmpty() const {
    3367           0 :     return routing_instances_.empty();
    3368             : }
    3369             : 
    3370           0 : const RoutingInstance::RouteTargetList &BgpXmppChannel::GetSubscribedRTargets(
    3371             :         RoutingInstance *instance) const {
    3372             :     SubscribedRoutingInstanceList::const_iterator it =
    3373           0 :         routing_instances_.find(instance);
    3374           0 :     assert(it != routing_instances_.end());
    3375           0 :     return it->second.targets;
    3376             : }

Generated by: LCOV version 1.14