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>