BUG 2854 : Do not add empty read write transactions to the replicable journal
[controller.git] / opendaylight / md-sal / messagebus-impl / src / main / java / org / opendaylight / controller / messagebus / app / impl / NetconfEventSource.java
1 /*
2  * Copyright (c) 2013 Cisco Systems, Inc. and others.  All rights reserved.
3  *
4  * This program and the accompanying materials are made available under the
5  * terms of the Eclipse Public License v1.0 which accompanies this distribution,
6  * and is available at http://www.eclipse.org/legal/epl-v10.html
7  */
8
9 package org.opendaylight.controller.messagebus.app.impl;
10
11 import com.google.common.base.Preconditions;
12 import java.util.ArrayList;
13 import java.util.List;
14 import java.util.Map;
15 import java.util.Set;
16 import java.util.concurrent.Future;
17 import java.util.regex.Pattern;
18 import org.opendaylight.controller.md.sal.binding.api.DataChangeListener;
19 import org.opendaylight.controller.md.sal.common.api.data.AsyncDataChangeEvent;
20 import org.opendaylight.controller.mdsal.MdSAL;
21 import org.opendaylight.controller.sal.core.api.notify.NotificationListener;
22 import org.opendaylight.yang.gen.v1.urn.cisco.params.xml.ns.yang.messagebus.eventaggregator.rev141202.NotificationPattern;
23 import org.opendaylight.yang.gen.v1.urn.cisco.params.xml.ns.yang.messagebus.eventaggregator.rev141202.TopicNotification;
24 import org.opendaylight.yang.gen.v1.urn.cisco.params.xml.ns.yang.messagebus.eventsource.rev141202.EventSourceService;
25 import org.opendaylight.yang.gen.v1.urn.cisco.params.xml.ns.yang.messagebus.eventsource.rev141202.JoinTopicInput;
26 import org.opendaylight.yang.gen.v1.urn.cisco.params.xml.ns.yang.messagebus.eventsource.rev141202.JoinTopicOutput;
27 import org.opendaylight.yang.gen.v1.urn.ietf.params.xml.ns.netconf.notification._1._0.rev080714.CreateSubscriptionInput;
28 import org.opendaylight.yang.gen.v1.urn.ietf.params.xml.ns.netconf.notification._1._0.rev080714.CreateSubscriptionInputBuilder;
29 import org.opendaylight.yang.gen.v1.urn.ietf.params.xml.ns.netconf.notification._1._0.rev080714.NotificationsService;
30 import org.opendaylight.yang.gen.v1.urn.ietf.params.xml.ns.netconf.notification._1._0.rev080714.StreamNameType;
31 import org.opendaylight.yang.gen.v1.urn.opendaylight.netconf.node.inventory.rev140108.NetconfNode;
32 import org.opendaylight.yangtools.concepts.ListenerRegistration;
33 import org.opendaylight.yangtools.yang.binding.DataObject;
34 import org.opendaylight.yangtools.yang.binding.InstanceIdentifier;
35 import org.opendaylight.yangtools.yang.common.QName;
36 import org.opendaylight.yangtools.yang.common.RpcResult;
37 import org.opendaylight.yangtools.yang.data.api.CompositeNode;
38 import org.opendaylight.yangtools.yang.data.impl.ImmutableCompositeNode;
39 import org.opendaylight.yangtools.yang.model.api.NotificationDefinition;
40 import org.slf4j.Logger;
41 import org.slf4j.LoggerFactory;
42
43 public class NetconfEventSource implements EventSourceService, NotificationListener, DataChangeListener {
44     private static final Logger LOGGER = LoggerFactory.getLogger(NetconfEventSource.class);
45
46     private final MdSAL mdSal;
47     private final String nodeId;
48
49     private final List<String> activeStreams = new ArrayList<>();
50
51     private final Map<String, String> urnPrefixToStreamMap;
52
53     public NetconfEventSource(final MdSAL mdSal, final String nodeId, final Map<String, String> streamMap) {
54         Preconditions.checkNotNull(mdSal);
55         Preconditions.checkNotNull(nodeId);
56
57         this.mdSal = mdSal;
58         this.nodeId = nodeId;
59         this.urnPrefixToStreamMap = streamMap;
60
61         LOGGER.info("NetconfEventSource [{}] created.", nodeId);
62     }
63
64     @Override
65     public Future<RpcResult<JoinTopicOutput>> joinTopic(final JoinTopicInput input) {
66         final NotificationPattern notificationPattern = input.getNotificationPattern();
67
68         // FIXME: default language should already be regex
69         final String regex = Util.wildcardToRegex(notificationPattern.getValue());
70
71         final Pattern pattern = Pattern.compile(regex);
72         List<QName> matchingNotifications = Util.expandQname(availableNotifications(), pattern);
73         registerNotificationListener(matchingNotifications);
74         return null;
75     }
76
77     private List<QName> availableNotifications() {
78         // FIXME: use SchemaContextListener to get changes asynchronously
79         Set<NotificationDefinition> availableNotifications = mdSal.getSchemaContext(nodeId).getNotifications();
80         List<QName> qNs = new ArrayList<>(availableNotifications.size());
81         for (NotificationDefinition nd : availableNotifications) {
82             qNs.add(nd.getQName());
83         }
84
85         return qNs;
86     }
87
88     private void registerNotificationListener(final List<QName> notificationsToSubscribe) {
89         for (QName qName : notificationsToSubscribe) {
90             startSubscription(qName);
91             // FIXME: do not lose this registration
92             final ListenerRegistration<NotificationListener> reg = mdSal.addNotificationListener(nodeId, qName, this);
93         }
94     }
95
96     private synchronized void startSubscription(final QName qName) {
97         String streamName = resolveStream(qName);
98
99         if (streamIsActive(streamName) == false) {
100             LOGGER.info("Stream {} is not active on node {}. Will subscribe.", streamName, nodeId);
101             startSubscription(streamName);
102         }
103     }
104
105     private synchronized void resubscribeToActiveStreams() {
106         for (String streamName : activeStreams) {
107             startSubscription(streamName);
108         }
109     }
110
111     private synchronized void startSubscription(final String streamName) {
112         CreateSubscriptionInput subscriptionInput = getSubscriptionInput(streamName);
113         mdSal.getRpcService(nodeId, NotificationsService.class).createSubscription(subscriptionInput);
114         activeStreams.add(streamName);
115     }
116
117     private static CreateSubscriptionInput getSubscriptionInput(final String streamName) {
118         CreateSubscriptionInputBuilder csib = new CreateSubscriptionInputBuilder();
119         csib.setStream(new StreamNameType(streamName));
120         return csib.build();
121     }
122
123     private String resolveStream(final QName qName) {
124         String streamName = null;
125
126         for (Map.Entry<String, String> entry : urnPrefixToStreamMap.entrySet()) {
127             String nameSpace = qName.getNamespace().toString();
128             String urnPrefix = entry.getKey();
129             if( nameSpace.startsWith(urnPrefix) ) {
130                 streamName = entry.getValue();
131                 break;
132             }
133         }
134
135         return streamName;
136     }
137
138     private boolean streamIsActive(final String streamName) {
139         return activeStreams.contains(streamName);
140     }
141
142     // PASS
143     @Override public Set<QName> getSupportedNotifications() {
144         return null;
145     }
146
147     @Override
148     public void onNotification(final CompositeNode notification) {
149         LOGGER.info("NetconfEventSource {} received notification {}. Will publish to MD-SAL.", nodeId, notification);
150         ImmutableCompositeNode payload = ImmutableCompositeNode.builder()
151                 .setQName(QName.create(TopicNotification.QNAME, "payload"))
152                 .add(notification).toInstance();
153         ImmutableCompositeNode icn = ImmutableCompositeNode.builder()
154                 .setQName(TopicNotification.QNAME)
155                 .add(payload)
156                 .addLeaf("event-source", nodeId)
157                 .toInstance();
158
159         mdSal.publishNotification(icn);
160     }
161
162     @Override
163     public void onDataChanged(final AsyncDataChangeEvent<InstanceIdentifier<?>, DataObject> change) {
164         boolean wasConnected = false;
165         boolean nowConnected = false;
166
167         for (Map.Entry<InstanceIdentifier<?>, DataObject> changeEntry : change.getOriginalData().entrySet()) {
168             if ( isNetconfNode(changeEntry) ) {
169                 NetconfNode nn = (NetconfNode)changeEntry.getValue();
170                 wasConnected = nn.isConnected();
171             }
172         }
173
174         for (Map.Entry<InstanceIdentifier<?>, DataObject> changeEntry : change.getUpdatedData().entrySet()) {
175             if ( isNetconfNode(changeEntry) ) {
176                 NetconfNode nn = (NetconfNode)changeEntry.getValue();
177                 nowConnected = nn.isConnected();
178             }
179         }
180
181         if (wasConnected == false && nowConnected == true) {
182             resubscribeToActiveStreams();
183         }
184     }
185
186     private static boolean isNetconfNode(final Map.Entry<InstanceIdentifier<?>, DataObject> changeEntry )  {
187         return NetconfNode.class.equals(changeEntry.getKey().getTargetType());
188     }
189
190 }