-
Notifications
You must be signed in to change notification settings - Fork 175
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
implement migration for adapter and pipeline element configurations (#…
…2077) * 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 646e792. * 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 * Improvements to extension migration (#2101) * Extract MigrationResource logic into smaller units * Use single request for submitting migrations from extensions to core * refactor: remove dead code * style: change formatting --------- Co-authored-by: bossenti <[email protected]> * refactor: remove redundant assignments * refactor: provide migrators as interface --------- Co-authored-by: Dominik Riemer <[email protected]>
- Loading branch information
1 parent
1a7a17e
commit a34a4dd
Showing
76 changed files
with
2,225 additions
and
123 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
113 changes: 113 additions & 0 deletions
113
...in/java/org/apache/streampipes/connect/management/management/AdapterMigrationManager.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,113 @@ | ||
/* | ||
* 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.connect.management.management; | ||
|
||
import org.apache.streampipes.commons.exceptions.connect.AdapterException; | ||
import org.apache.streampipes.manager.migration.AbstractMigrationManager; | ||
import org.apache.streampipes.manager.migration.IMigrationHandler; | ||
import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceRegistration; | ||
import org.apache.streampipes.model.migration.ModelMigratorConfig; | ||
import org.apache.streampipes.storage.api.IAdapterStorage; | ||
|
||
import org.apache.commons.lang3.StringUtils; | ||
import org.slf4j.Logger; | ||
import org.slf4j.LoggerFactory; | ||
|
||
import java.util.List; | ||
|
||
public class AdapterMigrationManager extends AbstractMigrationManager implements IMigrationHandler { | ||
|
||
private static final Logger LOG = LoggerFactory.getLogger(AdapterMigrationManager.class); | ||
|
||
private final IAdapterStorage adapterStorage; | ||
|
||
public AdapterMigrationManager(IAdapterStorage adapterStorage) { | ||
this.adapterStorage = adapterStorage; | ||
} | ||
|
||
@Override | ||
public void handleMigrations(SpServiceRegistration extensionsServiceConfig, | ||
List<ModelMigratorConfig> migrationConfigs) { | ||
|
||
LOG.info("Received {} migrations from extension service {}.", | ||
migrationConfigs.size(), | ||
extensionsServiceConfig.getServiceUrl()); | ||
LOG.info("Updating adapter descriptions by replacement..."); | ||
updateDescriptions(migrationConfigs, extensionsServiceConfig.getServiceUrl()); | ||
LOG.info("Adapter descriptions are up to date."); | ||
|
||
LOG.info("Checking migrations for existing adapters in StreamPipes Core ..."); | ||
for (var migrationConfig : migrationConfigs) { | ||
LOG.info("Searching for assets of '{}'", migrationConfig.targetAppId()); | ||
LOG.debug("Searching for assets of '{}' with config {}", migrationConfig.targetAppId(), migrationConfig); | ||
var adapterDescriptions = adapterStorage.getAdaptersByAppId(migrationConfig.targetAppId()); | ||
LOG.info("Found {} instances for appId '{}'", adapterDescriptions.size(), migrationConfig.targetAppId()); | ||
for (var adapterDescription : adapterDescriptions) { | ||
|
||
var adapterVersion = adapterDescription.getVersion(); | ||
|
||
if (adapterVersion == migrationConfig.fromVersion()) { | ||
LOG.info("Migration is required for adapter '{}'. Migrating from version '{}' to '{}' ...", | ||
adapterDescription.getElementId(), | ||
adapterVersion, migrationConfig.toVersion() | ||
); | ||
|
||
var migrationResult = performMigration( | ||
adapterDescription, | ||
migrationConfig, | ||
String.format("%s/%s/adapter", | ||
extensionsServiceConfig.getServiceUrl(), | ||
MIGRATION_ENDPOINT | ||
) | ||
); | ||
|
||
if (migrationResult.success()) { | ||
LOG.info("Migration successfully performed by extensions service. Updating adapter description ..."); | ||
LOG.debug( | ||
"Migration was performed by extensions service '{}'", | ||
extensionsServiceConfig.getServiceUrl()); | ||
|
||
adapterStorage.updateAdapter(migrationResult.element()); | ||
LOG.info("Adapter description is updated - Migration successfully completed at Core."); | ||
} else { | ||
LOG.error("Migration failed with the following reason: {}", migrationResult.message()); | ||
LOG.error( | ||
"Migration for adapter '{}' failed - Stopping adapter ...", | ||
migrationResult.element().getElementId() | ||
); | ||
try { | ||
WorkerRestClient.stopStreamAdapter(extensionsServiceConfig.getServiceUrl(), adapterDescription); | ||
} catch (AdapterException e) { | ||
LOG.error("Stopping adapter failed: {}", StringUtils.join(e.getStackTrace(), "\n")); | ||
} | ||
LOG.info("Adapter successfully stopped."); | ||
} | ||
} else { | ||
LOG.info( | ||
"Migration is not applicable for adapter '{}' because of a version mismatch - " | ||
+ "adapter version: '{}', migration starts at: '{}'", | ||
adapterDescription.getElementId(), | ||
adapterVersion, | ||
migrationConfig.fromVersion() | ||
); | ||
} | ||
} | ||
} | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
25 changes: 25 additions & 0 deletions
25
...s-api/src/main/java/org/apache/streampipes/extensions/api/migration/DataSinkMigrator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,25 @@ | ||
/* | ||
* 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.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> { | ||
} |
25 changes: 25 additions & 0 deletions
25
...s-api/src/main/java/org/apache/streampipes/extensions/api/migration/IAdapterMigrator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,25 @@ | ||
/* | ||
* 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.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> { | ||
} |
26 changes: 26 additions & 0 deletions
26
...src/main/java/org/apache/streampipes/extensions/api/migration/IDataProcessorMigrator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,26 @@ | ||
/* | ||
* 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.extensions.api.migration; | ||
|
||
import org.apache.streampipes.extensions.api.extractor.IDataProcessorParameterExtractor; | ||
import org.apache.streampipes.model.graph.DataProcessorInvocation; | ||
|
||
public interface IDataProcessorMigrator | ||
extends IModelMigrator<DataProcessorInvocation, IDataProcessorParameterExtractor> { | ||
} |
Oops, something went wrong.