+public class BGPPeer implements BGPSessionListener, Peer, BGPPeerRuntimeMXBean, TransactionChainListener {
+
+ private static final Logger LOG = LoggerFactory.getLogger(BGPPeer.class);
+
+ @GuardedBy("this")
+ private final Set<TablesKey> tables = new HashSet<>();
+ @GuardedBy("this")
+ private BGPSession session;
+ @GuardedBy("this")
+ private byte[] rawIdentifier;
+ @GuardedBy("this")
+ private DOMTransactionChain chain;
+ @GuardedBy("this")
+ private AdjRibInWriter ribWriter;
+ @GuardedBy("this")
+ private EffectiveRibInWriter effRibInWriter;
+
+ private final RIB rib;
+ private final String name;
+ private BGPPeerRuntimeRegistrator registrator;
+ private BGPPeerRuntimeRegistration runtimeReg;
+ private final Map<TablesKey, AdjRibOutListener> adjRibOutListenerSet = new HashMap();
+ private final RpcProviderRegistry rpcRegistry;
+ private RoutedRpcRegistration<BgpPeerRpcService> rpcRegistration;
+ private final PeerRole peerRole;
+ private final Optional<SimpleRoutingPolicy> simpleRoutingPolicy;
+ private final BGPPeerStats peerStats;
+ private YangInstanceIdentifier peerIId;
+ private final Set<AbstractRegistration> tableRegistration = new HashSet<>();
+
+ public BGPPeer(final String name, final RIB rib, final PeerRole role, final SimpleRoutingPolicy peerStatus, final RpcProviderRegistry rpcRegistry) {
+ this.peerRole = role;
+ this.simpleRoutingPolicy = Optional.ofNullable(peerStatus);
+ this.rib = Preconditions.checkNotNull(rib);
+ this.name = name;
+ this.rpcRegistry = rpcRegistry;
+ this.peerStats = new BGPPeerStatsImpl(this.name, this.tables);
+ this.chain = rib.createPeerChain(this);
+ }
+
+ public BGPPeer(final String name, final RIB rib, final PeerRole role, final RpcProviderRegistry rpcRegistry) {
+ this(name, rib, role, null, rpcRegistry);
+ }
+
+ public void instantiateServiceInstance() {
+ // add current peer to "configured BGP peer" stats
+ this.rib.getRenderStats().getConfiguredPeerCounter().increaseCount();
+ this.ribWriter = AdjRibInWriter.create(this.rib.getYangRibId(), this.peerRole, this.simpleRoutingPolicy, this.chain);
+ }
+
+ // FIXME ListenableFuture<?> should be used once closeServiceInstance uses wildcard too
+ @Override
+ public synchronized ListenableFuture<Void> close() {
+ final ListenableFuture<Void> future = releaseConnection();
+ this.chain.close();
+ return future;
+ }
+
+ @Override
+ public void onMessage(final BGPSession session, final Notification msg) throws BGPDocumentedException {
+ if (!(msg instanceof Update) && !(msg instanceof RouteRefresh)) {
+ LOG.info("Ignoring unhandled message class {}", msg.getClass());
+ return;
+ }
+ if (msg instanceof Update) {
+ onUpdateMessage((Update) msg);
+ } else {
+ onRouteRefreshMessage((RouteRefresh) msg, session);
+ }
+ }
+
+ private void onRouteRefreshMessage(final RouteRefresh message, final BGPSession session) {
+ final Class<? extends AddressFamily> rrAfi = message.getAfi();
+ final Class<? extends SubsequentAddressFamily> rrSafi = message.getSafi();
+
+ final TablesKey key = new TablesKey(rrAfi, rrSafi);
+ final AdjRibOutListener listener = this.adjRibOutListenerSet.get(key);
+ if (listener != null) {
+ listener.close();
+ this.adjRibOutListenerSet.remove(key);
+ createAdjRibOutListener(RouterIds.createPeerId(session.getBgpId()), key, listener.isMpSupported());
+ } else {
+ LOG.info("Ignoring RouteRefresh message. Afi/Safi is not supported: {}, {}.", rrAfi, rrSafi);
+ }
+ }
+
+ /**
+ * Check for presence of well known mandatory attribute LOCAL_PREF in Update message
+ *
+ * @param message Update message
+ * @throws BGPDocumentedException
+ */
+ private void checkMandatoryAttributesPresence(final Update message) throws BGPDocumentedException {
+ if (MessageUtil.isAnyNlriPresent(message)) {
+ final Attributes attrs = message.getAttributes();
+ if (this.peerRole == PeerRole.Ibgp && (attrs == null || attrs.getLocalPref() == null)) {
+ throw new BGPDocumentedException(BGPError.MANDATORY_ATTR_MISSING_MSG + "LOCAL_PREF",
+ BGPError.WELL_KNOWN_ATTR_MISSING,
+ new byte[] { LocalPreferenceAttributeParser.TYPE });
+ }
+ }
+ }
+
+ /**
+ * Process Update message received.
+ * Calls {@link #checkMandatoryAttributesPresence(Update)} to check for presence of mandatory attributes.
+ *
+ * @param message Update message
+ * @throws BGPDocumentedException
+ */
+ private void onUpdateMessage(final Update message) throws BGPDocumentedException {
+ checkMandatoryAttributesPresence(message);
+
+ // update AdjRibs
+ final Attributes attrs = message.getAttributes();
+ MpReachNlri mpReach = null;
+ final boolean isAnyNlriAnnounced = message.getNlri() != null;
+ if (isAnyNlriAnnounced) {
+ mpReach = prefixesToMpReach(message);
+ } else {
+ mpReach = MessageUtil.getMpReachNlri(attrs);
+ }
+ if (mpReach != null) {
+ this.ribWriter.updateRoutes(mpReach, nextHopToAttribute(attrs, mpReach));
+ }
+ MpUnreachNlri mpUnreach = null;
+ if (message.getWithdrawnRoutes() != null) {
+ mpUnreach = prefixesToMpUnreach(message, isAnyNlriAnnounced);
+ } else {
+ mpUnreach = MessageUtil.getMpUnreachNlri(attrs);
+ }
+ if (mpUnreach != null) {
+ this.ribWriter.removeRoutes(mpUnreach);
+ }
+ }
+
+ private static Attributes nextHopToAttribute(final Attributes attrs, final MpReachNlri mpReach) {
+ if (attrs.getCNextHop() == null && mpReach.getCNextHop() != null) {
+ final AttributesBuilder attributesBuilder = new AttributesBuilder(attrs);
+ attributesBuilder.setCNextHop(mpReach.getCNextHop());
+ return attributesBuilder.build();
+ }
+ return attrs;
+ }
+
+ /**
+ * Creates MPReach for the prefixes to be handled in the same way as linkstate routes
+ *
+ * @param message Update message containing prefixes in NLRI
+ * @return MpReachNlri with prefixes from the nlri field
+ */
+ private static MpReachNlri prefixesToMpReach(final Update message) {
+ final List<Ipv4Prefixes> prefixes = new ArrayList<>();
+ for (final Ipv4Prefix p : message.getNlri().getNlri()) {
+ prefixes.add(new Ipv4PrefixesBuilder().setPrefix(p).build());
+ }
+ final MpReachNlriBuilder b = new MpReachNlriBuilder().setAfi(Ipv4AddressFamily.class).setSafi(
+ UnicastSubsequentAddressFamily.class).setAdvertizedRoutes(
+ new AdvertizedRoutesBuilder().setDestinationType(
+ new DestinationIpv4CaseBuilder().setDestinationIpv4(
+ new DestinationIpv4Builder().setIpv4Prefixes(prefixes).build()).build()).build());
+ if (message.getAttributes() != null) {
+ b.setCNextHop(message.getAttributes().getCNextHop());
+ }
+ return b.build();
+ }
+
+ /**
+ * Create MPUnreach for the prefixes to be handled in the same way as linkstate routes
+ *
+ * @param message Update message containing withdrawn routes
+ * @param isAnyNlriAnnounced
+ * @return MpUnreachNlri with prefixes from the withdrawn routes field
+ */
+ private static MpUnreachNlri prefixesToMpUnreach(final Update message, final boolean isAnyNlriAnnounced) {
+ final List<Ipv4Prefixes> prefixes = new ArrayList<>();
+ for (final Ipv4Prefix p : message.getWithdrawnRoutes().getWithdrawnRoutes()) {
+ boolean nlriAnounced = false;
+ if(isAnyNlriAnnounced) {
+ nlriAnounced = message.getNlri().getNlri().contains(p);
+ }
+
+ if(!nlriAnounced) {
+ prefixes.add(new Ipv4PrefixesBuilder().setPrefix(p).build());
+ }
+ }
+ return new MpUnreachNlriBuilder().setAfi(Ipv4AddressFamily.class).setSafi(UnicastSubsequentAddressFamily.class).setWithdrawnRoutes(
+ new WithdrawnRoutesBuilder().setDestinationType(
+ new org.opendaylight.yang.gen.v1.urn.opendaylight.params.xml.ns.yang.bgp.inet.rev150305.update.attributes.mp.unreach.nlri.withdrawn.routes.destination.type.DestinationIpv4CaseBuilder().setDestinationIpv4(
+ new DestinationIpv4Builder().setIpv4Prefixes(prefixes).build()).build()).build()).build();
+ }
+
+ @Override
+ public synchronized void onSessionUp(final BGPSession session) {
+ final List<AddressFamilies> addPathTablesType = session.getAdvertisedAddPathTableTypes();
+ final Set<BgpTableType> advertizedTableTypes = session.getAdvertisedTableTypes();
+ LOG.info("Session with peer {} went up with tables {} and Add Path tables {}", this.name, advertizedTableTypes, addPathTablesType);
+ this.session = session;
+
+ this.rawIdentifier = InetAddresses.forString(session.getBgpId().getValue()).getAddress();
+ final PeerId peerId = RouterIds.createPeerId(session.getBgpId());
+
+ this.tables.addAll(advertizedTableTypes.stream().map(t -> new TablesKey(t.getAfi(), t.getSafi())).collect(Collectors.toList()));
+ final boolean announceNone = isAnnounceNone(this.simpleRoutingPolicy);
+ final Map<TablesKey, SendReceive> addPathTableMaps = mapTableTypesFamilies(addPathTablesType);
+ this.peerIId = this.rib.getYangRibId().node(org.opendaylight.yang.gen.v1.urn.opendaylight.params.xml.ns.yang.bgp.rib.rev130925.bgp.rib.rib.Peer.QNAME)
+ .node(IdentifierUtils.domPeerId(peerId));
+
+ if(!announceNone) {
+ createAdjRibOutListener(peerId);
+ }
+ this.tables.forEach(tablesKey -> {
+ final ExportPolicyPeerTracker exportTracker = this.rib.getExportPolicyPeerTracker(tablesKey);
+ if (exportTracker != null) {
+ this.tableRegistration.add(exportTracker.registerPeer(peerId, addPathTableMaps.get(tablesKey), this.peerIId, this.peerRole,
+ this.simpleRoutingPolicy));
+ }
+ });
+ addBgp4Support(peerId, announceNone);
+
+ if(!isLearnNone(this.simpleRoutingPolicy)) {
+ this.effRibInWriter = EffectiveRibInWriter.create(this.rib.getService(), this.rib.createPeerChain(this), this.peerIId,
+ this.rib.getImportPolicyPeerTracker(), this.rib.getRibSupportContext(), this.peerRole,
+ this.peerStats.getEffectiveRibInRouteCounters(), this.peerStats.getAdjRibInRouteCounters());
+ }
+ this.ribWriter = this.ribWriter.transform(peerId, this.rib.getRibSupportContext(), this.tables, addPathTableMaps);
+
+ // register BGP Peer stats
+ this.peerStats.getSessionEstablishedCounter().increaseCount();
+ if (this.registrator != null) {
+ this.runtimeReg = this.registrator.register(this);
+ }
+
+ if (this.rpcRegistry != null) {
+ this.rpcRegistration = this.rpcRegistry.addRoutedRpcImplementation(BgpPeerRpcService.class, new BgpPeerRpc(session, this.tables));
+ final KeyedInstanceIdentifier<org.opendaylight.yang.gen.v1.urn.opendaylight.params.xml.ns.yang.bgp.rib.rev130925.bgp.rib.rib.Peer, PeerKey> path =
+ this.rib.getInstanceIdentifier().child(org.opendaylight.yang.gen.v1.urn.opendaylight.params.xml.ns.yang.bgp.rib.rev130925.bgp.rib.rib.Peer.class, new PeerKey(peerId));
+ this.rpcRegistration.registerPath(PeerContext.class, path);
+ }
+
+ this.rib.getRenderStats().getConnectedPeerCounter().increaseCount();
+ }
+
+ private void createAdjRibOutListener(final PeerId peerId) {
+ this.tables.forEach(key->createAdjRibOutListener(peerId, key, true));
+ }
+
+ //try to add a support for old-school BGP-4, if peer did not advertise IPv4-Unicast MP capability
+ private void addBgp4Support(final PeerId peerId, final boolean announceNone) {
+ final TablesKey key = new TablesKey(Ipv4AddressFamily.class, UnicastSubsequentAddressFamily.class);
+ if (this.tables.add(key) && !announceNone) {
+ createAdjRibOutListener(peerId, key, false);
+ }
+ }
+
+ private void createAdjRibOutListener(final PeerId peerId, final TablesKey key, final boolean mpSupport) {
+ final RIBSupportContext context = this.rib.getRibSupportContext().getRIBSupportContext(key);
+
+ // not particularly nice
+ if (context != null && this.session instanceof BGPSessionImpl) {
+ this.adjRibOutListenerSet.put(key, AdjRibOutListener.create(peerId, key, this.rib.getYangRibId(), this.rib.getCodecsRegistry(),
+ context.getRibSupport(), this.rib.getService(), ((BGPSessionImpl) this.session).getLimiter(), mpSupport,
+ this.peerStats.getAdjRibOutRouteCounters().init(key)));
+ }
+ }
+
+ private ListenableFuture<Void> cleanup() {
+ // FIXME: BUG-196: support graceful
+ this.adjRibOutListenerSet.values().forEach(AdjRibOutListener::close);
+ this.adjRibOutListenerSet.clear();
+ if (this.effRibInWriter != null) {
+ this.effRibInWriter.close();
+ }
+ this.tables.clear();
+ if (this.ribWriter != null) {
+ return this.ribWriter.removePeer();
+ }
+ return Futures.immediateFuture(null);
+ }
+
+ @Override
+ public void onSessionDown(final BGPSession session, final Exception e) {
+ if(e.getMessage().equals(BGPSessionImpl.END_OF_INPUT)) {
+ LOG.info("Session with peer {} went down", this.name);
+ } else {
+ LOG.info("Session with peer {} went down", this.name, e);
+ }
+ releaseConnection();
+ }
+
+ @Override
+ public void onSessionTerminated(final BGPSession session, final BGPTerminationReason cause) {
+ LOG.info("Session with peer {} terminated: {}", this.name, cause);
+ releaseConnection();
+ }
+
+ @Override
+ public String toString() {
+ return addToStringAttributes(MoreObjects.toStringHelper(this)).toString();
+ }
+
+ protected ToStringHelper addToStringAttributes(final ToStringHelper toStringHelper) {
+ toStringHelper.add("name", this.name);
+ toStringHelper.add("tables", this.tables);
+ return toStringHelper;
+ }
+
+ @Override
+ public String getName() {
+ return this.name;
+ }
+
+ @Override
+ public synchronized ListenableFuture<Void> releaseConnection() {
+ if (this.rpcRegistration != null) {
+ this.rpcRegistration.close();
+ }
+ closeRegistration();
+ final ListenableFuture<Void> future = cleanup();
+ dropConnection();
+ return future;
+ }
+
+ private void closeRegistration() {
+ for (final AbstractRegistration tableCloseable : this.tableRegistration) {
+ tableCloseable.close();
+ }
+ this.tableRegistration.clear();
+ }
+
+ private void dropConnection() {
+ if (this.runtimeReg != null) {
+ this.runtimeReg.close();
+ this.runtimeReg = null;
+ }
+ if (this.session != null) {
+ try {
+ this.session.close();
+ } catch (final Exception e) {
+ LOG.warn("Error closing session with peer", e);
+ }
+ this.session = null;
+
+ this.rib.getRenderStats().getConnectedPeerCounter().decreaseCount();
+ }
+ }
+
+ @Override
+ public synchronized byte[] getRawIdentifier() {
+ return Arrays.copyOf(this.rawIdentifier, this.rawIdentifier.length);
+ }
+
+ @Override
+ public void resetSession() {
+ releaseConnection();
+ }
+
+ @Override
+ public void resetStats() {
+ if (this.session instanceof BGPSessionStats) {
+ ((BGPSessionStats) this.session).resetBgpSessionStats();
+ }
+ }
+
+ @Override
+ public BgpSessionState getBgpSessionState() {
+ if (this.session instanceof BGPSessionStats) {
+ return ((BGPSessionStats) this.session).getBgpSessionState();
+ }
+ return new BgpSessionState();
+ }
+
+ @Override
+ public synchronized BgpPeerState getBgpPeerState() {
+ return this.peerStats.getBgpPeerState();
+ }
+
+ @Override
+ public void onTransactionChainFailed(final TransactionChain<?, ?> chain, final AsyncTransaction<?, ?> transaction, final Throwable cause) {
+ LOG.error("Transaction chain failed.", cause);
+ this.chain.close();
+ this.chain = this.rib.createPeerChain(this);
+ this.ribWriter = AdjRibInWriter.create(this.rib.getYangRibId(), this.peerRole, this.simpleRoutingPolicy, this.chain);
+ releaseConnection();
+ }
+
+ @Override
+ public void onTransactionChainSuccessful(final TransactionChain<?, ?> chain) {
+ LOG.debug("Transaction chain {} successfull.", chain);
+ }
+
+ @Override
+ public void markUptodate(final TablesKey tablesKey) {
+ this.ribWriter.markTableUptodate(tablesKey);
+ }
+
+ private static Map<TablesKey, SendReceive> mapTableTypesFamilies(final List<AddressFamilies> addPathTablesType) {
+ return ImmutableMap.copyOf(addPathTablesType.stream().collect(Collectors.toMap(af -> new TablesKey(af.getAfi(), af.getSafi()),
+ BgpAddPathTableType::getSendReceive)));
+ }