2 * ============LICENSE_START=======================================================
3 * BBS-RELOCATION-CPE-AUTHENTICATION-HANDLER
4 * ================================================================================
5 * Copyright (C) 2019 NOKIA Intellectual Property. All rights reserved.
6 * ================================================================================
7 * Licensed under the Apache License, Version 2.0 (the "License");
8 * you may not use this file except in compliance with the License.
9 * You may obtain a copy of the License at
11 * http://www.apache.org/licenses/LICENSE-2.0
13 * Unless required by applicable law or agreed to in writing, software
14 * distributed under the License is distributed on an "AS IS" BASIS,
15 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16 * See the License for the specific language governing permissions and
17 * limitations under the License.
18 * ============LICENSE_END=========================================================
21 package org.onap.bbs.event.processor.pipelines;
23 import static org.onap.bbs.event.processor.config.ApplicationConstants.CONSUME_REREGISTRATION_TASK_NAME;
24 import static org.onap.bbs.event.processor.config.ApplicationConstants.DCAE_BBS_EVENT_PROCESSOR_MS_INSTANCE;
25 import static org.onap.bbs.event.processor.config.ApplicationConstants.RETRIEVE_HSI_CFS_SERVICE_INSTANCE_TASK_NAME;
26 import static org.onap.bbs.event.processor.config.ApplicationConstants.RETRIEVE_PNF_TASK_NAME;
27 import static org.onap.dcaegen2.services.sdk.rest.services.model.logging.MdcVariables.INSTANCE_UUID;
28 import static org.onap.dcaegen2.services.sdk.rest.services.model.logging.MdcVariables.RESPONSE_CODE;
30 import java.time.Duration;
31 import java.time.Instant;
32 import java.util.HashMap;
33 import java.util.List;
35 import java.util.UUID;
36 import java.util.concurrent.TimeoutException;
38 import javax.net.ssl.SSLException;
40 import org.onap.bbs.event.processor.config.ApplicationConfiguration;
41 import org.onap.bbs.event.processor.exceptions.EmptyDmaapResponseException;
42 import org.onap.bbs.event.processor.model.ControlLoopPublisherDmaapModel;
43 import org.onap.bbs.event.processor.model.ImmutableControlLoopPublisherDmaapModel;
44 import org.onap.bbs.event.processor.model.PnfAaiObject;
45 import org.onap.bbs.event.processor.model.ReRegistrationConsumerDmaapModel;
46 import org.onap.bbs.event.processor.model.RelationshipListAaiObject;
47 import org.onap.bbs.event.processor.model.ServiceInstanceAaiObject;
48 import org.onap.bbs.event.processor.tasks.AaiClientTask;
49 import org.onap.bbs.event.processor.tasks.DmaapPublisherTask;
50 import org.onap.bbs.event.processor.tasks.DmaapReRegistrationConsumerTask;
51 import org.onap.dcaegen2.services.sdk.rest.services.adapters.http.HttpResponse;
52 import org.slf4j.Logger;
53 import org.slf4j.LoggerFactory;
55 import org.springframework.beans.factory.annotation.Autowired;
56 import org.springframework.stereotype.Component;
57 import org.springframework.util.StringUtils;
59 import reactor.core.publisher.Flux;
60 import reactor.core.publisher.Mono;
63 public class ReRegistrationPipeline {
65 private static final Logger LOGGER = LoggerFactory.getLogger(ReRegistrationPipeline.class);
67 private static final String POLICY_NAME = "Nomadic_ONT";
69 private DmaapReRegistrationConsumerTask consumerTask;
70 private DmaapPublisherTask publisherTask;
71 private AaiClientTask aaiClientTask;
73 private ApplicationConfiguration configuration;
75 private Map<String, String> mdcContextMap;
78 ReRegistrationPipeline(ApplicationConfiguration configuration,
79 DmaapReRegistrationConsumerTask consumerTask,
80 DmaapPublisherTask publisherTask,
81 AaiClientTask aaiClientTask,
82 Map<String, String> mdcContextMap) {
83 this.configuration = configuration;
84 this.consumerTask = consumerTask;
85 this.publisherTask = publisherTask;
86 this.aaiClientTask = aaiClientTask;
87 this.mdcContextMap = mdcContextMap;
91 * PNF re-registration processing pipeline for BBS uS.
93 public void processPnfReRegistrationEvents() {
94 MDC.setContextMap(mdcContextMap);
95 LOGGER.info("Process next Re-Registration events");
97 .subscribe(this::onSuccess, this::onError, this::onComplete);
98 LOGGER.trace("Reactive PNF Re-registration pipeline subscribed - Execution started");
101 Flux<HttpResponse> executePipeline() {
103 // Consume Re-Registration from DMaaP
104 consumeReRegistrationsFromDmaap()
105 // Fetch PNF from A&AI
106 .flatMap(this::fetchPnfFromAai)
107 // Fetch related HSI CFS instance from A&AI
108 .flatMap(this::fetchHsiCfsServiceInstanceFromAai)
109 // Trigger Policy for relocation
110 .flatMap(this::triggerPolicy);
113 private void onSuccess(HttpResponse responseCode) {
114 MDC.put(RESPONSE_CODE, String.valueOf(responseCode.statusCode()));
115 LOGGER.info("PNF Re-Registration event successfully handled. "
116 + "Publishing to DMaaP for Policy returned a status code of ({} {})",
117 responseCode.statusCode(), responseCode.statusReason());
118 MDC.remove(RESPONSE_CODE);
121 private void onError(Throwable throwable) {
122 LOGGER.error("Aborted PNF Re-Registration events processing. Error: {}", throwable.getMessage());
125 private void onComplete() {
126 LOGGER.info("PNF Re-Registration processing pipeline has been completed");
129 private Flux<PipelineState> consumeReRegistrationsFromDmaap() {
130 return Flux.defer(() -> {
131 MDC.put(INSTANCE_UUID, UUID.randomUUID().toString());
133 return consumerTask.execute(CONSUME_REREGISTRATION_TASK_NAME)
134 .timeout(Duration.ofSeconds(configuration.getPipelinesTimeoutInSeconds()))
136 if (e instanceof TimeoutException) {
137 LOGGER.warn("Timed out waiting for DMaaP response");
138 } else if (e instanceof EmptyDmaapResponseException) {
139 LOGGER.info("Nothing to consume from DMaaP");
141 LOGGER.error("DMaaP Consumer error: {}", e.getMessage());
142 LOGGER.debug("Error\n", e);
146 e -> e instanceof Exception,
149 // For each message, we have to keep separate state. This state will be enhanced
150 // in each step and handed off to the next processing step
151 PipelineState state = new PipelineState();
152 state.setReRegistrationEvent(event);
155 } catch (SSLException e) {
156 return Flux.error(e);
161 private Mono<PipelineState> fetchPnfFromAai(PipelineState state) {
163 ReRegistrationConsumerDmaapModel vesEvent = state.getReRegistrationEvent();
164 String pnfName = vesEvent.getCorrelationId();
165 String url = String.format("/aai/v14/network/pnfs/pnf/%s?depth=all", pnfName);
166 LOGGER.debug("Processing Step: Retrieve PNF. Url: ({})", url);
168 return aaiClientTask.executePnfRetrieval(RETRIEVE_PNF_TASK_NAME, url)
169 .timeout(Duration.ofSeconds(configuration.getPipelinesTimeoutInSeconds()))
170 .doOnError(TimeoutException.class,
171 e -> LOGGER.warn("Timed out waiting for A&AI response")
174 LOGGER.error("Error while retrieving PNF: {}", e.getMessage());
175 LOGGER.debug("Error\n", e);
179 e -> e instanceof Exception,
182 state.setPnfAaiObject(p);
187 private Mono<PipelineState> fetchHsiCfsServiceInstanceFromAai(PipelineState state) {
189 if (state == null || state.getPnfAaiObject() == null) {
193 // At this point, we have both the VES-event of the re-registration and the PNF object retrieved from A&AI
194 // We can check if this processing needs to continue in case of a true relocation
195 if (isNotReallyAnOntRelocation(state)) {
199 PnfAaiObject pnf = state.getPnfAaiObject();
200 // Assuming that the PNF will only have a single service-instance relationship pointing
201 // towards the HSI CFS service
202 String serviceInstanceId = pnf.getRelationshipListAaiObject().getRelationshipEntries()
204 .filter(e -> "service-instance".equals(e.getRelatedTo()))
205 .flatMap(e -> e.getRelationshipData().stream())
206 .filter(d -> "service-instance.service-instance-id".equals(d.getRelationshipKey()))
207 .map(RelationshipListAaiObject.RelationshipDataEntryAaiObject::getRelationshipValue)
208 .findFirst().orElse("");
210 if (StringUtils.isEmpty(serviceInstanceId)) {
211 LOGGER.error("Unable to retrieve HSI CFS service instance from PNF {}",
212 state.getPnfAaiObject().getPnfName());
216 String url = String.format("/aai/v14/nodes/service-instances/service-instance/%s?depth=all",
218 LOGGER.debug("Processing Step: Retrieve HSI CFS Service. Url: ({})", url);
219 return aaiClientTask.executeServiceInstanceRetrieval(RETRIEVE_HSI_CFS_SERVICE_INSTANCE_TASK_NAME, url)
220 .timeout(Duration.ofSeconds(configuration.getPipelinesTimeoutInSeconds()))
221 .doOnError(TimeoutException.class,
222 e -> LOGGER.warn("Timed out waiting for A&AI response")
225 LOGGER.error("Error while retrieving HSI CFS Service instance: {}", e.getMessage());
226 LOGGER.debug("Error\n", e);
230 e -> e instanceof Exception,
233 state.setHsiCfsServiceInstance(s);
238 private boolean isNotReallyAnOntRelocation(PipelineState state) {
239 List<RelationshipListAaiObject.RelationshipEntryAaiObject> relationshipEntries =
240 state.getPnfAaiObject().getRelationshipListAaiObject().getRelationshipEntries();
242 // If no logical-link, fail further processing
243 if (relationshipEntries.stream().noneMatch(e -> "logical-link".equals(e.getRelatedTo()))) {
244 LOGGER.warn("PNF {} does not have any logical-links bridged. Stop further processing",
245 state.getPnfAaiObject().getPnfName());
249 // Assuming PNF will only have one logical-link per BBS use case design
250 boolean isNotRelocation = relationshipEntries
252 .filter(e -> "logical-link".equals(e.getRelatedTo()))
253 .flatMap(e -> e.getRelationshipData().stream())
254 .anyMatch(d -> d.getRelationshipValue()
255 .equals(state.getReRegistrationEvent().getAttachmentPoint()));
258 if (isNotRelocation) {
259 LOGGER.warn("Not a Relocation for PNF {} with attachment point {}",
260 state.getPnfAaiObject().getPnfName(),
261 state.getReRegistrationEvent().getAttachmentPoint());
263 return isNotRelocation;
266 private Mono<HttpResponse> triggerPolicy(PipelineState state) {
268 if (state == null || state.getHsiCfsServiceInstance() == null) {
272 ControlLoopPublisherDmaapModel event = buildTriggeringPolicyEvent(state);
273 return publisherTask.execute(event)
274 .timeout(Duration.ofSeconds(configuration.getPipelinesTimeoutInSeconds()))
275 .doOnError(TimeoutException.class,
276 e -> LOGGER.warn("Timed out waiting for DMaaP confirmation")
278 .doOnError(e -> LOGGER.error("Error while triggering Policy: {}", e.getMessage()))
280 e -> e instanceof Exception,
284 private ControlLoopPublisherDmaapModel buildTriggeringPolicyEvent(PipelineState state) {
286 String cfsServiceInstanceId = state.getHsiCfsServiceInstance().getServiceInstanceId();
288 String attachmentPoint = state.getReRegistrationEvent().getAttachmentPoint();
289 String remoteId = state.getReRegistrationEvent().getRemoteId();
290 String cvlan = state.getReRegistrationEvent().getCVlan();
291 String svlan = state.getReRegistrationEvent().getSVlan();
293 Map<String, String> enrichmentData = new HashMap<>();
294 enrichmentData.put("service-information.hsia-cfs-service-instance-id", cfsServiceInstanceId);
296 enrichmentData.put("attachmentPoint", attachmentPoint);
297 enrichmentData.put("remoteId", remoteId);
298 enrichmentData.put("cvlan", cvlan);
299 enrichmentData.put("svlan", svlan);
301 ControlLoopPublisherDmaapModel triggerEvent = ImmutableControlLoopPublisherDmaapModel.builder()
302 .closedLoopEventClient(DCAE_BBS_EVENT_PROCESSOR_MS_INSTANCE)
303 .policyVersion(configuration.getPolicyVersion())
304 .policyName(POLICY_NAME)
305 .policyScope(configuration.getReRegistrationCloseLoopPolicyScope())
306 .targetType(configuration.getCloseLoopTargetType())
307 .aaiEnrichmentData(enrichmentData)
308 .closedLoopAlarmStart(Instant.now().getEpochSecond())
309 .closedLoopEventStatus(configuration.getCloseLoopEventStatus())
310 .closedLoopControlName(configuration.getReRegistrationCloseLoopControlName())
311 .version(configuration.getCloseLoopVersion())
312 .target(configuration.getCloseLoopTarget())
313 .requestId(UUID.randomUUID().toString())
314 .originator(configuration.getCloseLoopOriginator())
316 LOGGER.debug("Processing Step: Publish for Policy");
317 LOGGER.trace("Trigger Policy event: ({})",triggerEvent);
321 private static class PipelineState {
323 private ReRegistrationConsumerDmaapModel reRegistrationEvent;
324 private PnfAaiObject pnfAaiObject;
325 private ServiceInstanceAaiObject hsiCfsServiceInstance;
327 ReRegistrationConsumerDmaapModel getReRegistrationEvent() {
328 return reRegistrationEvent;
331 void setReRegistrationEvent(ReRegistrationConsumerDmaapModel reRegistrationEvent) {
332 this.reRegistrationEvent = reRegistrationEvent;
335 PnfAaiObject getPnfAaiObject() {
339 void setPnfAaiObject(PnfAaiObject pnfAaiObject) {
340 this.pnfAaiObject = pnfAaiObject;
343 ServiceInstanceAaiObject getHsiCfsServiceInstance() {
344 return hsiCfsServiceInstance;
347 void setHsiCfsServiceInstance(ServiceInstanceAaiObject hsiCfsServiceInstance) {
348 this.hsiCfsServiceInstance = hsiCfsServiceInstance;