This is an automated email from the ASF dual-hosted git repository.

riemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/dev by this push:
     new 552f3fed3 Improve startup behaviour (#2105)
552f3fed3 is described below

commit 552f3fed38abb53c8ea2657b3e4bc242215a67fa
Author: Dominik Riemer <[email protected]>
AuthorDate: Thu Nov 2 18:53:17 2023 +0100

    Improve startup behaviour (#2105)
    
    * feat(#2002): Align adapter registration with other pipeline elements
    
    * Fix checkstyle
    
    * style: remove trailing whitespace
    
    * Add initial draft of migration concept
    
    * refactor: fix logger configuration
    
    * refactor: extend storage implementations by method to get all instances 
by the app id
    
    * refactor: update generated typescript model
    
    * feat: add version to models & builders
    
    * refactor: implement string representation of Notification
    
    * feat: implement data model for migration
    
    * feat: register migrations at service
    
    * Revert "refactor: implement string representation of Notification"
    
    This reverts commit 646e792d50b3022748e99be738c390270b87c345.
    
    * refactor: use correct Notification class
    
    * feat: introduce migrate extensions resource
    
    * feat: introduce migrate adapter endpoint
    
    * feat: implement adapter migration at the core
    
    * remove data lake migration
    
    * ensure order & uniqueness of migrations
    
    * remove redundant exception
    
    * remove redundant exception
    
    * add tests
    
    * remove outdated test
    
    * refactor: separate adapter migration from pipeline element migrations
    
    * refactor: move MigrationResult to StreamPipes model
    
    * refactor: migration result
    
    * refactor: introduce generic migration request
    
    * feat: send migration requests to core
    
    * feat: process migrations at core
    
    * refactor: remove legacy generic
    
    * refactor: introduce versioned StreamPipes entity
    
    * refactor: remove deprecated generic type
    
    * feat: implement migration for processing elements & data sinks
    
    * docs: add endpoint documentation
    
    * refactor: move to correct module
    
    * feature: add update for descriptions
    
    * refactor: adapt ProcessingElementBuilder to be capable of versions
    
    * refactor: minor improvements
    
    * style: fix checkstyle issues
    
    * refactor: remove legacy type definition
    
    * refactor: update generated TS models
    
    * fix: add missing license header
    
    * Fix adapter model migration, add OPC adapter migration as sample
    
    * Fix typo
    
    * Extract MigrationResource logic into smaller units
    
    * Use single request for submitting migrations from extensions to core
    
    * Improve execution order of migrations and service startup tasks
    
    * Improve exception logging
    
    * Improve pipeline health check
    
    * Properly execute migration of adapter models
    
    * Fix adapter model migration
    
    * Improve service health check
    
    * Add more checks to adapter migration
    
    * Simplify migration request handling
    
    * Fix registration
    
    ---------
    
    Co-authored-by: bossenti <[email protected]>
---
 .../management/health/AdapterHealthCheck.java      |  12 +-
 .../management/management/WorkerRestClient.java    |  11 +-
 .../extensions/api/migration/IAdapterMigrator.java |   2 +-
 ...ataSinkMigrator.java => IDataSinkMigrator.java} |   2 +-
 .../model/configuration/SpCoreConfiguration.java   |   9 ++
 .../configuration/SpCoreConfigurationStatus.java   |  10 +-
 .../adapter/migration/MigrationHelpers.java        |   4 +
 .../svcdiscovery/SpServiceRegistration.java        |  18 +--
 .../extensions/svcdiscovery/SpServiceStatus.java   |  11 +-
 .../model/graph/DataProcessorInvocation.java       |   1 +
 .../model/graph/DataSinkInvocation.java            |   1 +
 .../ExtensionsServiceEndpointGenerator.java        |   2 +-
 .../manager/health/CoreServiceStatusManager.java   |  59 +++++++++
 .../manager/health/PipelineHealthCheck.java        |  21 +--
 .../manager/health/ServiceHealthCheck.java         |  31 ++---
 .../manager/health/ServiceRegistrationManager.java | 104 +++++++++++++++
 .../migration/AdapterDescriptionMigration093.java  |  73 +++++++++++
 .../AdapterDescriptionMigration093Provider.java    |  27 ++--
 .../migration/PipelineElementMigrationManager.java | 143 +++++++++++----------
 .../manager/setup/SpCoreConfigurationStep.java     |  10 +-
 .../manager/setup/StreamPipesEnvChecker.java       |   2 +-
 .../migration/DataSinkMigrationResource.java       |   4 +-
 .../rest/impl/admin/MigrationResource.java         |  60 +++++++--
 .../impl/admin/ServiceRegistrationResource.java    |  12 +-
 .../streampipes/service/core/PostStartupTask.java  |   7 +-
 .../service/core/StreamPipesCoreApplication.java   |  53 +++++---
 .../core/migrations/v093/AdapterMigration.java     |  29 +++--
 .../migrations/v093/ConsulConfigMigration.java     |   2 +
 .../svcdiscovery/SpServiceDiscoveryCore.java       |   3 +-
 .../service/extensions/CoreRequestSubmitter.java   |  54 ++++++++
 .../extensions/ExtensionsModelSubmitter.java       |   8 +-
 .../StreamPipesExtensionsServiceBase.java          |  30 ++---
 .../storage/api/ISpCoreConfigurationStorage.java   |   2 +
 .../couchdb/impl/CoreConfigurationStorageImpl.java |   5 +
 .../src/lib/model/gen/streampipes-model.ts         |  13 +-
 .../registered-extensions-services.component.html  |   7 +-
 36 files changed, 625 insertions(+), 217 deletions(-)

diff --git 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
index 07d70cf7d..e59187ea6 100644
--- 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
+++ 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/health/AdapterHealthCheck.java
@@ -34,7 +34,7 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 
-public class AdapterHealthCheck {
+public class AdapterHealthCheck implements Runnable {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(AdapterHealthCheck.class);
 
@@ -52,6 +52,11 @@ public class AdapterHealthCheck {
     this.adapterMasterManagement = adapterMasterManagement;
   }
 
+  @Override
+  public void run() {
+    this.checkAndRestoreAdapters();
+  }
+
   /**
    * In this method it is checked which adapters are currently running.
    * Then it calls all workers to validate if the adapter instance is
@@ -114,7 +119,7 @@ public class AdapterHealthCheck {
         allRunningInstancesOfOneWorker.forEach(adapterDescription ->
             
allRunningInstancesAdapterDescription.remove(adapterDescription.getElementId()));
       } catch (AdapterException e) {
-        e.printStackTrace();
+        LOG.info("Could not recover adapter at endpoint {} due to {}", 
adapterEndpointUrl, e.getMessage());
       }
     });
 
@@ -130,10 +135,9 @@ public class AdapterHealthCheck {
           
this.adapterMasterManagement.startStreamAdapter(adapterDescription.getElementId());
         }
       } catch (AdapterException e) {
-        LOG.warn("Could not start adapter {}", adapterDescription.getName(), 
e);
+        LOG.warn("Could not start adapter {} ({})", 
adapterDescription.getName(), e.getMessage());
       }
     }
 
   }
-
 }
diff --git 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
index a551201e2..c71545907 100644
--- 
a/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
+++ 
b/streampipes-connect-management/src/main/java/org/apache/streampipes/connect/management/management/WorkerRestClient.java
@@ -73,14 +73,12 @@ public class WorkerRestClient {
 
   public static List<AdapterDescription> 
getAllRunningAdapterInstanceDescriptions(String url) throws AdapterException {
     try {
-      LOG.info("Requesting all running adapter description instances: " + url);
       var responseString = ExtensionServiceExecutions
               .extServiceGetRequest(url)
               .execute().returnContent().asString();
 
       return JacksonSerializer.getObjectMapper().readValue(responseString, 
List.class);
     } catch (IOException e) {
-      LOG.error("List of running adapters could not be fetched", e);
       throw new AdapterException("List of running adapters could not be 
fetched from: " + url);
     }
   }
@@ -112,9 +110,6 @@ public class WorkerRestClient {
         var exception = getSerializer().readValue(responseString, 
AdapterException.class);
         throw new AdapterException(exception.getMessage(), 
exception.getCause());
       }
-
-      LOG.info("Adapter {} on endpoint: " + url + " with Response: ", 
ad.getName() + responseString);
-
     } catch (IOException e) {
       LOG.error("Adapter was not {} successfully", action, e);
       throw new AdapterException("Adapter was not " + action + " successfully 
with url " + url, e);
@@ -153,8 +148,7 @@ public class WorkerRestClient {
         throw new SpConfigurationException(exception.getMessage(), 
exception.getCause());
       }
     } catch (IOException e) {
-      e.printStackTrace();
-      throw new AdapterException("Could not resolve runtime configurations 
from " + url);
+      throw new AdapterException("Could not resolve runtime configurations 
from " + url, e);
     }
   }
 
@@ -178,11 +172,10 @@ public class WorkerRestClient {
     String url = baseUrl + "/assets/icon";
 
     try {
-      byte[] responseString = Request.Get(url)
+      return Request.Get(url)
               .connectTimeout(1000)
               .socketTimeout(100000)
               .execute().returnContent().asBytes();
-      return responseString;
     } catch (IOException e) {
       LOG.error(e.getMessage());
       throw new AdapterException("Could not get icon endpoint: " + url);
diff --git 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/IAdapterMigrator.java
 
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/IAdapterMigrator.java
index e7a0996dd..64ead7f63 100644
--- 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/IAdapterMigrator.java
+++ 
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/IAdapterMigrator.java
@@ -21,5 +21,5 @@ package org.apache.streampipes.extensions.api.migration;
 import 
org.apache.streampipes.extensions.api.extractor.IStaticPropertyExtractor;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 
-public  interface IAdapterMigrator extends IModelMigrator<AdapterDescription, 
IStaticPropertyExtractor> {
+public interface IAdapterMigrator extends IModelMigrator<AdapterDescription, 
IStaticPropertyExtractor> {
 }
diff --git 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
 
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/IDataSinkMigrator.java
similarity index 90%
copy from 
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
copy to 
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/IDataSinkMigrator.java
index e30ab00c2..d1ce71f89 100644
--- 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
+++ 
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/IDataSinkMigrator.java
@@ -21,5 +21,5 @@ package org.apache.streampipes.extensions.api.migration;
 import 
org.apache.streampipes.extensions.api.extractor.IDataSinkParameterExtractor;
 import org.apache.streampipes.model.graph.DataSinkInvocation;
 
-public interface DataSinkMigrator extends IModelMigrator<DataSinkInvocation, 
IDataSinkParameterExtractor> {
+public interface IDataSinkMigrator extends IModelMigrator<DataSinkInvocation, 
IDataSinkParameterExtractor> {
 }
diff --git 
a/streampipes-model/src/main/java/org/apache/streampipes/model/configuration/SpCoreConfiguration.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/configuration/SpCoreConfiguration.java
index 3e3bd251d..6aad62a17 100644
--- 
a/streampipes-model/src/main/java/org/apache/streampipes/model/configuration/SpCoreConfiguration.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/configuration/SpCoreConfiguration.java
@@ -34,6 +34,7 @@ public class SpCoreConfiguration {
   private GeneralConfig generalConfig;
 
   private boolean isConfigured;
+  private SpCoreConfigurationStatus serviceStatus;
 
   private String assetDir;
   private String filesDir;
@@ -120,4 +121,12 @@ public class SpCoreConfiguration {
   public void setEmailTemplateConfig(EmailTemplateConfig emailTemplateConfig) {
     this.emailTemplateConfig = emailTemplateConfig;
   }
+
+  public SpCoreConfigurationStatus getServiceStatus() {
+    return this.serviceStatus;
+  }
+
+  public void setServiceStatus(SpCoreConfigurationStatus serviceStatus) {
+    this.serviceStatus = serviceStatus;
+  }
 }
diff --git 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/configuration/SpCoreConfigurationStatus.java
similarity index 72%
copy from 
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
copy to 
streampipes-model/src/main/java/org/apache/streampipes/model/configuration/SpCoreConfigurationStatus.java
index e30ab00c2..394b9a8ab 100644
--- 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/configuration/SpCoreConfigurationStatus.java
@@ -16,10 +16,10 @@
  *
  */
 
-package org.apache.streampipes.extensions.api.migration;
+package org.apache.streampipes.model.configuration;
 
-import 
org.apache.streampipes.extensions.api.extractor.IDataSinkParameterExtractor;
-import org.apache.streampipes.model.graph.DataSinkInvocation;
-
-public interface DataSinkMigrator extends IModelMigrator<DataSinkInvocation, 
IDataSinkParameterExtractor> {
+public enum SpCoreConfigurationStatus {
+  INSTALLING,
+  MIGRATING,
+  READY
 }
diff --git 
a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/migration/MigrationHelpers.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/migration/MigrationHelpers.java
index d518086e9..2343d41bb 100644
--- 
a/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/migration/MigrationHelpers.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/connect/adapter/migration/MigrationHelpers.java
@@ -40,6 +40,10 @@ public class MigrationHelpers {
     return adapter.get(REV).getAsString();
   }
 
+  public String getAppId(JsonObject adapter) {
+    return 
adapter.get("properties").getAsJsonObject().get(APP_ID).getAsString();
+  }
+
   public void updateType(JsonObject adapter,
                          String typeFieldName) {
     adapter.add(typeFieldName, new JsonPrimitive(AdapterModels.NEW_MODEL));
diff --git 
a/streampipes-model/src/main/java/org/apache/streampipes/model/extensions/svcdiscovery/SpServiceRegistration.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/extensions/svcdiscovery/SpServiceRegistration.java
index b1490adfa..8edd9932f 100644
--- 
a/streampipes-model/src/main/java/org/apache/streampipes/model/extensions/svcdiscovery/SpServiceRegistration.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/extensions/svcdiscovery/SpServiceRegistration.java
@@ -36,8 +36,8 @@ public class SpServiceRegistration {
   private int port;
   private List<SpServiceTag> tags;
   private String healthCheckPath;
-  private boolean healthy = true;
   private long firstTimeSeenUnhealthy = 0;
+  private SpServiceStatus status = SpServiceStatus.REGISTERED;
 
   public SpServiceRegistration() {
   }
@@ -133,14 +133,6 @@ public class SpServiceRegistration {
     this.rev = rev;
   }
 
-  public boolean isHealthy() {
-    return healthy;
-  }
-
-  public void setHealthy(boolean healthy) {
-    this.healthy = healthy;
-  }
-
   public String getScheme() {
     return scheme;
   }
@@ -168,4 +160,12 @@ public class SpServiceRegistration {
   public void setSvcType(String svcType) {
     this.svcType = svcType;
   }
+
+  public SpServiceStatus getStatus() {
+    return status;
+  }
+
+  public void setStatus(SpServiceStatus status) {
+    this.status = status;
+  }
 }
diff --git 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/extensions/svcdiscovery/SpServiceStatus.java
similarity index 72%
rename from 
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
rename to 
streampipes-model/src/main/java/org/apache/streampipes/model/extensions/svcdiscovery/SpServiceStatus.java
index e30ab00c2..56b388d16 100644
--- 
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/extensions/svcdiscovery/SpServiceStatus.java
@@ -16,10 +16,11 @@
  *
  */
 
-package org.apache.streampipes.extensions.api.migration;
+package org.apache.streampipes.model.extensions.svcdiscovery;
 
-import 
org.apache.streampipes.extensions.api.extractor.IDataSinkParameterExtractor;
-import org.apache.streampipes.model.graph.DataSinkInvocation;
-
-public interface DataSinkMigrator extends IModelMigrator<DataSinkInvocation, 
IDataSinkParameterExtractor> {
+public enum SpServiceStatus {
+  REGISTERED,
+  MIGRATING,
+  HEALTHY,
+  UNHEALTHY
 }
diff --git 
a/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataProcessorInvocation.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataProcessorInvocation.java
index 642134bf1..64aed3b49 100644
--- 
a/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataProcessorInvocation.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataProcessorInvocation.java
@@ -75,6 +75,7 @@ public class DataProcessorInvocation extends 
InvocableStreamPipesEntity implemen
   public DataProcessorInvocation(DataProcessorDescription sepa, String domId) {
     this(sepa);
     this.dom = domId;
+    this.serviceTagPrefix = SpServiceTagPrefix.DATA_PROCESSOR;
   }
 
   public DataProcessorInvocation() {
diff --git 
a/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataSinkInvocation.java
 
b/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataSinkInvocation.java
index 235ca33b4..4476c6651 100644
--- 
a/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataSinkInvocation.java
+++ 
b/streampipes-model/src/main/java/org/apache/streampipes/model/graph/DataSinkInvocation.java
@@ -59,6 +59,7 @@ public class DataSinkInvocation extends 
InvocableStreamPipesEntity {
   public DataSinkInvocation(DataSinkDescription sec, String domId) {
     this(sec);
     this.setDom(domId);
+    this.serviceTagPrefix = SpServiceTagPrefix.DATA_SINK;
   }
 
   public DataSinkInvocation() {
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
index 970d81a0d..cb2d28408 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/execution/endpoint/ExtensionsServiceEndpointGenerator.java
@@ -72,7 +72,7 @@ public class ExtensionsServiceEndpointGenerator {
 
   private String selectService() throws NoServiceEndpointsAvailableException {
     List<String> serviceEndpoints = getServiceEndpoints();
-    if (serviceEndpoints.size() > 0) {
+    if (!serviceEndpoints.isEmpty()) {
       return getServiceEndpoints().get(0);
     } else {
       LOG.error("Could not find any service endpoints for appId {}, serviceTag 
{}", appId,
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/CoreServiceStatusManager.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/CoreServiceStatusManager.java
new file mode 100644
index 000000000..3d2c407de
--- /dev/null
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/CoreServiceStatusManager.java
@@ -0,0 +1,59 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.manager.health;
+
+import org.apache.streampipes.model.configuration.SpCoreConfiguration;
+import org.apache.streampipes.model.configuration.SpCoreConfigurationStatus;
+import org.apache.streampipes.storage.api.ISpCoreConfigurationStorage;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class CoreServiceStatusManager {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(CoreServiceStatusManager.class);
+
+  private final ISpCoreConfigurationStorage storage;
+
+  public CoreServiceStatusManager(ISpCoreConfigurationStorage storage) {
+    this.storage = storage;
+  }
+
+  public boolean existsConfig() {
+    return storage.exists();
+  }
+
+  public boolean isCoreReady() {
+    return existsConfig() && storage.get().getServiceStatus() == 
SpCoreConfigurationStatus.READY;
+  }
+
+  public void updateCoreStatus(SpCoreConfigurationStatus status) {
+    var config = storage.get();
+    config.setServiceStatus(status);
+    storage.updateElement(config);
+    logService(config);
+  }
+
+  private void logService(SpCoreConfiguration coreConfig) {
+    LOG.info(
+        "Core is now in {} state",
+        coreConfig.getServiceStatus()
+    );
+  }
+}
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineHealthCheck.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineHealthCheck.java
index d8064055a..66a448be0 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineHealthCheck.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/PipelineHealthCheck.java
@@ -45,6 +45,8 @@ import java.util.Set;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Collectors;
 
+import static 
org.apache.streampipes.manager.pipeline.PipelineManager.getPipeline;
+
 public class PipelineHealthCheck implements Runnable {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(PipelineHealthCheck.class);
@@ -68,7 +70,7 @@ public class PipelineHealthCheck implements Runnable {
     pipelinesStats.setRunningPipelines(runningPipelines.size());
     pipelinesStats.setStoppedPipelines(pipelinesStats.getAllPipelines() - 
pipelinesStats.getRunningPipelines());
 
-    if (runningPipelines.size() > 0) {
+    if (!runningPipelines.isEmpty()) {
       Map<String, List<InvocableStreamPipesEntity>> endpointMap = 
generateEndpointMap();
       List<String> allRunningInstances = 
findRunningInstances(endpointMap.keySet());
 
@@ -115,15 +117,16 @@ public class PipelineHealthCheck implements Runnable {
           }
         });
         if (shouldUpdatePipeline.get()) {
-          if (failedInstances.size() > 0) {
-            pipeline.setHealthStatus(PipelineHealthStatus.FAILURE);
+          var currentPipeline = getPipeline(pipeline.getPipelineId());
+          if (!failedInstances.isEmpty()) {
+            currentPipeline.setHealthStatus(PipelineHealthStatus.FAILURE);
             pipelinesStats.failedIncrease();
-          } else if (recoveredInstances.size() > 0) {
-            pipeline.setHealthStatus(PipelineHealthStatus.REQUIRES_ATTENTION);
+          } else if (!recoveredInstances.isEmpty()) {
+            
currentPipeline.setHealthStatus(PipelineHealthStatus.REQUIRES_ATTENTION);
             pipelinesStats.attentionRequiredIncrease();
           }
-          pipeline.setPipelineNotifications(pipelineNotifications);
-          
StorageDispatcher.INSTANCE.getNoSqlStore().getPipelineStorageAPI().updatePipeline(pipeline);
+          currentPipeline.setPipelineNotifications(pipelineNotifications);
+          
StorageDispatcher.INSTANCE.getNoSqlStore().getPipelineStorageAPI().updatePipeline(currentPipeline);
         }
       });
       int healthNum = pipelinesStats.getRunningPipelines() - 
pipelinesStats.getFailedPipelines()
@@ -233,13 +236,11 @@ public class PipelineHealthCheck implements Runnable {
   }
 
   private List<Pipeline> getAllPipelines() {
-    List<Pipeline> allPipelines = StorageDispatcher
+    return StorageDispatcher
             .INSTANCE
             .getNoSqlStore()
             .getPipelineStorageAPI()
             .getAllPipelines();
-
-    return allPipelines;
   }
 
   private int getElementsCount(List<Pipeline> allPipelines){
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
index 71842d592..46d2792ea 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceHealthCheck.java
@@ -21,7 +21,7 @@ package org.apache.streampipes.manager.health;
 
 import org.apache.streampipes.manager.execution.ExtensionServiceExecutions;
 import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
-import org.apache.streampipes.storage.api.CRUDStorage;
+import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceStatus;
 import org.apache.streampipes.storage.management.StorageDispatcher;
 
 import org.apache.http.HttpStatus;
@@ -35,12 +35,13 @@ public class ServiceHealthCheck implements Runnable {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(ServiceHealthCheck.class);
 
-  private static final int MAX_UNHEALTHY_DURATION_BEFORE_REMOVAL_MS = 60000;
+  private static final int MAX_UNHEALTHY_DURATION_BEFORE_REMOVAL_MS = 20000;
 
-  private final CRUDStorage<String, SpServiceRegistration> storage;
+  private final ServiceRegistrationManager serviceRegistrationManager;
 
   public ServiceHealthCheck() {
-    this.storage = 
StorageDispatcher.INSTANCE.getNoSqlStore().getExtensionsServiceStorage();
+    var storage = 
StorageDispatcher.INSTANCE.getNoSqlStore().getExtensionsServiceStorage();
+    this.serviceRegistrationManager = new ServiceRegistrationManager(storage);
   }
 
   @Override
@@ -58,9 +59,8 @@ public class ServiceHealthCheck implements Runnable {
       if (response.returnResponse().getStatusLine().getStatusCode() != 
HttpStatus.SC_OK) {
         processUnhealthyService(service);
       } else {
-        if (!service.isHealthy()) {
-          service.setHealthy(true);
-          updateService(service);
+        if (service.getStatus() == SpServiceStatus.UNHEALTHY) {
+          serviceRegistrationManager.applyServiceStatus(service.getSvcId(), 
SpServiceStatus.HEALTHY);
         }
       }
     } catch (IOException e) {
@@ -69,15 +69,16 @@ public class ServiceHealthCheck implements Runnable {
   }
 
   private void processUnhealthyService(SpServiceRegistration service) {
-    if (service.isHealthy()) {
-      service.setHealthy(false);
-      service.setFirstTimeSeenUnhealthy(System.currentTimeMillis());
-      updateService(service);
+    if (service.getStatus() == SpServiceStatus.HEALTHY) {
+      serviceRegistrationManager.applyServiceStatus(
+          service.getSvcId(),
+          SpServiceStatus.UNHEALTHY,
+          System.currentTimeMillis());
     }
     if (shouldDeleteService(service)) {
       LOG.info("Removing service {} which has been unhealthy for more than {} 
seconds.",
           service.getSvcId(), MAX_UNHEALTHY_DURATION_BEFORE_REMOVAL_MS / 1000);
-      storage.deleteElement(service);
+      serviceRegistrationManager.removeService(service.getSvcId());
     }
   }
 
@@ -86,15 +87,11 @@ public class ServiceHealthCheck implements Runnable {
     return (currentTimeMillis - service.getFirstTimeSeenUnhealthy() > 
MAX_UNHEALTHY_DURATION_BEFORE_REMOVAL_MS);
   }
 
-  private void updateService(SpServiceRegistration service) {
-    storage.updateElement(service);
-  }
-
   private String makeHealthCheckUrl(SpServiceRegistration service) {
     return service.getServiceUrl() + service.getHealthCheckPath();
   }
 
   private List<SpServiceRegistration> getRegisteredServices() {
-    return storage.getAll();
+    return serviceRegistrationManager.getAllServices();
   }
 }
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceRegistrationManager.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceRegistrationManager.java
new file mode 100644
index 000000000..f29b39d11
--- /dev/null
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/health/ServiceRegistrationManager.java
@@ -0,0 +1,104 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.manager.health;
+
+import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
+import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceStatus;
+import org.apache.streampipes.storage.api.CRUDStorage;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.List;
+
+public class ServiceRegistrationManager {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(ServiceRegistrationManager.class);
+
+  private final CRUDStorage<String, SpServiceRegistration> storage;
+
+  public ServiceRegistrationManager(CRUDStorage<String, SpServiceRegistration> 
storage) {
+    this.storage = storage;
+  }
+
+  public void applyServiceStatus(String serviceId,
+                                 SpServiceStatus status,
+                                 long firstTimeSeenUnhealthy) {
+    var serviceRegistration = storage.getElementById(serviceId);
+    serviceRegistration.setFirstTimeSeenUnhealthy(firstTimeSeenUnhealthy);
+    applyServiceStatus(status, serviceRegistration);
+  }
+
+  public void applyServiceStatus(String serviceId,
+                                 SpServiceStatus status) {
+    var serviceRegistration = storage.getElementById(serviceId);
+    applyServiceStatus(status, serviceRegistration);
+  }
+
+  private void applyServiceStatus(SpServiceStatus status,
+                                  SpServiceRegistration serviceRegistration) {
+    serviceRegistration.setStatus(status);
+    storage.updateElement(serviceRegistration);
+    logService(serviceRegistration);
+  }
+
+  public void addService(SpServiceRegistration serviceRegistration,
+                         SpServiceStatus status) {
+    serviceRegistration.setStatus(status);
+    storage.createElement(serviceRegistration);
+    logService(serviceRegistration);
+  }
+
+  public List<SpServiceRegistration> getAllServices() {
+    return storage.getAll();
+  }
+
+  public SpServiceRegistration getService(String serviceId) {
+    return storage.getElementById(serviceId);
+  }
+
+  public boolean isAnyServiceMigrating() {
+    return storage.getAll()
+        .stream()
+        .anyMatch(service -> service.getStatus() == SpServiceStatus.MIGRATING);
+  }
+
+  public void removeService(String serviceId) {
+    var serviceRegistration = storage.getElementById(serviceId);
+    storage.deleteElement(serviceRegistration);
+    LOG.info(
+        "Service {} (id={}) has been removed",
+        serviceRegistration.getSvcGroup(),
+        serviceRegistration.getSvcId())
+    ;
+  }
+
+  public SpServiceStatus getServiceStatus(String serviceId) {
+    return storage.getElementById(serviceId).getStatus();
+  }
+
+  private void logService(SpServiceRegistration serviceRegistration) {
+    LOG.info(
+        "Service {} (id={}) is now in {} state",
+        serviceRegistration.getSvcGroup(),
+        serviceRegistration.getSvcId(),
+        serviceRegistration.getStatus()
+    );
+  }
+}
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AdapterDescriptionMigration093.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AdapterDescriptionMigration093.java
new file mode 100644
index 000000000..d10d27999
--- /dev/null
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AdapterDescriptionMigration093.java
@@ -0,0 +1,73 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.manager.migration;
+
+import org.apache.streampipes.commons.exceptions.SepaParseException;
+import org.apache.streampipes.manager.endpoint.HttpJsonParser;
+import org.apache.streampipes.manager.operations.Operations;
+import org.apache.streampipes.manager.util.AuthTokenUtils;
+import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
+import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceTagPrefix;
+import org.apache.streampipes.storage.api.IAdapterStorage;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.net.URI;
+
+import static 
org.apache.streampipes.manager.migration.MigrationUtils.getRequestUrl;
+
+public class AdapterDescriptionMigration093 extends AbstractMigrationManager {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(AdapterDescriptionMigration093.class);
+
+  private final IAdapterStorage adapterDescriptionStorage;
+
+  public AdapterDescriptionMigration093(IAdapterStorage 
adapterDescriptionStorage) {
+    this.adapterDescriptionStorage = adapterDescriptionStorage;
+  }
+
+  public void reinstallAdapters(SpServiceRegistration extensionsServiceConfig) 
{
+    var migrationProvider = AdapterDescriptionMigration093Provider.INSTANCE;
+    if (migrationProvider.hasAppIdsToReinstall()) {
+      var appIdsToReinstall = migrationProvider.getAppIdsToReinstall();
+      var serviceUrl = extensionsServiceConfig.getServiceUrl();
+      extensionsServiceConfig.getTags()
+          .stream()
+          .filter(tag -> tag.getPrefix() == SpServiceTagPrefix.ADAPTER)
+          .filter(tag -> appIdsToReinstall.contains(tag.getValue()))
+          .forEach(tag -> {
+            var appId = tag.getValue();
+            try {
+              if 
(adapterDescriptionStorage.getAdaptersByAppId(appId).isEmpty()) {
+                var requestUrl = getRequestUrl(SpServiceTagPrefix.ADAPTER, 
appId, serviceUrl);
+                var entityPayload = 
HttpJsonParser.getContentFromUrl(URI.create(requestUrl));
+                Operations.verifyAndAddElement(
+                    entityPayload,
+                    AuthTokenUtils.getAuthTokenForCurrentUser(),
+                    true);
+              }
+            } catch (IOException | SepaParseException e) {
+              LOG.warn("Could not reinstall adapter description {}", appId);
+            }
+          });
+    }
+  }
+}
diff --git 
a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/ISpCoreConfigurationStorage.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AdapterDescriptionMigration093Provider.java
similarity index 60%
copy from 
streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/ISpCoreConfigurationStorage.java
copy to 
streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AdapterDescriptionMigration093Provider.java
index 829dc09ce..8ad025d5c 100644
--- 
a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/ISpCoreConfigurationStorage.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/AdapterDescriptionMigration093Provider.java
@@ -16,21 +16,30 @@
  *
  */
 
-package org.apache.streampipes.storage.api;
-
-import org.apache.streampipes.model.configuration.SpCoreConfiguration;
+package org.apache.streampipes.manager.migration;
 
+import java.util.ArrayList;
 import java.util.List;
 
-public interface ISpCoreConfigurationStorage {
+public enum AdapterDescriptionMigration093Provider {
+
+  INSTANCE;
 
-  List<SpCoreConfiguration> getAll();
+  private final List<String> appIdsToReinstall;
 
-  void createElement(SpCoreConfiguration element);
+  AdapterDescriptionMigration093Provider() {
+    this.appIdsToReinstall = new ArrayList<>();
+  }
 
-  SpCoreConfiguration get();
+  public void addAppId(String appId) {
+    this.appIdsToReinstall.add(appId);
+  }
 
-  SpCoreConfiguration updateElement(SpCoreConfiguration element);
+  public List<String> getAppIdsToReinstall() {
+    return appIdsToReinstall;
+  }
 
-  void deleteElement();
+  public boolean hasAppIdsToReinstall() {
+    return !appIdsToReinstall.isEmpty();
+  }
 }
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/PipelineElementMigrationManager.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/PipelineElementMigrationManager.java
index 28338e1ac..166d24c8e 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/PipelineElementMigrationManager.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/migration/PipelineElementMigrationManager.java
@@ -38,6 +38,7 @@ import org.slf4j.LoggerFactory;
 
 import java.util.ArrayList;
 import java.util.List;
+import java.util.stream.Stream;
 
 import static 
org.apache.streampipes.manager.migration.MigrationUtils.getApplicableMigration;
 
@@ -60,77 +61,89 @@ public class PipelineElementMigrationManager extends 
AbstractMigrationManager im
   @Override
   public void handleMigrations(SpServiceRegistration extensionsServiceConfig,
                                List<ModelMigratorConfig> migrationConfigs) {
+    if (!migrationConfigs.isEmpty()) {
+      LOG.info("Updating pipeline element descriptions by replacement...");
+      updateDescriptions(migrationConfigs, 
extensionsServiceConfig.getServiceUrl());
+      LOG.info("Pipeline element descriptions are up to date.");
 
-    LOG.info("Updating pipeline element descriptions by replacement...");
-    updateDescriptions(migrationConfigs, 
extensionsServiceConfig.getServiceUrl());
-    LOG.info("Pipeline element descriptions are up to date.");
-
-    LOG.info("Received {} pipeline element migrations from extension service 
{}.",
-        migrationConfigs.size(),
-        extensionsServiceConfig.getServiceUrl());
-    var availablePipelines = pipelineStorage.getAllPipelines();
-    if (!availablePipelines.isEmpty()) {
-      LOG.info("Found {} available pipelines. Checking pipelines for 
applicable migrations...",
-          availablePipelines.size()
-      );
-    }
+      LOG.info("Received {} pipeline element migrations from extension service 
{}.",
+          migrationConfigs.size(),
+          extensionsServiceConfig.getServiceUrl());
+      var availablePipelines = pipelineStorage.getAllPipelines();
+      if (!availablePipelines.isEmpty()) {
+        LOG.info("Found {} available pipelines. Checking pipelines for 
applicable migrations...",
+            availablePipelines.size()
+        );
+      }
 
-    for (var pipeline : availablePipelines) {
-      List<MigrationResult<?>> failedMigrations = new ArrayList<>();
-
-      var migratedDataProcessors = pipeline.getSepas()
-          .stream()
-          .map(processor -> {
-            if (getApplicableMigration(processor, 
migrationConfigs).isPresent()) {
-              return migratePipelineElement(
-                  processor,
-                  migrationConfigs,
-                  String.format("%s/%s/processor",
-                      extensionsServiceConfig.getServiceUrl(),
-                      MIGRATION_ENDPOINT
-                  ),
-                  failedMigrations
-              );
-            } else {
-              LOG.info("No migration applicable for data processor '{}'.", 
processor.getElementId());
-              return processor;
-            }
-          })
-          .toList();
-      pipeline.setSepas(migratedDataProcessors);
-
-      var migratedDataSinks = pipeline.getActions()
-          .stream()
-          .map(sink -> {
-            if (getApplicableMigration(sink, migrationConfigs).isPresent()) {
-              return migratePipelineElement(
-                  sink,
-                  migrationConfigs,
-                  String.format("%s/%s/sink",
-                      extensionsServiceConfig.getServiceUrl(),
-                      MIGRATION_ENDPOINT
-                  ),
-                  failedMigrations
-              );
-            } else {
-              LOG.info("No migration applicable for data sink '{}'.", 
sink.getElementId());
-              return sink;
-            }
-          })
-          .toList();
-      pipeline.setActions(migratedDataSinks);
-
-      pipelineStorage.updatePipeline(pipeline);
-
-      if (failedMigrations.isEmpty()) {
-        LOG.info("Migration for pipeline successfully completed.");
-      } else {
-        // pass most recent version of pipeline
-        
handleFailedMigrations(pipelineStorage.getPipeline(pipeline.getPipelineId()), 
failedMigrations);
+      for (var pipeline : availablePipelines) {
+        if (shouldMigratePipeline(pipeline, migrationConfigs)) {
+          List<MigrationResult<?>> failedMigrations = new ArrayList<>();
+
+          var migratedDataProcessors = pipeline.getSepas()
+              .stream()
+              .map(processor -> {
+                if (getApplicableMigration(processor, 
migrationConfigs).isPresent()) {
+                  return migratePipelineElement(
+                      processor,
+                      migrationConfigs,
+                      String.format("%s/%s/processor",
+                          extensionsServiceConfig.getServiceUrl(),
+                          MIGRATION_ENDPOINT
+                      ),
+                      failedMigrations
+                  );
+                } else {
+                  LOG.info("No migration applicable for data processor '{}'.", 
processor.getElementId());
+                  return processor;
+                }
+              })
+              .toList();
+          pipeline.setSepas(migratedDataProcessors);
+
+          var migratedDataSinks = pipeline.getActions()
+              .stream()
+              .map(sink -> {
+                if (getApplicableMigration(sink, 
migrationConfigs).isPresent()) {
+                  return migratePipelineElement(
+                      sink,
+                      migrationConfigs,
+                      String.format("%s/%s/sink",
+                          extensionsServiceConfig.getServiceUrl(),
+                          MIGRATION_ENDPOINT
+                      ),
+                      failedMigrations
+                  );
+                } else {
+                  LOG.info("No migration applicable for data sink '{}'.", 
sink.getElementId());
+                  return sink;
+                }
+              })
+              .toList();
+          pipeline.setActions(migratedDataSinks);
+
+          pipelineStorage.updatePipeline(pipeline);
+
+          if (failedMigrations.isEmpty()) {
+            LOG.info("Migration for pipeline successfully completed.");
+          } else {
+            // pass most recent version of pipeline
+            
handleFailedMigrations(pipelineStorage.getPipeline(pipeline.getPipelineId()), 
failedMigrations);
+          }
+        }
       }
+    } else {
+      LOG.info("No pipeline element migrations to perform");
     }
   }
 
+  private boolean shouldMigratePipeline(Pipeline pipeline,
+                                        List<ModelMigratorConfig> 
migrationConfigs) {
+    return Stream
+        .concat(pipeline.getSepas().stream(), pipeline.getActions().stream())
+        .anyMatch(element -> getApplicableMigration(element, 
migrationConfigs).isPresent());
+  }
+
   /**
    * Takes care about the failed migrations of pipeline elements.
    * This includes the following steps:
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java
index 187d677ef..b5e0a8bd3 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/SpCoreConfigurationStep.java
@@ -19,14 +19,22 @@
 package org.apache.streampipes.manager.setup;
 
 import org.apache.streampipes.model.configuration.DefaultSpCoreConfiguration;
+import org.apache.streampipes.model.configuration.SpCoreConfigurationStatus;
 import org.apache.streampipes.storage.management.StorageDispatcher;
 
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
 public class SpCoreConfigurationStep extends InstallationStep {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(SpCoreConfigurationStep.class);
+
   @Override
   public void install() {
     var coreCfg = new DefaultSpCoreConfiguration().make();
-
+    coreCfg.setServiceStatus(SpCoreConfigurationStatus.INSTALLING);
     
StorageDispatcher.INSTANCE.getNoSqlStore().getSpCoreConfigurationStorage().createElement(coreCfg);
+    LOG.info("Core is now in {} state", coreCfg.getServiceStatus());
     new StreamPipesEnvChecker().updateEnvironmentVariables();
   }
 
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
index c4fc1bba0..70e5a8aa5 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/setup/StreamPipesEnvChecker.java
@@ -52,7 +52,7 @@ public class StreamPipesEnvChecker {
         .getNoSqlStore()
         .getSpCoreConfigurationStorage();
 
-    if (configStorage.getAll().size() > 0) {
+    if (configStorage.exists()) {
       this.coreConfig = configStorage.get();
 
       LOG.info("Checking and updating environment variables...");
diff --git 
a/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/migration/DataSinkMigrationResource.java
 
b/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/migration/DataSinkMigrationResource.java
index 98569ffae..6d1844743 100644
--- 
a/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/migration/DataSinkMigrationResource.java
+++ 
b/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/migration/DataSinkMigrationResource.java
@@ -19,7 +19,7 @@
 package org.apache.streampipes.rest.extensions.migration;
 
 import 
org.apache.streampipes.extensions.api.extractor.IDataSinkParameterExtractor;
-import org.apache.streampipes.extensions.api.migration.DataSinkMigrator;
+import org.apache.streampipes.extensions.api.migration.IDataSinkMigrator;
 import org.apache.streampipes.model.extensions.migration.MigrationRequest;
 import org.apache.streampipes.model.graph.DataSinkInvocation;
 import org.apache.streampipes.rest.security.AuthConstants;
@@ -42,7 +42,7 @@ import jakarta.ws.rs.core.Response;
 public class DataSinkMigrationResource extends MigrateExtensionsResource<
         DataSinkInvocation,
         IDataSinkParameterExtractor,
-        DataSinkMigrator
+    IDataSinkMigrator
         > {
   @POST
   @Consumes(MediaType.APPLICATION_JSON)
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java
index 13c87875f..6b28afe0a 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/MigrationResource.java
@@ -18,9 +18,14 @@
 
 package org.apache.streampipes.rest.impl.admin;
 
+import org.apache.streampipes.config.backend.BackendConfig;
 import 
org.apache.streampipes.connect.management.management.AdapterMigrationManager;
+import org.apache.streampipes.manager.health.CoreServiceStatusManager;
+import org.apache.streampipes.manager.health.ServiceRegistrationManager;
+import org.apache.streampipes.manager.migration.AdapterDescriptionMigration093;
 import 
org.apache.streampipes.manager.migration.PipelineElementMigrationManager;
 import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
+import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceStatus;
 import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceTagPrefix;
 import org.apache.streampipes.model.migration.ModelMigratorConfig;
 import 
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
@@ -36,6 +41,8 @@ import io.swagger.v3.oas.annotations.Parameter;
 import io.swagger.v3.oas.annotations.enums.ParameterIn;
 import io.swagger.v3.oas.annotations.responses.ApiResponse;
 import org.apache.http.HttpStatus;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 import org.springframework.security.access.prepost.PreAuthorize;
 import org.springframework.stereotype.Component;
 
@@ -53,9 +60,12 @@ import java.util.List;
 @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
 public class MigrationResource extends AbstractAuthGuardedRestResource {
 
+  private static final Logger LOG = 
LoggerFactory.getLogger(MigrationResource.class);
+
   private final CRUDStorage<String, SpServiceRegistration> 
extensionsServiceStorage =
       getNoSqlStorage().getExtensionsServiceStorage();
 
+  private final IAdapterStorage adapterDescriptionStorage = 
getNoSqlStorage().getAdapterDescriptionStorage();
   private final IAdapterStorage adapterStorage = 
getNoSqlStorage().getAdapterInstanceStorage();
 
   private final IDataProcessorStorage dataProcessorStorage = 
getNoSqlStorage().getDataProcessorStorage();
@@ -63,6 +73,10 @@ public class MigrationResource extends 
AbstractAuthGuardedRestResource {
   private final IDataSinkStorage dataSinkStorage = 
getNoSqlStorage().getDataSinkStorage();
   private final IPipelineStorage pipelineStorage = 
getNoSqlStorage().getPipelineStorageAPI();
 
+  private final CoreServiceStatusManager coreServiceStatusManager = new 
CoreServiceStatusManager(
+      getNoSqlStorage().getSpCoreConfigurationStorage()
+  );
+
   @POST
   @Path("{serviceId}")
   @Consumes(MediaType.APPLICATION_JSON)
@@ -88,24 +102,42 @@ public class MigrationResource extends 
AbstractAuthGuardedRestResource {
       )
       List<ModelMigratorConfig> migrationConfigs) {
 
-    var extensionsServiceConfig = 
extensionsServiceStorage.getElementById(serviceId);
-    var adapterMigrations = filterConfigs(migrationConfigs, 
List.of(SpServiceTagPrefix.ADAPTER));
-    var pipelineElementMigrations = filterConfigs(
-        migrationConfigs,
-        List.of(SpServiceTagPrefix.DATA_PROCESSOR, 
SpServiceTagPrefix.DATA_SINK)
-    );
-
-    new 
AdapterMigrationManager(adapterStorage).handleMigrations(extensionsServiceConfig,
 adapterMigrations);
-    new PipelineElementMigrationManager(
-        pipelineStorage,
-        dataProcessorStorage,
-        dataSinkStorage)
-        .handleMigrations(extensionsServiceConfig, pipelineElementMigrations);
+    var serviceManager = new 
ServiceRegistrationManager(extensionsServiceStorage);
+    var extensionsServiceConfig = serviceManager.getService(serviceId);
+    if (BackendConfig.INSTANCE.isConfigured()) {
+      new 
AdapterDescriptionMigration093(adapterDescriptionStorage).reinstallAdapters(extensionsServiceConfig);
+      if (!migrationConfigs.isEmpty()) {
+        if (serviceManager.isAnyServiceMigrating() || !isCoreReady()) {
+          LOG.info("Refusing migration request since precondition is not 
met.");
+          return Response.status(HttpStatus.SC_CONFLICT).build();
+        } else {
+          serviceManager.applyServiceStatus(serviceId, 
SpServiceStatus.MIGRATING);
+          var adapterMigrations = filterConfigs(migrationConfigs, 
List.of(SpServiceTagPrefix.ADAPTER));
+          var pipelineElementMigrations = filterConfigs(
+              migrationConfigs,
+              List.of(SpServiceTagPrefix.DATA_PROCESSOR, 
SpServiceTagPrefix.DATA_SINK)
+          );
+
+          new 
AdapterMigrationManager(adapterStorage).handleMigrations(extensionsServiceConfig,
 adapterMigrations);
+          new PipelineElementMigrationManager(
+              pipelineStorage,
+              dataProcessorStorage,
+              dataSinkStorage)
+              .handleMigrations(extensionsServiceConfig, 
pipelineElementMigrations);
+        }
+      }
+    }
+    new ServiceRegistrationManager(extensionsServiceStorage)
+        .applyServiceStatus(extensionsServiceConfig.getSvcId(), 
SpServiceStatus.HEALTHY);
     return ok();
   }
 
+  private boolean isCoreReady() {
+    return coreServiceStatusManager.isCoreReady();
+  }
+
   private List<ModelMigratorConfig> filterConfigs(List<ModelMigratorConfig> 
migrationConfigs,
-                                                       
List<SpServiceTagPrefix> modelTypes) {
+                                                  List<SpServiceTagPrefix> 
modelTypes) {
     return migrationConfigs
         .stream()
         .filter(config -> modelTypes.stream().anyMatch(modelType -> modelType 
== config.modelType()))
diff --git 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ServiceRegistrationResource.java
 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ServiceRegistrationResource.java
index f89cf34ac..d97653d86 100644
--- 
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ServiceRegistrationResource.java
+++ 
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/admin/ServiceRegistrationResource.java
@@ -18,11 +18,15 @@
 
 package org.apache.streampipes.rest.impl.admin;
 
+import org.apache.streampipes.manager.health.ServiceRegistrationManager;
 import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
+import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceStatus;
 import 
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
 import org.apache.streampipes.rest.security.AuthConstants;
 import org.apache.streampipes.storage.api.CRUDStorage;
 
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 import org.springframework.security.access.prepost.PreAuthorize;
 import org.springframework.stereotype.Component;
 
@@ -40,6 +44,8 @@ import jakarta.ws.rs.core.Response;
 @PreAuthorize(AuthConstants.IS_ADMIN_ROLE)
 public class ServiceRegistrationResource extends 
AbstractAuthGuardedRestResource {
 
+  private static final Logger LOG = 
LoggerFactory.getLogger(ServiceRegistrationResource.class);
+
   private final CRUDStorage<String, SpServiceRegistration> 
extensionsServiceStorage =
       getNoSqlStorage().getExtensionsServiceStorage();
 
@@ -52,7 +58,8 @@ public class ServiceRegistrationResource extends 
AbstractAuthGuardedRestResource
   @POST
   @Consumes(MediaType.APPLICATION_JSON)
   public Response registerService(SpServiceRegistration serviceRegistration) {
-    extensionsServiceStorage.createElement(serviceRegistration);
+    new ServiceRegistrationManager(extensionsServiceStorage)
+        .addService(serviceRegistration, SpServiceStatus.REGISTERED);
     return ok();
   }
 
@@ -60,8 +67,7 @@ public class ServiceRegistrationResource extends 
AbstractAuthGuardedRestResource
   @Path("/{serviceId}")
   public Response unregisterService(@PathParam("serviceId") String serviceId) {
     try {
-      var serviceRegistration = 
extensionsServiceStorage.getElementById(serviceId);
-      extensionsServiceStorage.deleteElement(serviceRegistration);
+      new 
ServiceRegistrationManager(extensionsServiceStorage).removeService(serviceId);
       return ok();
     } catch (IllegalArgumentException e) {
       return badRequest("Could not find registered service with id " + 
serviceId);
diff --git 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
index 9355c6b89..244281674 100644
--- 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
+++ 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/PostStartupTask.java
@@ -44,13 +44,13 @@ public class PostStartupTask implements Runnable {
   private static final int MAX_PIPELINE_START_RETRIES = 3;
   private static final int WAIT_TIME_AFTER_FAILURE_IN_SECONDS = 10;
 
-  private final List<Pipeline> allPipelines;
+  private final IPipelineStorage pipelineStorage;
   private final Map<String, Integer> failedPipelines = new HashMap<>();
   private final ScheduledExecutorService executorService;
   private final WorkerAdministrationManagement workerAdministrationManagement;
 
-  public PostStartupTask(List<Pipeline> allPipelines) {
-    this.allPipelines = allPipelines;
+  public PostStartupTask(IPipelineStorage pipelineStorage) {
+    this.pipelineStorage = pipelineStorage;
     this.executorService = Executors.newSingleThreadScheduledExecutor();
     this.workerAdministrationManagement = new WorkerAdministrationManagement();
   }
@@ -76,6 +76,7 @@ public class PostStartupTask implements Runnable {
   }
 
   private void startAllPreviouslyStoppedPipelines() {
+    var allPipelines = pipelineStorage.getAllPipelines();
     LOG.info("Checking for orphaned pipelines...");
     List<Pipeline> orphanedPipelines = allPipelines
         .stream()
diff --git 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
index e03d3291d..49b218f73 100644
--- 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
+++ 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/StreamPipesCoreApplication.java
@@ -18,6 +18,8 @@
 package org.apache.streampipes.service.core;
 
 import org.apache.streampipes.config.backend.BackendConfig;
+import org.apache.streampipes.connect.management.health.AdapterHealthCheck;
+import org.apache.streampipes.manager.health.CoreServiceStatusManager;
 import org.apache.streampipes.manager.health.PipelineHealthCheck;
 import org.apache.streampipes.manager.health.ServiceHealthCheck;
 import 
org.apache.streampipes.manager.monitoring.pipeline.ExtensionsServiceLogExecutor;
@@ -30,6 +32,7 @@ import 
org.apache.streampipes.messaging.kafka.SpKafkaProtocolFactory;
 import org.apache.streampipes.messaging.mqtt.SpMqttProtocolFactory;
 import org.apache.streampipes.messaging.nats.SpNatsProtocolFactory;
 import org.apache.streampipes.messaging.pulsar.SpPulsarProtocolFactory;
+import org.apache.streampipes.model.configuration.SpCoreConfigurationStatus;
 import org.apache.streampipes.model.pipeline.Pipeline;
 import org.apache.streampipes.model.pipeline.PipelineOperationStatus;
 import org.apache.streampipes.rest.security.SpPermissionEvaluator;
@@ -37,6 +40,7 @@ import 
org.apache.streampipes.service.base.BaseNetworkingConfig;
 import org.apache.streampipes.service.base.StreamPipesServiceBase;
 import org.apache.streampipes.service.core.migrations.MigrationsHandler;
 import org.apache.streampipes.storage.api.IPipelineStorage;
+import org.apache.streampipes.storage.api.ISpCoreConfigurationStorage;
 import org.apache.streampipes.storage.couchdb.utils.CouchDbViewGenerator;
 import org.apache.streampipes.storage.management.StorageDispatcher;
 
@@ -72,11 +76,13 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
   private static final int LOG_FETCH_INTERVAL = 60;
   private static final TimeUnit LOG_FETCH_UNIT = TimeUnit.SECONDS;
 
-  private static final int HEALTH_CHECK_INTERVAL = 60;
+  private static final int HEALTH_CHECK_INTERVAL = 30;
   private static final TimeUnit HEALTH_CHECK_UNIT = TimeUnit.SECONDS;
 
-  private static final int SERVICE_HEALTH_CHECK_INTERVAL = 60;
-  private static final TimeUnit SERVICE_HEALTH_CHECK_UNIT = TimeUnit.SECONDS;
+  private final ISpCoreConfigurationStorage coreConfigStorage = 
StorageDispatcher.INSTANCE
+      .getNoSqlStore().getSpCoreConfigurationStorage();
+
+  private final CoreServiceStatusManager coreStatusManager = new 
CoreServiceStatusManager(coreConfigStorage);
 
   public static void main(String[] args) {
     StreamPipesCoreApplication application = new StreamPipesCoreApplication();
@@ -109,9 +115,7 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
   @PostConstruct
   public void init() {
     var executorService = Executors.newSingleThreadScheduledExecutor();
-    var healthCheckExecutorService = 
Executors.newSingleThreadScheduledExecutor();
     var logCheckExecutorService = Executors.newSingleThreadScheduledExecutor();
-    var serviceHealthCheckExecutorService = 
Executors.newSingleThreadScheduledExecutor();
 
     new StreamPipesEnvChecker().updateEnvironmentVariables();
     new CouchDbViewGenerator().createGenericDatabaseIfNotExists();
@@ -119,22 +123,20 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
     if (!isConfigured()) {
       doInitialSetup();
     } else {
+      // Check needs to be present since core configuration is part of 
migration
+      if (coreConfigStorage.exists()) {
+        
coreStatusManager.updateCoreStatus(SpCoreConfigurationStatus.MIGRATING);
+      }
       new MigrationsHandler().performMigrations();
     }
+    coreStatusManager.updateCoreStatus(SpCoreConfigurationStatus.READY);
 
-    executorService.schedule(new PostStartupTask(getAllPipelines()), 10, 
TimeUnit.SECONDS);
+    executorService.schedule(new PostStartupTask(getPipelineStorage()), 10, 
TimeUnit.SECONDS);
 
-    LOG.info("Service health check will run every {} seconds", 
SERVICE_HEALTH_CHECK_INTERVAL);
-    serviceHealthCheckExecutorService.scheduleAtFixedRate(new 
ServiceHealthCheck(),
-        SERVICE_HEALTH_CHECK_INTERVAL,
-        SERVICE_HEALTH_CHECK_INTERVAL,
-        SERVICE_HEALTH_CHECK_UNIT);
-
-    LOG.info("Pipeline health check will run every {} seconds", 
HEALTH_CHECK_INTERVAL);
-    healthCheckExecutorService.scheduleAtFixedRate(new PipelineHealthCheck(),
-        HEALTH_CHECK_INTERVAL,
-        HEALTH_CHECK_INTERVAL,
-        HEALTH_CHECK_UNIT);
+    scheduleHealthChecks(List.of(
+        new ServiceHealthCheck(),
+        new PipelineHealthCheck(),
+        new AdapterHealthCheck()));
 
     LOG.info("Extensions logs will be fetched every {} seconds", 
LOG_FETCH_INTERVAL);
     logCheckExecutorService.scheduleAtFixedRate(new 
ExtensionsServiceLogExecutor(),
@@ -143,6 +145,21 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
         LOG_FETCH_UNIT);
   }
 
+  private void scheduleHealthChecks(List<Runnable> checks) {
+    var healthCheckExecutorService = 
Executors.newSingleThreadScheduledExecutor();
+    checks.forEach(check -> {
+      LOG.info(
+          "Health check {} configured to run every {} {}",
+          check.getClass().getCanonicalName(),
+          HEALTH_CHECK_INTERVAL,
+          HEALTH_CHECK_UNIT);
+      healthCheckExecutorService.scheduleAtFixedRate(check,
+          HEALTH_CHECK_INTERVAL,
+          HEALTH_CHECK_INTERVAL,
+          HEALTH_CHECK_UNIT);
+    });
+  }
+
   private boolean isConfigured() {
     return BackendConfig.INSTANCE.isConfigured();
   }
@@ -163,8 +180,6 @@ public class StreamPipesCoreApplication extends 
StreamPipesServiceBase {
     }
   }
 
-
-
   @PreDestroy
   public void onExit() {
     LOG.info("Shutting down StreamPipes...");
diff --git 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/AdapterMigration.java
 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/AdapterMigration.java
index 94d2d7eec..29fd205ab 100644
--- 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/AdapterMigration.java
+++ 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/AdapterMigration.java
@@ -18,6 +18,7 @@
 
 package org.apache.streampipes.service.core.migrations.v093;
 
+import 
org.apache.streampipes.manager.migration.AdapterDescriptionMigration093Provider;
 import org.apache.streampipes.model.connect.adapter.migration.MigrationHelpers;
 import 
org.apache.streampipes.model.connect.adapter.migration.utils.AdapterModels;
 import org.apache.streampipes.service.core.migrations.Migration;
@@ -45,7 +46,7 @@ public class AdapterMigration implements Migration {
   private final CouchDbClient adapterInstanceClient;
   private final CouchDbClient adapterDescriptionClient;
   private final List<JsonObject> adaptersToMigrate;
-  private final List<JsonObject> adapterDescriptionsToMigrate;
+  private final List<JsonObject> adapterDescriptionsToRemove;
 
   private final MigrationHelpers helpers;
 
@@ -54,7 +55,7 @@ public class AdapterMigration implements Migration {
     this.adapterInstanceClient = Utils.getCouchDbAdapterInstanceClient();
     this.adapterDescriptionClient = Utils.getCouchDbAdapterDescriptionClient();
     this.adaptersToMigrate = new ArrayList<>();
-    this.adapterDescriptionsToMigrate = new ArrayList<>();
+    this.adapterDescriptionsToRemove = new ArrayList<>();
     this.helpers = new MigrationHelpers();
   }
 
@@ -64,9 +65,9 @@ public class AdapterMigration implements Migration {
     var adapterDescriptionUri = getAllDocsUri(adapterDescriptionClient);
 
     findDocsToMigrate(adapterInstanceClient, adapterInstanceUri, 
adaptersToMigrate);
-    findDocsToMigrate(adapterDescriptionClient, adapterDescriptionUri, 
adapterDescriptionsToMigrate);
+    findDocsToMigrate(adapterDescriptionClient, adapterDescriptionUri, 
adapterDescriptionsToRemove);
 
-    return !adaptersToMigrate.isEmpty() || 
!adapterDescriptionsToMigrate.isEmpty();
+    return !adaptersToMigrate.isEmpty() || 
!adapterDescriptionsToRemove.isEmpty();
   }
 
   private void findDocsToMigrate(CouchDbClient adapterClient,
@@ -89,15 +90,19 @@ public class AdapterMigration implements Migration {
   public void executeMigration() {
     var adapterInstanceBackupClient = 
Utils.getCouchDbAdapterInstanceBackupClient();
 
-    adapterDescriptionsToMigrate.forEach(ad -> {
+    LOG.info("Deleting {} adapter descriptions, which will be regenerated 
after migration",
+        adapterDescriptionsToRemove.size());
+
+    adapterDescriptionsToRemove.forEach(ad -> {
+      String docId = helpers.getDocId(ad);
       var adapterType = ad.get("type").getAsString();
-      var appId = ad.get("appId");
-      if (isSetAdapter(adapterType)) {
-        LOG.info("Deleting adapter description data set {}", appId);
-        adapterDescriptionClient.remove(helpers.getDocId(ad), 
helpers.getRev(ad));
-      } else {
-        LOG.info("Migrating adapter description {} to new adapter model", 
appId);
-        getAdapterMigrator(adapterType).migrate(adapterDescriptionClient, ad);
+      String rev = helpers.getRev(ad);
+      String appId = helpers.getAppId(ad);
+      if (!isSetAdapter(adapterType)) {
+        AdapterDescriptionMigration093Provider.INSTANCE.addAppId(appId);
+      }
+      if (docId != null && rev != null) {
+        adapterDescriptionClient.remove(docId, rev);
       }
     });
 
diff --git 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/ConsulConfigMigration.java
 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/ConsulConfigMigration.java
index 5da0d0bef..53857b12c 100644
--- 
a/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/ConsulConfigMigration.java
+++ 
b/streampipes-service-core/src/main/java/org/apache/streampipes/service/core/migrations/v093/ConsulConfigMigration.java
@@ -22,6 +22,7 @@ import org.apache.streampipes.config.backend.BackendConfig;
 import org.apache.streampipes.config.backend.BackendConfigKeys;
 import org.apache.streampipes.model.configuration.DefaultMessagingSettings;
 import org.apache.streampipes.model.configuration.SpCoreConfiguration;
+import org.apache.streampipes.model.configuration.SpCoreConfigurationStatus;
 import org.apache.streampipes.service.core.migrations.Migration;
 import org.apache.streampipes.storage.api.ISpCoreConfigurationStorage;
 import org.apache.streampipes.storage.management.StorageDispatcher;
@@ -71,6 +72,7 @@ public class ConsulConfigMigration implements Migration {
     newConf.setFilesDir(currConf.getFilesDir());
 
     newConf.setMessagingSettings(messagingSettings);
+    newConf.setServiceStatus(SpCoreConfigurationStatus.MIGRATING);
     storage.createElement(newConf);
   }
 
diff --git 
a/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
 
b/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
index 8170c36ad..0455dc1d8 100644
--- 
a/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
+++ 
b/streampipes-service-discovery/src/main/java/org/apache/streampipes/svcdiscovery/SpServiceDiscoveryCore.java
@@ -19,6 +19,7 @@
 package org.apache.streampipes.svcdiscovery;
 
 import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
+import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceStatus;
 import org.apache.streampipes.storage.api.CRUDStorage;
 import org.apache.streampipes.storage.management.StorageDispatcher;
 import org.apache.streampipes.svcdiscovery.api.ISpServiceDiscovery;
@@ -62,7 +63,7 @@ public class SpServiceDiscoveryCore implements 
ISpServiceDiscovery {
         .stream()
         .filter(service -> allFiltersSupported(service, filterByTags))
         .filter(service -> !restrictToHealthy
-            || service.isHealthy())
+            || service.getStatus() != SpServiceStatus.UNHEALTHY)
         .map(this::makeServiceUrl)
         .collect(Collectors.toList());
   }
diff --git 
a/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/CoreRequestSubmitter.java
 
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/CoreRequestSubmitter.java
new file mode 100644
index 000000000..ab7f631d0
--- /dev/null
+++ 
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/CoreRequestSubmitter.java
@@ -0,0 +1,54 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.service.extensions;
+
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.TimeUnit;
+import java.util.function.Supplier;
+
+public class CoreRequestSubmitter {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(CoreRequestSubmitter.class);
+
+  private static final int RETRY_INTERVAL_SECONDS = 3;
+
+  public void submitRepeatedRequest(Supplier<Boolean> request,
+                                    String successMessage,
+                                    String failureMessage) {
+    try {
+      request.get();
+      LOG.info(successMessage);
+    } catch (SpRuntimeException e) {
+      LOG.warn(
+          failureMessage + " Trying again in {} seconds",
+          RETRY_INTERVAL_SECONDS
+      );
+      try {
+        TimeUnit.SECONDS.sleep(RETRY_INTERVAL_SECONDS);
+        submitRepeatedRequest(request, successMessage, failureMessage);
+      } catch (InterruptedException ex) {
+        throw new RuntimeException(ex);
+      }
+    }
+  }
+}
diff --git 
a/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/ExtensionsModelSubmitter.java
 
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/ExtensionsModelSubmitter.java
index cbdb3992f..24f128cc4 100644
--- 
a/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/ExtensionsModelSubmitter.java
+++ 
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/ExtensionsModelSubmitter.java
@@ -52,7 +52,13 @@ public abstract class ExtensionsModelSubmitter extends 
StreamPipesExtensionsServ
 
     // register all migrations at StreamPipes Core
     var migrationConfigs = 
serviceDef.getMigrators().stream().map(IModelMigrator::config).toList();
-    client.adminApi().registerMigrations(migrationConfigs, serviceId());
+    new CoreRequestSubmitter().submitRepeatedRequest(
+        () -> {
+          client.adminApi().registerMigrations(migrationConfigs, serviceId());
+          return true;
+        },
+        "Successfully sent migration request",
+        "Core currently doesn't accept migration requests.");
 
     // initialize all function instances
     
StreamPipesFunctionHandler.INSTANCE.initializeFunctions(serviceDef.getServiceGroup());
diff --git 
a/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/StreamPipesExtensionsServiceBase.java
 
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/StreamPipesExtensionsServiceBase.java
index c13b1153a..0913fd258 100644
--- 
a/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/StreamPipesExtensionsServiceBase.java
+++ 
b/streampipes-service-extensions/src/main/java/org/apache/streampipes/service/extensions/StreamPipesExtensionsServiceBase.java
@@ -19,7 +19,6 @@
 package org.apache.streampipes.service.extensions;
 
 import org.apache.streampipes.client.StreamPipesClient;
-import org.apache.streampipes.commons.exceptions.SpRuntimeException;
 import 
org.apache.streampipes.extensions.management.client.StreamPipesClientResolver;
 import org.apache.streampipes.extensions.management.init.DeclarersSingleton;
 import org.apache.streampipes.extensions.management.model.SpServiceDefinition;
@@ -40,7 +39,6 @@ import jakarta.annotation.PreDestroy;
 import java.net.UnknownHostException;
 import java.util.ArrayList;
 import java.util.List;
-import java.util.concurrent.TimeUnit;
 
 public abstract class StreamPipesExtensionsServiceBase extends 
StreamPipesServiceBase {
 
@@ -96,23 +94,17 @@ public abstract class StreamPipesExtensionsServiceBase 
extends StreamPipesServic
   }
 
   private void registerService(SpServiceRegistration serviceRegistration) {
-    StreamPipesClient client = new 
StreamPipesClientResolver().makeStreamPipesClientInstance();
-    try {
-      client.adminApi().registerService(serviceRegistration);
-      LOG.info("Successfully registered service at core.");
-    } catch (SpRuntimeException e) {
-      LOG.warn(
-          "Could not register at core at url {}. Trying again in {} seconds",
-          client.getConnectionConfig().getBaseUrl(),
-          RETRY_INTERVAL_SECONDS
-      );
-      try {
-        TimeUnit.SECONDS.sleep(RETRY_INTERVAL_SECONDS);
-        registerService(serviceRegistration);
-      } catch (InterruptedException ex) {
-        throw new RuntimeException(ex);
-      }
-    }
+    var client = new 
StreamPipesClientResolver().makeStreamPipesClientInstance();
+    new CoreRequestSubmitter().submitRepeatedRequest(
+        () -> {
+          client.adminApi().registerService(serviceRegistration);
+          return true;
+        },
+        "Successfully registered service at core.",
+        String.format(
+            "Could not register service at core at url %s",
+            client.getConnectionConfig().getBaseUrl()
+        ));
   }
 
   protected List<SpServiceTag> getServiceTags() {
diff --git 
a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/ISpCoreConfigurationStorage.java
 
b/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/ISpCoreConfigurationStorage.java
index 829dc09ce..2d0f8ae65 100644
--- 
a/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/ISpCoreConfigurationStorage.java
+++ 
b/streampipes-storage-api/src/main/java/org/apache/streampipes/storage/api/ISpCoreConfigurationStorage.java
@@ -24,6 +24,8 @@ import java.util.List;
 
 public interface ISpCoreConfigurationStorage {
 
+  boolean exists();
+
   List<SpCoreConfiguration> getAll();
 
   void createElement(SpCoreConfiguration element);
diff --git 
a/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/impl/CoreConfigurationStorageImpl.java
 
b/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/impl/CoreConfigurationStorageImpl.java
index 7ff0a99c0..9b25b5ac8 100644
--- 
a/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/impl/CoreConfigurationStorageImpl.java
+++ 
b/streampipes-storage-couchdb/src/main/java/org/apache/streampipes/storage/couchdb/impl/CoreConfigurationStorageImpl.java
@@ -33,6 +33,11 @@ public class CoreConfigurationStorageImpl extends 
AbstractDao<SpCoreConfiguratio
     super(Utils::getCouchDbGeneralConfigStorage, SpCoreConfiguration.class);
   }
 
+  @Override
+  public boolean exists() {
+    return !findAll().isEmpty();
+  }
+
   @Override
   public List<SpCoreConfiguration> getAll() {
     return findAll();
diff --git 
a/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
 
b/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
index 47f06b81f..85f9a230b 100644
--- 
a/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
+++ 
b/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
@@ -16,10 +16,11 @@
  * specific language governing permissions and limitations
  * under the License.
  */
+
 /* tslint:disable */
 /* eslint-disable */
 // @ts-nocheck
-// Generated using typescript-generator version 3.2.1263 on 2023-10-27 
10:43:45.
+// Generated using typescript-generator version 3.2.1263 on 2023-10-30 
22:49:29.
 
 export class NamedStreamPipesEntity {
     '@class':
@@ -3599,12 +3600,12 @@ export class SpServiceConfiguration {
 export class SpServiceRegistration {
     firstTimeSeenUnhealthy: number;
     healthCheckPath: string;
-    healthy: boolean;
     host: string;
     port: number;
     rev: string;
     scheme: string;
     serviceUrl: string;
+    status: SpServiceStatus;
     svcGroup: string;
     svcId: string;
     svcType: string;
@@ -3620,12 +3621,12 @@ export class SpServiceRegistration {
         const instance = target || new SpServiceRegistration();
         instance.firstTimeSeenUnhealthy = data.firstTimeSeenUnhealthy;
         instance.healthCheckPath = data.healthCheckPath;
-        instance.healthy = data.healthy;
         instance.host = data.host;
         instance.port = data.port;
         instance.rev = data.rev;
         instance.scheme = data.scheme;
         instance.serviceUrl = data.serviceUrl;
+        instance.status = data.status;
         instance.svcGroup = data.svcGroup;
         instance.svcId = data.svcId;
         instance.svcType = data.svcType;
@@ -4105,6 +4106,12 @@ export type SpProtocol = 'KAFKA' | 'JMS' | 'MQTT' | 
'NATS' | 'PULSAR';
 
 export type SpQueryStatus = 'OK' | 'TOO_MUCH_DATA';
 
+export type SpServiceStatus =
+    | 'REGISTERED'
+    | 'MIGRATING'
+    | 'HEALTHY'
+    | 'UNHEALTHY';
+
 export type SpServiceTagPrefix =
     | 'SYSTEM'
     | 'SP_GROUP'
diff --git 
a/ui/src/app/configuration/extensions-service-management/registered-extensions-services/registered-extensions-services.component.html
 
b/ui/src/app/configuration/extensions-service-management/registered-extensions-services/registered-extensions-services.component.html
index 32f153408..c0a1e74b5 100644
--- 
a/ui/src/app/configuration/extensions-service-management/registered-extensions-services/registered-extensions-services.component.html
+++ 
b/ui/src/app/configuration/extensions-service-management/registered-extensions-services/registered-extensions-services.component.html
@@ -38,11 +38,14 @@
                 mat-cell
                 *matCellDef="let element"
             >
-                <span *ngIf="element.healthy" fxLayoutAlign="center center">
+                <span
+                    *ngIf="element.status === 'HEALTHY'"
+                    fxLayoutAlign="center center"
+                >
                     <mat-icon class="service-icon-passing">lens</mat-icon>
                 </span>
                 <span
-                    *ngIf="element.healthy === false"
+                    *ngIf="element.status !== 'HEALTHY'"
                     fxLayoutAlign="center center"
                 >
                     <mat-icon class="service-icon-critical">lens</mat-icon>

Reply via email to