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

riemer pushed a commit to branch fix-service-registration
in repository https://gitbox.apache.org/repos/asf/streampipes.git

commit 49378b4c103d2ed0a793b9cff14f2456f8cc5b6b
Author: Dominik Riemer <[email protected]>
AuthorDate: Mon Nov 6 21:31:19 2023 +0100

    fix(#2002) Retry service registration in case services are removed before 
migration
---
 .../rest/impl/admin/MigrationResource.java         | 65 ++++++++++++----------
 .../service/extensions/CoreRequestSubmitter.java   | 37 ++++++++++++
 .../extensions/ExtensionsModelSubmitter.java       | 12 ++--
 .../StreamPipesExtensionsServiceBase.java          | 17 ++----
 4 files changed, 80 insertions(+), 51 deletions(-)

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 4d2d66202..340f57d3e 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
@@ -103,39 +103,44 @@ public class MigrationResource extends 
AbstractAuthGuardedRestResource {
       List<ModelMigratorConfig> migrationConfigs) {
 
     var serviceManager = new 
ServiceRegistrationManager(extensionsServiceStorage);
-    var extensionsServiceConfig = serviceManager.getService(serviceId);
-    if (!CoreInitialInstallationProgress.INSTANCE.isInitiallyInstalling()) {
-      new 
AdapterDescriptionMigration093(adapterDescriptionStorage).reinstallAdapters(extensionsServiceConfig);
-      if (!migrationConfigs.isEmpty()) {
-        var anyServiceMigrating = serviceManager.isAnyServiceMigrating();
-        var coreReady = isCoreReady();
-        if (anyServiceMigrating || !coreReady) {
-          LOG.info(
-              "Refusing migration request since precondition is not met 
(anyServiceMigratione={}, coreReady={}.",
-              anyServiceMigrating,
-              coreReady
-          );
-          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);
+    try {
+      var extensionsServiceConfig = serviceManager.getService(serviceId);
+      if (!CoreInitialInstallationProgress.INSTANCE.isInitiallyInstalling()) {
+        new 
AdapterDescriptionMigration093(adapterDescriptionStorage).reinstallAdapters(extensionsServiceConfig);
+        if (!migrationConfigs.isEmpty()) {
+          var anyServiceMigrating = serviceManager.isAnyServiceMigrating();
+          var coreReady = isCoreReady();
+          if (anyServiceMigrating || !coreReady) {
+            LOG.info(
+                "Refusing migration request since precondition is not met 
(anyServiceMigratione={}, coreReady={}.",
+                anyServiceMigrating,
+                coreReady
+            );
+            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();
+    } catch (IllegalArgumentException e) {
+      LOG.warn("Refusing migration request since the service {} is not 
registered.", serviceId);
+      return notFound();
     }
-    new ServiceRegistrationManager(extensionsServiceStorage)
-        .applyServiceStatus(extensionsServiceConfig.getSvcId(), 
SpServiceStatus.HEALTHY);
-    return ok();
   }
 
   private boolean isCoreReady() {
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
index ab7f631d0..d5c2e01af 100644
--- 
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
@@ -18,11 +18,15 @@
 
 package org.apache.streampipes.service.extensions;
 
+import org.apache.streampipes.client.api.IStreamPipesClient;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
+import org.apache.streampipes.model.migration.ModelMigratorConfig;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.List;
 import java.util.concurrent.TimeUnit;
 import java.util.function.Supplier;
 
@@ -51,4 +55,37 @@ public class CoreRequestSubmitter {
       }
     }
   }
+
+  public void submitRegistrationRequest(IStreamPipesClient client,
+                                        SpServiceRegistration serviceReg) {
+    submitRepeatedRequest(
+        () -> {
+          client.adminApi().registerService(serviceReg);
+          return true;
+        },
+        "Successfully registered service at core.",
+        String.format(
+            "Could not register service at core at url %s",
+            client.getConnectionConfig().getBaseUrl()
+        ));
+  }
+
+  public void submitMigrationRequest(IStreamPipesClient client,
+                                     List<ModelMigratorConfig> 
migrationConfigs,
+                                     String serviceId,
+                                     SpServiceRegistration serviceReg) {
+    submitRepeatedRequest(
+        () -> {
+          try {
+            client.adminApi().registerMigrations(migrationConfigs, serviceId);
+            return true;
+          } catch (RuntimeException e) {
+            submitRegistrationRequest(client, serviceReg);
+            submitMigrationRequest(client, migrationConfigs, serviceId, 
serviceReg);
+            return true;
+          }
+        },
+        "Successfully sent migration request",
+        "Core currently doesn't accept migration requests.");
+  }
 }
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 24f128cc4..ac1aed13b 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
@@ -22,6 +22,7 @@ import 
org.apache.streampipes.extensions.api.migration.IModelMigrator;
 import 
org.apache.streampipes.extensions.management.client.StreamPipesClientResolver;
 import org.apache.streampipes.extensions.management.init.DeclarersSingleton;
 import org.apache.streampipes.extensions.management.model.SpServiceDefinition;
+import 
org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration;
 import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceTag;
 import 
org.apache.streampipes.service.extensions.function.StreamPipesFunctionHandler;
 import org.apache.streampipes.service.extensions.security.WebSecurityConfig;
@@ -47,18 +48,13 @@ public abstract class ExtensionsModelSubmitter extends 
StreamPipesExtensionsServ
   }
 
   @Override
-  public void afterServiceRegistered(SpServiceDefinition serviceDef) {
+  public void afterServiceRegistered(SpServiceDefinition serviceDef,
+                                     SpServiceRegistration serviceReg) {
     StreamPipesClient client = new 
StreamPipesClientResolver().makeStreamPipesClientInstance();
 
     // register all migrations at StreamPipes Core
     var migrationConfigs = 
serviceDef.getMigrators().stream().map(IModelMigrator::config).toList();
-    new CoreRequestSubmitter().submitRepeatedRequest(
-        () -> {
-          client.adminApi().registerMigrations(migrationConfigs, serviceId());
-          return true;
-        },
-        "Successfully sent migration request",
-        "Core currently doesn't accept migration requests.");
+    new CoreRequestSubmitter().submitMigrationRequest(client, 
migrationConfigs, serviceId(), serviceReg);
 
     // 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 0913fd258..11c351c8f 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
@@ -43,7 +43,6 @@ import java.util.List;
 public abstract class StreamPipesExtensionsServiceBase extends 
StreamPipesServiceBase {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(StreamPipesExtensionsServiceBase.class);
-  private static final int RETRY_INTERVAL_SECONDS = 3;
 
   public void init() {
     SpServiceDefinition serviceDef = provideServiceDefinition();
@@ -67,7 +66,8 @@ public abstract class StreamPipesExtensionsServiceBase 
extends StreamPipesServic
 
   public abstract SpServiceDefinition provideServiceDefinition();
 
-  public abstract void afterServiceRegistered(SpServiceDefinition serviceDef);
+  public abstract void afterServiceRegistered(SpServiceDefinition serviceDef,
+                                              SpServiceRegistration 
serviceReg);
 
   public void startExtensionsService(Class<?> serviceClass,
                                      SpServiceDefinition serviceDef,
@@ -90,21 +90,12 @@ public abstract class StreamPipesExtensionsServiceBase 
extends StreamPipesServic
         networkingConfig
     );
 
-    this.afterServiceRegistered(serviceDef);
+    this.afterServiceRegistered(serviceDef, req);
   }
 
   private void registerService(SpServiceRegistration serviceRegistration) {
     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()
-        ));
+    new CoreRequestSubmitter().submitRegistrationRequest(client, 
serviceRegistration);
   }
 
   protected List<SpServiceTag> getServiceTags() {

Reply via email to