diff --git a/cepheus-cep/src/main/java/com/orange/cepheus/cep/SubscriptionManager.java b/cepheus-cep/src/main/java/com/orange/cepheus/cep/SubscriptionManager.java index 1e9db5d0..3cdfd44a 100644 --- a/cepheus-cep/src/main/java/com/orange/cepheus/cep/SubscriptionManager.java +++ b/cepheus-cep/src/main/java/com/orange/cepheus/cep/SubscriptionManager.java @@ -258,7 +258,25 @@ private void subscribeProvider(Provider provider, SubscribeContext subscribeCont logger.warn("Error during subscription for {}", provider.getUrl(), throwable); }); } - + + /** + * Unsubscribe from Provider of the configuration that is being deleted + * @param configuration the configuration that is being deleted + */ + public void unsubscribe(Configuration configuration) { + + List eventTypeIns = Collections.emptyList(); + eventTypeIns = configuration.getEventTypeIns(); + + for (EventTypeIn eventType : eventTypeIns) { + for (Provider provider : eventType.getProviders()) { + if (provider != null) { + unsubscribeProvider(provider); + } + } + } + } + /** * Unsubscribe from a provider * @param provider the provider to unusubscribe from diff --git a/cepheus-cep/src/main/java/com/orange/cepheus/cep/controller/AdminController.java b/cepheus-cep/src/main/java/com/orange/cepheus/cep/controller/AdminController.java index d4a8caa3..04dcf686 100755 --- a/cepheus-cep/src/main/java/com/orange/cepheus/cep/controller/AdminController.java +++ b/cepheus-cep/src/main/java/com/orange/cepheus/cep/controller/AdminController.java @@ -144,6 +144,13 @@ public synchronized ResponseEntity removeConfiguration() throws P if (tenantFilter != null) { tenantFilter.removeTenant(configurationId); } + + //Retrieve the configuration + final Configuration configuration = complexEventProcessor.getConfiguration(); + //Unsubscribe from Provider of the configuration + if (configuration != null) { + subscriptionManager.unsubscribe(configuration); + } // Reset the CEP complexEventProcessor.reset(); // Delete the persisted configuration