This is an automated email from the ASF dual-hosted git repository.
wiener pushed a commit to branch edge-extensions
in repository https://gitbox.apache.org/repos/asf/incubator-streampipes.git
The following commit(s) were added to refs/heads/edge-extensions by this push:
new 0a2af10 minor refactoring on container model and node info description
0a2af10 is described below
commit 0a2af10635c4bf4a95b8b9d087c1502337b4ea22
Author: Patrick Wiener <[email protected]>
AuthorDate: Tue Jun 1 17:38:16 2021 +0200
minor refactoring on container model and node info description
---
docker-save.sh | 6 +-
.../model/base/ConsumableStreamPipesEntity.java | 2 +-
.../model/node/NodeInfoDescription.java | 13 +++
.../model/node/NodeInfoDescriptionBuilder.java | 9 ++
.../model/node/container/ContainerEnvBuilder.java | 8 +-
.../{DockerContainer.java => ContainerEnvVar.java} | 54 ++++++++---
.../model/node/container/ContainerLabel.java | 83 ++++++++++++++++
.../model/node/container/ContainerLabels.java | 14 ++-
.../model/node/container/DeploymentContainer.java | 73 ++++++++------
.../model/node/container/DockerContainer.java | 4 +-
.../node/container/DockerContainerBuilder.java | 6 +-
.../streampipes/model/node/container/Ports.java | 7 +-
.../api/ContainerDeploymentResource.java | 4 +-
.../node/controller/config/EnvConfigParam.java | 3 +-
.../node/controller/config/NodeConfiguration.java | 13 +++
.../management/NodeControllerSubmitter.java | 14 +--
.../controller/management/node/NodeConstants.java | 37 ++++---
.../controller/management/node/NodeManager.java | 3 +-
.../offloading/OffloadingPolicyManager.java | 2 +
.../docker/AbstractStreamPipesDockerContainer.java | 8 +-
.../orchestrator/docker/utils/DockerUtils.java | 14 ++-
.../management/resource/utils/ResourceUtils.java | 49 +++++++---
.../node/controller/storage/MapDBImpl.java | 9 +-
.../apache/streampipes/vocabulary/StreamPipes.java | 32 ++++---
.../node-configuration-details.component.html | 106 ++++++++++++++++++++-
.../node-configuration-details.component.ts | 37 ++++++-
.../node-configuration.component.html | 18 +++-
.../node-configuration.component.ts | 2 +-
ui/src/app/core-model/gen/streampipes-model.ts | 56 +++++++++--
.../migrate-pipeline-processors.component.html | 74 ++++++++++----
.../migrate-pipeline-processors.component.scss | 2 +-
.../migrate-pipeline-processors.component.ts | 72 ++++++++++++--
.../node-tag-selector.component.ts | 31 +++---
.../save-pipeline/save-pipeline.component.html | 32 +++----
.../save-pipeline/save-pipeline.component.scss | 2 +-
.../services/pipeline-operations.service.ts | 2 +-
36 files changed, 699 insertions(+), 202 deletions(-)
diff --git a/docker-save.sh b/docker-save.sh
index 45dc594..4885424 100755
--- a/docker-save.sh
+++ b/docker-save.sh
@@ -61,10 +61,10 @@ docker_save_bundle(){
docker save ${docker_img_edge[@]} -o $dir/$docker_bundled_edge_tar
elif [ "$2" == "armv7" ]; then
echo "Save edge images (armv7) to tar ..."
- docker save ${docker_img_edge_arm[@]} -o
$dir/docker_bundled_edge_arm_tar
+ docker save ${docker_img_edge_arm[@]} -o
$dir/$docker_bundled_edge_arm_tar
elif [ "$2" == "aarch64" ]; then
- echo "Save edge images (aarch64) to tar ..."
- docker save ${docker_img_edge_aarch64[@]} -o
$dir/docker_bundled_edge_aarch64_tar
+ echo "Save edge images (aarch64) to tar ..."
${docker_img_edge_aarch64[@]}
+ docker save ${docker_img_edge_aarch64[@]} -o
$dir/$docker_bundled_edge_aarch64_tar
fi
elif [ "$1" == "core" ]; then
echo "Save core images to tar ..."
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/base/ConsumableStreamPipesEntity.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/base/ConsumableStreamPipesEntity.java
index f7936fb..6236014 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/base/ConsumableStreamPipesEntity.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/base/ConsumableStreamPipesEntity.java
@@ -49,7 +49,7 @@ public abstract class ConsumableStreamPipesEntity extends
NamedStreamPipesEntity
@OneToMany(fetch = FetchType.EAGER,
cascade = {CascadeType.ALL})
- @RdfProperty(StreamPipes.HAS_NODE_RESOURCE_PROPERTY)
+ @RdfProperty(StreamPipes.HAS_NODE_RESOURCE_REQUIREMENT)
protected List<NodeResourceRequirement> resourceRequirements;
@OneToOne(fetch = FetchType.EAGER,
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescription.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescription.java
index 8cf0547..58acca0 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescription.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescription.java
@@ -85,6 +85,11 @@ public class NodeInfoDescription extends
UnnamedStreamPipesEntity {
@RdfProperty(StreamPipes.HAS_CONTAINER)
private List<DeploymentContainer> registeredContainers;
+ @OneToMany(fetch = FetchType.EAGER,
+ cascade = {CascadeType.ALL})
+ @RdfProperty(StreamPipes.HAS_DEPLOYMENT_CONTAINER)
+ private List<DeploymentContainer> deploymentContainers;
+
@RdfProperty(StreamPipes.HAS_SUPPORTED_ELEMENTS)
private List<String> supportedElements;
@@ -208,4 +213,12 @@ public class NodeInfoDescription extends
UnnamedStreamPipesEntity {
public void setLastHeartBeatTime(long lastHeartBeatTime) {
this.lastHeartBeatTime = lastHeartBeatTime;
}
+
+ public List<DeploymentContainer> getDeploymentContainers() {
+ return deploymentContainers;
+ }
+
+ public void setDeploymentContainers(List<DeploymentContainer>
deploymentContainers) {
+ this.deploymentContainers = deploymentContainers;
+ }
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescriptionBuilder.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescriptionBuilder.java
index 4e7e810..ba42b91 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescriptionBuilder.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/NodeInfoDescriptionBuilder.java
@@ -17,6 +17,7 @@
*/
package org.apache.streampipes.model.node;
+import com.ibm.dtfj.corereaders.zos.dumpreader.AddressRange;
import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
import org.apache.streampipes.model.grounding.MqttTransportProtocol;
import org.apache.streampipes.model.grounding.TransportProtocol;
@@ -35,6 +36,7 @@ public class NodeInfoDescriptionBuilder {
private final StaticNodeMetadata staticNodeMetadata;
private NodeResource nodeResources;
private List<DeploymentContainer> registeredContainers;
+ private List<DeploymentContainer> deploymentContainers;
private final List<String> supportedElements;
public NodeInfoDescriptionBuilder(String id) {
@@ -43,6 +45,7 @@ public class NodeInfoDescriptionBuilder {
this.staticNodeMetadata = new StaticNodeMetadata();
this.nodeResources = new NodeResource();
this.registeredContainers = new ArrayList<>();
+ this.deploymentContainers = new ArrayList<>();
this.supportedElements = new ArrayList<>();
}
@@ -118,6 +121,11 @@ public class NodeInfoDescriptionBuilder {
return this;
}
+ public NodeInfoDescriptionBuilder
withAutoDeploymentContainers(List<DeploymentContainer> deploymentContainers) {
+ this.deploymentContainers = deploymentContainers;
+ return this;
+ }
+
public NodeInfoDescriptionBuilder withNodeResources(NodeResource
nodeResources) {
this.nodeResources = nodeResources;
return this;
@@ -136,6 +144,7 @@ public class NodeInfoDescriptionBuilder {
nodeInfoDescription.setStaticNodeMetadata(staticNodeMetadata);
nodeInfoDescription.setNodeBroker(nodeBroker);
nodeInfoDescription.setRegisteredContainers(registeredContainers);
+ nodeInfoDescription.setDeploymentContainers(deploymentContainers);
nodeInfoDescription.setNodeResources(nodeResources);
nodeInfoDescription.setSupportedElements(supportedElements);
nodeInfoDescription.setActive(true);
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerEnvBuilder.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerEnvBuilder.java
index 00b8346..e5e2d0a 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerEnvBuilder.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerEnvBuilder.java
@@ -22,7 +22,7 @@ import java.util.List;
public class ContainerEnvBuilder {
- private final List<String> envVariables;
+ private final List<ContainerEnvVar> envVariables;
public ContainerEnvBuilder() {
this.envVariables = new ArrayList<>();
@@ -32,17 +32,17 @@ public class ContainerEnvBuilder {
return new ContainerEnvBuilder();
}
- public ContainerEnvBuilder addNodeEnvs(List<String> nodeEnvVariables) {
+ public ContainerEnvBuilder addNodeEnvs(List<ContainerEnvVar>
nodeEnvVariables) {
this.envVariables.addAll(nodeEnvVariables);
return this;
}
public ContainerEnvBuilder add(String key, String value) {
- this.envVariables.add(String.format("%s=%s", key, value));
+ this.envVariables.add(new ContainerEnvVar(key,value));
return this;
}
- public List<String> build() {
+ public List<ContainerEnvVar> build() {
return envVariables;
}
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainer.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerEnvVar.java
similarity index 51%
copy from
streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainer.java
copy to
streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerEnvVar.java
index 68c7182..2e156d3 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainer.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerEnvVar.java
@@ -17,32 +17,60 @@
*/
package org.apache.streampipes.model.node.container;
+import io.fogsy.empire.annotations.RdfProperty;
import io.fogsy.empire.annotations.RdfsClass;
+import org.apache.streampipes.model.base.UnnamedStreamPipesEntity;
import org.apache.streampipes.model.shared.annotation.TsModel;
import org.apache.streampipes.vocabulary.StreamPipes;
import javax.persistence.Entity;
-import java.util.List;
-import java.util.Map;
-@RdfsClass(StreamPipes.DEPLOYMENT_DOCKER_CONTAINER)
+
+@RdfsClass(StreamPipes.CONTAINER_ENV_VAR)
@Entity
@TsModel
-public class DockerContainer extends DeploymentContainer {
+public class ContainerEnvVar extends UnnamedStreamPipesEntity {
+
+ @RdfProperty(StreamPipes.HAS_CONTAINER_ENV_KEY)
+ private String key;
+
+ @RdfProperty(StreamPipes.HAS_CONTAINER_ENV_VALUE)
+ private String value;
+
+ public ContainerEnvVar() {
+ super();
+ }
+
+ public ContainerEnvVar(String key, String value) {
+ this.key = key;
+ this.value = value;
+ }
- public DockerContainer(String elementId) {
+ public ContainerEnvVar(UnnamedStreamPipesEntity other, String key, String
value) {
+ super(other);
+ this.key = key;
+ this.value = value;
+ }
+
+ public ContainerEnvVar(String elementId, String key, String value) {
super(elementId);
+ this.key = key;
+ this.value = value;
}
- public DockerContainer() {
- super();
+ public String getKey() {
+ return key;
+ }
+
+ public void setKey(String key) {
+ this.key = key;
+ }
+
+ public String getValue() {
+ return value;
}
- public DockerContainer(String imageTag, String containerName, String
serviceId, String[] containerPorts,
- List<String> envVars, Map<String, String> labels,
List<String> volumes,
- List<String> supportedArchitectures, List<String>
supportedOperatingSystemTypes,
- List<String> dependsOnContainers) {
- super(imageTag, containerName, serviceId, containerPorts, envVars,
labels, volumes, supportedArchitectures,
- supportedOperatingSystemTypes, dependsOnContainers);
+ public void setValue(String value) {
+ this.value = value;
}
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerLabel.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerLabel.java
new file mode 100644
index 0000000..0397bb7
--- /dev/null
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerLabel.java
@@ -0,0 +1,83 @@
+/*
+ * 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.model.node.container;
+
+
+import io.fogsy.empire.annotations.RdfProperty;
+import io.fogsy.empire.annotations.RdfsClass;
+import org.apache.streampipes.model.base.UnnamedStreamPipesEntity;
+import org.apache.streampipes.model.shared.annotation.TsModel;
+import org.apache.streampipes.vocabulary.StreamPipes;
+
+import javax.persistence.Entity;
+import java.util.Map;
+
+@RdfsClass(StreamPipes.CONTAINER_LABEL)
+@Entity
+@TsModel
+public class ContainerLabel extends UnnamedStreamPipesEntity {
+
+ @RdfProperty(StreamPipes.HAS_CONTAINER_LABEL_KEY)
+ private String key;
+
+ @RdfProperty(StreamPipes.HAS_CONTAINER_LABEL_VALUE)
+ private String value;
+
+ public ContainerLabel() {
+ super();
+ }
+
+ public ContainerLabel(String key, String value) {
+ super();
+ this.key = key;
+ this.value = value;
+ }
+
+ public ContainerLabel(UnnamedStreamPipesEntity other, String key, String
value) {
+ super(other);
+ this.key = key;
+ this.value = value;
+ }
+
+ public ContainerLabel(String elementId, String key, String value) {
+ super(elementId);
+ this.key = key;
+ this.value = value;
+ }
+
+ public ContainerLabel(Map.Entry<String, String> entity) {
+ this.key = entity.getKey();
+ this.value = entity.getValue();
+ }
+
+ public String getKey() {
+ return key;
+ }
+
+ public void setKey(String key) {
+ this.key = key;
+ }
+
+ public String getValue() {
+ return value;
+ }
+
+ public void setValue(String value) {
+ this.value = value;
+ }
+}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerLabels.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerLabels.java
index 9ec07e8..f65bbe3 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerLabels.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/ContainerLabels.java
@@ -17,15 +17,13 @@
*/
package org.apache.streampipes.model.node.container;
-import java.util.HashMap;
-import java.util.Map;
+import java.util.*;
public class ContainerLabels {
- public static Map<String, String> with (String containerId, String
nodeType, ContainerType containerType) {
- return new HashMap<String, String>() {{
- put("org.apache.streampipes.service.id", containerId);
- put("org.apache.streampipes.node.type",nodeType);
- put("org.apache.streampipes.container.type", containerType.name());
- }};
+ public static List<ContainerLabel> with (String containerId, String
nodeType, ContainerType containerType) {
+ return Arrays.asList(
+ new ContainerLabel("org.apache.streampipes.service.id",
containerId),
+ new
ContainerLabel("org.apache.streampipes.node.type",nodeType),
+ new ContainerLabel("org.apache.streampipes.container.type",
containerType.name().toLowerCase()));
}
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DeploymentContainer.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DeploymentContainer.java
index ec90467..3907397 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DeploymentContainer.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DeploymentContainer.java
@@ -22,13 +22,12 @@ import io.fogsy.empire.annotations.RdfProperty;
import io.fogsy.empire.annotations.RdfsClass;
import org.apache.streampipes.model.base.UnnamedStreamPipesEntity;
import org.apache.streampipes.model.shared.annotation.TsModel;
+import org.apache.streampipes.vocabulary.RDFS;
import org.apache.streampipes.vocabulary.StreamPipes;
import javax.persistence.*;
-import java.util.ArrayList;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
+import java.util.*;
+import java.util.stream.Collectors;
@RdfsClass(StreamPipes.DEPLOYMENT_CONTAINER)
@Entity
@@ -38,52 +37,56 @@ import java.util.Map;
@TsModel
public abstract class DeploymentContainer extends UnnamedStreamPipesEntity {
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_IMAGE_TAG)
+ protected static final String prefix = "urn:streampipes.org:spi:";
+
+ @RdfProperty(StreamPipes.HAS_IMAGE_TAG)
private String imageTag;
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_NAME)
+ @RdfProperty(StreamPipes.HAS_CONTAINER_NAME)
private String containerName;
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_SERVICE_ID)
+ @RdfProperty(StreamPipes.HAS_CONTAINER_SERVICE_ID)
private String serviceId;
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_PORTS)
- private String[] containerPorts;
+ @OneToMany(fetch = FetchType.EAGER,
+ cascade = {CascadeType.ALL})
+ @RdfProperty(StreamPipes.HAS_CONTAINER_PORTS)
+ private List<String> containerPorts;
@OneToMany(fetch = FetchType.EAGER,
cascade = {CascadeType.ALL})
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_ENV_VARS)
- private List<String> envVars;
+ @RdfProperty(StreamPipes.HAS_CONTAINER_ENV_VARS)
+ private List<ContainerEnvVar> envVars;
@OneToMany(fetch = FetchType.EAGER,
cascade = {CascadeType.ALL})
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_LABELS)
- private Map<String, String> labels;
+ @RdfProperty(StreamPipes.HAS_CONTAINER_LABELS)
+ private List<ContainerLabel> labels;
@OneToMany(fetch = FetchType.EAGER,
cascade = {CascadeType.ALL})
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_VOLUMES)
+ @RdfProperty(StreamPipes.HAS_CONTAINER_VOLUMES)
private List<String> volumes;
@OneToMany(fetch = FetchType.EAGER,
cascade = {CascadeType.ALL})
- @RdfProperty(StreamPipes.DEPLOYMENT_SUPPORTED_ARCHITECTURES)
+ @RdfProperty(StreamPipes.HAS_SUPPORTED_ARCHITECTURES)
private List<String> supportedArchitectures;
@OneToMany(fetch = FetchType.EAGER,
cascade = {CascadeType.ALL})
- @RdfProperty(StreamPipes.DEPLOYMENT_SUPPORTED_OS_TYPES)
+ @RdfProperty(StreamPipes.HAS_SUPPORTED_OS_TYPES)
private List<String> supportedOperatingSystemTypes;
@OneToMany(fetch = FetchType.EAGER,
cascade = {CascadeType.ALL})
- @RdfProperty(StreamPipes.DEPLOYMENT_CONTAINER_DEPENDENCIES)
+ @RdfProperty(StreamPipes.HAS_CONTAINER_DEPENDENCIES)
private List<String> dependsOnContainers;
public DeploymentContainer() {
this.envVars = new ArrayList<>();
- this.labels = new HashMap<>();
+ this.labels = new ArrayList<>();
this.volumes = new ArrayList<>();
this.dependsOnContainers = new ArrayList<>();
this.supportedArchitectures = new ArrayList<>();
@@ -93,7 +96,7 @@ public abstract class DeploymentContainer extends
UnnamedStreamPipesEntity {
public DeploymentContainer(String elementId) {
super(elementId);
this.envVars = new ArrayList<>();
- this.labels = new HashMap<>();
+ this.labels = new ArrayList<>();
this.volumes = new ArrayList<>();
this.dependsOnContainers = new ArrayList<>();
this.supportedArchitectures = new ArrayList<>();
@@ -102,12 +105,23 @@ public abstract class DeploymentContainer extends
UnnamedStreamPipesEntity {
public DeploymentContainer(DeploymentContainer other) {
super(other);
- }
-
- public DeploymentContainer(String imageTag, String containerName, String
serviceId, String[] containerPorts,
- List<String> envVars, Map<String, String>
labels, List<String> volumes,
+ this.imageTag = other.getImageTag();
+ this.containerName = other.getContainerName();
+ this.serviceId = other.getServiceId();
+ this.containerPorts = other.getContainerPorts();
+ this.envVars = other.getEnvVars();
+ this.labels = other.getLabels();
+ this.volumes = other.getVolumes();
+ this.supportedArchitectures = other.getSupportedArchitectures();
+ this.supportedOperatingSystemTypes =
other.getSupportedOperatingSystemTypes();
+ this.dependsOnContainers = other.getDependsOnContainers();
+ }
+
+ public DeploymentContainer(String imageTag, String containerName, String
serviceId, List<String> containerPorts,
+ List<ContainerEnvVar> envVars,
List<ContainerLabel> labels, List<String> volumes,
List<String> supportedArchitectures,
List<String> supportedOperatingSystemsTypes,
List<String> dependsOnContainers) {
+ super(prefix + UUID.randomUUID().toString());
this.imageTag = imageTag;
this.containerName = containerName;
this.serviceId = serviceId;
@@ -143,27 +157,27 @@ public abstract class DeploymentContainer extends
UnnamedStreamPipesEntity {
this.serviceId = serviceId;
}
- public String[] getContainerPorts() {
+ public List<String> getContainerPorts() {
return containerPorts;
}
- public void setContainerPorts(String[] containerPorts) {
+ public void setContainerPorts(List<String> containerPorts) {
this.containerPorts = containerPorts;
}
- public List<String> getEnvVars() {
+ public List<ContainerEnvVar> getEnvVars() {
return envVars;
}
- public void setEnvVars(List<String> envVars) {
+ public void setEnvVars(List<ContainerEnvVar> envVars) {
this.envVars = envVars;
}
- public Map<String, String> getLabels() {
+ public List<ContainerLabel> getLabels() {
return labels;
}
- public void setLabels(Map<String, String> labels) {
+ public void setLabels(List<ContainerLabel> labels) {
this.labels = labels;
}
@@ -198,5 +212,4 @@ public abstract class DeploymentContainer extends
UnnamedStreamPipesEntity {
public void setSupportedOperatingSystemTypes(List<String>
supportedOperatingSystemTypes) {
this.supportedOperatingSystemTypes = supportedOperatingSystemTypes;
}
-
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainer.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainer.java
index 68c7182..db20469 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainer.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainer.java
@@ -38,8 +38,8 @@ public class DockerContainer extends DeploymentContainer {
super();
}
- public DockerContainer(String imageTag, String containerName, String
serviceId, String[] containerPorts,
- List<String> envVars, Map<String, String> labels,
List<String> volumes,
+ public DockerContainer(String imageTag, String containerName, String
serviceId, List<String> containerPorts,
+ List<ContainerEnvVar> envVars, List<ContainerLabel>
labels, List<String> volumes,
List<String> supportedArchitectures, List<String>
supportedOperatingSystemTypes,
List<String> dependsOnContainers) {
super(imageTag, containerName, serviceId, containerPorts, envVars,
labels, volumes, supportedArchitectures,
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainerBuilder.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainerBuilder.java
index 86a59b6..316b404 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainerBuilder.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/DockerContainerBuilder.java
@@ -45,17 +45,17 @@ public class DockerContainerBuilder {
return this;
}
- public DockerContainerBuilder withExposedPorts(String[] ports) {
+ public DockerContainerBuilder withExposedPorts(List<String> ports) {
this.dockerContainer.setContainerPorts(ports);
return this;
}
- public DockerContainerBuilder withEnvironmentVariables(List<String> envs) {
+ public DockerContainerBuilder
withEnvironmentVariables(List<ContainerEnvVar> envs) {
this.dockerContainer.setEnvVars(envs);
return this;
}
- public DockerContainerBuilder withLabels(Map<String, String> labels) {
+ public DockerContainerBuilder withLabels(List<ContainerLabel> labels) {
this.dockerContainer.setLabels(labels);
return this;
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/Ports.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/Ports.java
index 55a64bb..b079d02 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/Ports.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/node/container/Ports.java
@@ -17,8 +17,11 @@
*/
package org.apache.streampipes.model.node.container;
+import java.util.Arrays;
+import java.util.List;
+
public class Ports {
- public static String[] withMapping(String... port) {
- return port;
+ public static List<String> withMapping(String... port) {
+ return Arrays.asList(port);
}
}
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/api/ContainerDeploymentResource.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/api/ContainerDeploymentResource.java
index 50c398c..5520321 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/api/ContainerDeploymentResource.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/api/ContainerDeploymentResource.java
@@ -39,9 +39,9 @@ public class ContainerDeploymentResource extends
AbstractResource {
}
@GET
- @Path("/registered")
+ @Path("/catalog")
@Produces(MediaType.APPLICATION_JSON)
- public javax.ws.rs.core.Response getAllRegisteredContainer(){
+ public javax.ws.rs.core.Response getAllAvailableDeploymentContainers(){
return
ok(DockerContainerDeclarerSingleton.getInstance().getAllDockerContainerAsList());
}
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/EnvConfigParam.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/EnvConfigParam.java
index 0854466..2c709de 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/EnvConfigParam.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/EnvConfigParam.java
@@ -49,7 +49,8 @@ public enum EnvConfigParam {
NODE_ACCESSIBLE_FIELD_DEVICE("SP_NODE_ACCESSIBLE_FIELD_DEVICE", ""),
CONSUL_LOCATION("CONSUL_LOCATION", "consul"),
SUPPORTED_PIPELINE_ELEMENTS("SP_SUPPORTED_PIPELINE_ELEMENTS", ""),
- AUTO_OFFLOADING_STRATEGY("SP_AUTO_OFFLOADING_STRATEGY", "default");
+ AUTO_OFFLOADING_STRATEGY("SP_AUTO_OFFLOADING_STRATEGY", "default"),
+ NODE_STORAGE_PATH("SP_NODE_STORAGE_PATH", "/var/lib/streampipes");
private final String environmentKey;
private final String defaultValue;
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/NodeConfiguration.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/NodeConfiguration.java
index 872464f..44b1291 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/NodeConfiguration.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/config/NodeConfiguration.java
@@ -59,6 +59,7 @@ public final class NodeConfiguration {
private static int relayEventBufferSize;
private static String consulHost;
private static OffloadingStrategyType autoOffloadingStrategy;
+ private static String nodeStoragePath;
private static HashMap<String, String> configMap;
@@ -282,6 +283,14 @@ public final class NodeConfiguration {
NodeConfiguration.autoOffloadingStrategy = autoOffloadingStrategy;
}
+ public static String getNodeStoragePath() {
+ return nodeStoragePath;
+ }
+
+ public static void setNodeStoragePath(String nodeStoragePath) {
+ NodeConfiguration.nodeStoragePath = nodeStoragePath;
+ }
+
public static HashMap<String, String> getConfigMap() {
return configMap;
}
@@ -492,6 +501,10 @@ public final class NodeConfiguration {
configMap.put(envKey, value);
setAutoOffloadingStrategy(OffloadingStrategyType.fromString(value));
break;
+ case NODE_STORAGE_PATH:
+ configMap.put(envKey, value);
+ setNodeStoragePath(value);
+ break;
default:
throw new IllegalArgumentException("Invalid environment
config param: " + configParam);
}
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/NodeControllerSubmitter.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/NodeControllerSubmitter.java
index 9f4a669..741f96c 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/NodeControllerSubmitter.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/NodeControllerSubmitter.java
@@ -55,6 +55,13 @@ public abstract class NodeControllerSubmitter {
app.setDefaultProperties(Collections.singletonMap("server.port",
NodeConfiguration.getNodeControllerPort()));
app.run();
+ LOG.info("Register container descriptions");
+ DockerContainerDeclarerSingleton.getInstance()
+ .register(new DockerExtensionsContainer())
+ .register(new DockerMosquittoContainer())
+ .register(new DockerKafkaContainer())
+ .register(new DockerZookeeperContainer());
+
LOG.info("Load node info description");
NodeManager.getInstance().init();
@@ -67,13 +74,6 @@ public abstract class NodeControllerSubmitter {
if (!"true".equals(System.getenv("SP_DEBUG"))) {
- LOG.info("Register container descriptions");
- DockerContainerDeclarerSingleton.getInstance()
- .register(new DockerExtensionsContainer())
- .register(new DockerMosquittoContainer())
- .register(new DockerKafkaContainer())
- .register(new DockerZookeeperContainer());
-
LOG.info("Auto-deploy extensions and selected broker
container");
DockerEngineManager.getInstance().init();
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeConstants.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeConstants.java
index 965d230..270b449 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeConstants.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeConstants.java
@@ -17,6 +17,7 @@
*/
package org.apache.streampipes.node.controller.management.node;
+import org.apache.streampipes.model.node.container.ContainerLabel;
import org.apache.streampipes.model.node.container.DeploymentContainer;
import org.apache.streampipes.model.node.container.DockerContainer;
import org.apache.streampipes.model.node.resources.software.ContainerRuntime;
@@ -28,6 +29,7 @@ import
org.apache.streampipes.model.node.resources.hardware.DISK;
import org.apache.streampipes.model.node.resources.hardware.GPU;
import org.apache.streampipes.model.node.resources.hardware.MEM;
import org.apache.streampipes.node.controller.config.NodeConfiguration;
+import
org.apache.streampipes.node.controller.management.orchestrator.docker.DockerContainerDeclarerSingleton;
import
org.apache.streampipes.node.controller.management.orchestrator.docker.model.DockerInfo;
import
org.apache.streampipes.node.controller.management.orchestrator.docker.utils.DockerUtils;
import oshi.SystemInfo;
@@ -64,7 +66,9 @@ public class NodeConstants {
public static final List<String> SUPPORTED_PIPELINE_ELEMENTS =
NodeConfiguration.getSupportedPipelineElements();
public static final String NODE_MODEL =
!printComputerSystem(hal.getComputerSystem()).equals("") ?
printComputerSystem(hal.getComputerSystem()) : "n/a";
- public static final List<DeploymentContainer> REGISTERED_DOCKER_CONTAINER
= getRegisteredDockerContainer();
+ public static final List<DeploymentContainer>
REGISTERED_DEPLOYMENT_CONTAINERS = getRegisteredDeploymentContainers();
+ public static final List<DeploymentContainer> AUTO_DEPLOYMENT_CONTAINERS =
getAutoDeploymentContainer();
+
public static final String NODE_OPERATING_SYSTEM = docker.getOs();
public static final String NODE_KERNEL_VERSION = docker.getKernelVersion();
public static final ContainerRuntime NODE_CONTAINER_RUNTIME =
getContainerRuntime();
@@ -75,29 +79,38 @@ public class NodeConstants {
public static final DISK NODE_DISK = getNodeDisk();
public static final GPU NODE_GPU = getNodeGpu();
- private static List<DeploymentContainer> getRegisteredDockerContainer() {
+ private static List<DeploymentContainer>
getRegisteredDeploymentContainers() {
List<DeploymentContainer> containers = new ArrayList<>();
DockerUtils.getInstance().getRunningStreamPipesContainer()
- .forEach(rc -> {
- DockerContainer c = new DockerContainer();
- c.setContainerName(rc.names().get(0).replace("/", ""));
- c.setImageTag(rc.image());
+ .forEach(runningContainer -> {
+ DockerContainer container = new DockerContainer();
+
container.setContainerName(runningContainer.names().get(0).replace("/", ""));
+ container.setImageTag(runningContainer.image());
- Optional<String> serviceId =
rc.labels().entrySet().stream()
+ Optional<String> serviceId =
runningContainer.labels().entrySet().stream()
.filter(l ->
l.getKey().contains("org.apache.streampipes.service.id"))
.map(Map.Entry::getValue)
.findFirst();
- serviceId.ifPresent(c::setServiceId);
+ serviceId.ifPresent(container::setServiceId);
- c.setLabels(rc.labels().entrySet().stream()
- .filter(l ->
l.getKey().contains("org.apache.streampipes"))
- .collect(Collectors.toMap(Map.Entry::getKey,
Map.Entry::getValue)));
- containers.add(c);
+ container.setLabels(
+ runningContainer.labels().entrySet().stream()
+ .filter(l ->
l.getKey().contains("org.apache.streampipes"))
+ .map(ContainerLabel::new)
+ .collect(Collectors.toList()));
+ containers.add(container);
});
return containers;
}
+ private static List<DeploymentContainer> getAutoDeploymentContainer() {
+ List<DeploymentContainer> deploymentContainers = new ArrayList<>();
+
deploymentContainers.addAll(DockerContainerDeclarerSingleton.getInstance().getAllDockerContainerAsList());
+
+ return deploymentContainers;
+ }
+
private static CPU getNodeCpu(){
CPU cpu = new CPU();
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeManager.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeManager.java
index df205a0..fcd7b36 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeManager.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/node/NodeManager.java
@@ -60,7 +60,8 @@ public class NodeManager implements INodeManager {
new GeoLocation(),
NodeConstants.NODE_LOCATION_TAGS)
.withSupportedElements(NodeConstants.SUPPORTED_PIPELINE_ELEMENTS)
-
.withRegisteredContainers(NodeConstants.REGISTERED_DOCKER_CONTAINER)
+
.withRegisteredContainers(NodeConstants.REGISTERED_DEPLOYMENT_CONTAINERS)
+
.withAutoDeploymentContainers(NodeConstants.AUTO_DEPLOYMENT_CONTAINERS)
.withNodeResources(NodeResourceBuilder.create()
.hardwareResource(
NodeConstants.NODE_CPU,
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/offloading/OffloadingPolicyManager.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/offloading/OffloadingPolicyManager.java
index c9c51e1..6ac8234 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/offloading/OffloadingPolicyManager.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/offloading/OffloadingPolicyManager.java
@@ -58,6 +58,8 @@ public class OffloadingPolicyManager {
for(OffloadingStrategy strategy : offloadingStrategies){
strategy.getOffloadingPolicy().addValue(strategy.getResourceProperty().getProperty(rm));
if(strategy.getOffloadingPolicy().isViolated()){
+ String violatedProperty =
strategy.getResourceProperty().getClass().getSimpleName();
+ LOG.info("Violated policy for resource property {}",
violatedProperty);
violatedPolicies.add(strategy);
}
}
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/AbstractStreamPipesDockerContainer.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/AbstractStreamPipesDockerContainer.java
index 9b76e35..c4f63eb 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/AbstractStreamPipesDockerContainer.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/AbstractStreamPipesDockerContainer.java
@@ -17,6 +17,7 @@
*/
package org.apache.streampipes.node.controller.management.orchestrator.docker;
+import org.apache.streampipes.model.node.container.ContainerEnvVar;
import org.apache.streampipes.model.node.container.DockerContainer;
import org.apache.streampipes.node.controller.config.EnvConfigParam;
import org.apache.streampipes.node.controller.config.NodeConfiguration;
@@ -35,7 +36,7 @@ public abstract class AbstractStreamPipesDockerContainer {
NodeConfiguration.getStreampipesVersion() :
VersionUtils.getStreamPipesVersion();
}
- public List<String> generateStreamPipesNodeEnvs() {
+ public List<ContainerEnvVar> generateStreamPipesNodeEnvs() {
return new ArrayList<>(Arrays.asList(
toEnv(EnvConfigParam.NODE_CONTROLLER_ID.getEnvironmentKey(),
NodeConfiguration.getNodeControllerId()),
@@ -58,7 +59,8 @@ public abstract class AbstractStreamPipesDockerContainer {
// Helper
- private <T>String toEnv(String key, T value) {
- return String.format("%s=%s", key, value);
+ private <T>ContainerEnvVar toEnv(String key, T value) {
+ return new ContainerEnvVar(key, String.valueOf(value));
+ //return String.format("%s=%s", key, value);
}
}
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/utils/DockerUtils.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/utils/DockerUtils.java
index 0080900..387ad7c 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/utils/DockerUtils.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/orchestrator/docker/utils/DockerUtils.java
@@ -27,6 +27,7 @@ import com.spotify.docker.client.messages.*;
import
com.spotify.docker.client.shaded.com.google.common.collect.ImmutableList;
import com.spotify.docker.client.shaded.com.google.common.collect.Lists;
import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.model.node.container.ContainerLabel;
import org.apache.streampipes.model.node.container.DockerContainer;
import
org.apache.streampipes.node.controller.management.orchestrator.docker.model.DockerInfo;
import org.slf4j.Logger;
@@ -138,8 +139,11 @@ public class DockerUtils {
.hostname(p.getContainerName())
.tty(true)
.image(p.getImageTag())
- .labels(p.getLabels())
- .env(p.getEnvVars())
+ .labels(p.getLabels().stream()
+ .collect(Collectors.toMap(ContainerLabel::getKey,
ContainerLabel::getValue)))
+ .env(p.getEnvVars().stream()
+ .map(e -> String.format("%s=%s", e.getKey(),
e.getValue()))
+ .collect(Collectors.toList()))
.hostConfig(getHostConfig(SP_CONTAINER_NETWORK,
p.getContainerPorts(), p.getVolumes()))
.exposedPorts(modifyExposedPorts(p.getContainerPorts()))
.networkingConfig(getNetworkingConfig(SP_CONTAINER_NETWORK,
p.getContainerName()))
@@ -175,7 +179,7 @@ public class DockerUtils {
return Collections.emptyList();
}
- private HostConfig getHostConfig(String network, String[] ports) {
+ private HostConfig getHostConfig(String network, List<String> ports) {
return getHostConfig(network, ports, null);
}
@@ -183,7 +187,7 @@ public class DockerUtils {
return getHostConfig(network, null, null);
}
- private HostConfig getHostConfig(String network, String[] ports,
List<String> volumes) {
+ private HostConfig getHostConfig(String network, List<String> ports,
List<String> volumes) {
Map<String, List<PortBinding>> portBindings = new HashMap<>();
if (ports != null) {
for (String port : ports) {
@@ -350,7 +354,7 @@ public class DockerUtils {
return hasNvidiaRuntime;
}
- private String[] modifyExposedPorts(String[] containerPorts) {
+ private String[] modifyExposedPorts(List<String> containerPorts) {
List<String> modifyPorts = new ArrayList<>();
for (String port: containerPorts) {
modifyPorts.add(port + "/tcp");
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/resource/utils/ResourceUtils.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/resource/utils/ResourceUtils.java
index a584397..717b486 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/resource/utils/ResourceUtils.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/management/resource/utils/ResourceUtils.java
@@ -17,6 +17,8 @@
*/
package org.apache.streampipes.node.controller.management.resource.utils;
+import org.apache.streampipes.node.controller.config.NodeConfiguration;
+import org.apache.streampipes.node.controller.management.node.NodeManager;
import oshi.hardware.CentralProcessor;
import oshi.hardware.GlobalMemory;
import oshi.hardware.Sensors;
@@ -40,28 +42,45 @@ public class ResourceUtils {
public static Map<String, Map<String,Long>> getDiskUsage(FileSystem fs) {
List<OSFileStore> fsArray = fs.getFileStores();
Map<String, Map<String, Long>> diskUsage = new HashMap<>();
+
for(OSFileStore f : fsArray) {
String volume = f.getVolume();
- // has SATA disk
- if (volume.contains(FileSystemType.SDA.getName())){
- addDiskUsage(diskUsage, f);
- } else if (volume.contains(FileSystemType.NVME.getName())){
- addDiskUsage(diskUsage, f);
- } else if (volume.contains(FileSystemType.DISK.getName())){
- addDiskUsage(diskUsage, f);
- } else if (volume.contains(FileSystemType.ROOT.getName())){
- // Docker in RPi
- addDiskUsage(diskUsage, f);
- } else if (volume.contains(FileSystemType.MMCBLK.getName())){
- // Docker in Jetson Nano
- addDiskUsage(diskUsage, f);
- } else if (volume.contains(FileSystemType.SDB.getName())){
- addDiskUsage(diskUsage, f);
+ String mount = f.getMount();
+
+ if ("true".equals(System.getenv("SP_DEBUG"))) {
+ // TODO: find better approach to retrieve host volume in debug
setup
+ findVolumeAndAdd(volume, diskUsage, f);
+ } else {
+ // check mount for node storage path (default:
/var/lib/streampipes)
+ if (mount.equals(NodeConfiguration.getNodeStoragePath())) {
+ findVolumeAndAdd(volume, diskUsage, f);
+ }
}
+
}
return diskUsage.isEmpty() ? defaultDiskUsage() : diskUsage;
}
+ private static void findVolumeAndAdd(String volume, Map<String,
Map<String, Long>> diskUsage, OSFileStore f) {
+ // has SATA disk
+ if (volume.contains(FileSystemType.SDA.getName())){
+ addDiskUsage(diskUsage, f);
+ // Docker in Jetson Xavier with NVMe storage
+ } else if (volume.contains(FileSystemType.NVME.getName())){
+ addDiskUsage(diskUsage, f);
+ } else if (volume.contains(FileSystemType.DISK.getName())){
+ addDiskUsage(diskUsage, f);
+ } else if (volume.contains(FileSystemType.ROOT.getName())){
+ // Docker in RPi
+ addDiskUsage(diskUsage, f);
+ } else if (volume.contains(FileSystemType.MMCBLK.getName())){
+ // Docker in Jetson Nano
+ addDiskUsage(diskUsage, f);
+ } else if (volume.contains(FileSystemType.SDB.getName())){
+ addDiskUsage(diskUsage, f);
+ }
+ }
+
public static void addDiskUsage(Map<String, Map<String, Long>> m,
OSFileStore f) {
Map<String, Long> i = new HashMap<>();
i.put(DiskSpace.USABLE.getName(), f.getUsableSpace());
diff --git
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/storage/MapDBImpl.java
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/storage/MapDBImpl.java
index ebe5518..3647523 100644
---
a/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/storage/MapDBImpl.java
+++
b/streampipes-node-controller/src/main/java/org/apache/streampipes/node/controller/storage/MapDBImpl.java
@@ -17,6 +17,7 @@
*/
package org.apache.streampipes.node.controller.storage;
+import org.apache.streampipes.node.controller.config.NodeConfiguration;
import org.mapdb.DB;
import org.mapdb.DBMaker;
import org.mapdb.Serializer;
@@ -28,10 +29,8 @@ import java.util.stream.Collectors;
public class MapDBImpl implements CRUDStorage {
- private static final String DB_STORAGE_PATH = "/var/lib/streampipes/";
-
- private DB db;
- private ConcurrentMap<String, Object> map;
+ private final DB db;
+ private final ConcurrentMap<String, Object> map;
public MapDBImpl(File dbFile) {
if("true".equals(System.getenv("SP_DEBUG"))) {
@@ -41,7 +40,7 @@ public class MapDBImpl implements CRUDStorage {
.make();
} else {
db = DBMaker
- .fileDB(DB_STORAGE_PATH + dbFile)
+ .fileDB(NodeConfiguration.getNodeStoragePath() + dbFile)
.transactionEnable()
.closeOnJvmShutdown()
.make();
diff --git
a/streampipes-vocabulary/src/main/java/org/apache/streampipes/vocabulary/StreamPipes.java
b/streampipes-vocabulary/src/main/java/org/apache/streampipes/vocabulary/StreamPipes.java
index 51a6cb9..81fd4bd 100644
---
a/streampipes-vocabulary/src/main/java/org/apache/streampipes/vocabulary/StreamPipes.java
+++
b/streampipes-vocabulary/src/main/java/org/apache/streampipes/vocabulary/StreamPipes.java
@@ -361,17 +361,17 @@ public class StreamPipes {
public static final String SECRET_STATIC_PROPERTY = NS +
"SecretStaticProperty";
public static final String IS_ENCRYPTED = NS + "isEncrypted";
- public static final String DEPLOYMENT_DOCKER_CONTAINER = NS +
"pipelineElementDockerContainer";
- public static final String DEPLOYMENT_CONTAINER_IMAGE_TAG = NS +
"dockerContainerImageTag";
- public static final String DEPLOYMENT_CONTAINER_NAME = NS +
"dockerContainerName";
- public static final String DEPLOYMENT_CONTAINER_SERVICE_ID = NS +
"dockerContainerServiceId";
- public static final String DEPLOYMENT_CONTAINER_PORTS = NS +
"dockerContainerPorts";
- public static final String DEPLOYMENT_CONTAINER_ENV_VARS = NS +
"dockerContainerEnvVars";
- public static final String DEPLOYMENT_CONTAINER_LABELS = NS +
"dockerContainerLabels";
- public static final String DEPLOYMENT_CONTAINER_VOLUMES = NS +
"dockerContainerVolumes";
- public static final String DEPLOYMENT_CONTAINER_DEPENDENCIES = NS +
"dockerContainerDependencies";
- public static final String DEPLOYMENT_SUPPORTED_ARCHITECTURES = NS +
"dockerContainerSupportedArchitectures";
- public static final String DEPLOYMENT_SUPPORTED_OS_TYPES = NS +
"dockerContainerSupportedOperatingSystemTypes";
+ public static final String DEPLOYMENT_DOCKER_CONTAINER = NS +
"DockerContainer";
+ public static final String HAS_IMAGE_TAG = NS + "hasImageTag";
+ public static final String HAS_CONTAINER_NAME = NS + "hasContainerName";
+ public static final String HAS_CONTAINER_SERVICE_ID = NS +
"hasContainerServiceId";
+ public static final String HAS_CONTAINER_PORTS = NS + "hasContainerPorts";
+ public static final String HAS_CONTAINER_ENV_VARS = NS +
"hasContainerEnvironmentVariables";
+ public static final String HAS_CONTAINER_LABELS = NS + "hasContainerLabels";
+ public static final String HAS_CONTAINER_VOLUMES = NS +
"hasContainerVolumes";
+ public static final String HAS_CONTAINER_DEPENDENCIES = NS +
"hasContainerDependencies";
+ public static final String HAS_SUPPORTED_ARCHITECTURES = NS +
"hasSupportedArchitectures";
+ public static final String HAS_SUPPORTED_OS_TYPES = NS +
"hasSupportedOSTypes";
// UI Rendering
@@ -489,7 +489,7 @@ public class StreamPipes {
public static final String HAS_CONTAINER_RUNTIME_API_VERSION = NS +
"hasContainerRuntimeApiVersion";
public static final String HARDWARE_REQUIREMENT = NS + "HardwareRequirement";
public static final String HAS_GPU_REQUIREMENT = NS + "hasGpuRequirement";
- public static final String HAS_NODE_RESOURCE_PROPERTY = NS +
"hasNodeResourceProperty";
+ public static final String HAS_NODE_RESOURCE_REQUIREMENT = NS +
"hasNodeResourceRequirement";
public static final String NODE_RESOURCE_REQUIREMENT = NS +
"NodeResourceRequirement";
public static final String REQUIRES_RESOURCES = NS + "requiresResources";
public static final String REQUIRES_CPU_CORES = NS + "requiresCpuCores";
@@ -500,4 +500,12 @@ public class StreamPipes {
public static final String IS_INTERNALLY_MANAGED = NS +
"isInternallyManaged";
public static final String HAS_CORRESPONDING_ADAPTER_ID = NS +
"hasCorrespondingAdapterId";
public static final String IS_RECONFIGURABLE = NS + "isReconfigurable";
+ public static final String OPERATING_SYSTEM = NS + "OperatingSystem";
+ public static final String CONTAINER_LABEL = NS + "ContainerLabel";
+ public static final String HAS_CONTAINER_LABEL_KEY = NS +
"hasContainerLabelKey";
+ public static final String HAS_CONTAINER_LABEL_VALUE = NS +
"hasContainerLabelValue";
+ public static final String CONTAINER_ENV_VAR = NS + "ContainerEnvVar";
+ public static final String HAS_CONTAINER_ENV_KEY = NS + "hasContainerEnvKey";
+ public static final String HAS_CONTAINER_ENV_VALUE = NS +
"hasContainerEnvValue";
+ public static final String HAS_DEPLOYMENT_CONTAINER = NS +
"hasDeploymentContainer";
}
diff --git
a/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.html
b/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.html
index 37b068f..e167328 100644
---
a/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.html
+++
b/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.html
@@ -20,7 +20,10 @@
<div class="sp-dialog-content padding-20">
<div fxFlex="100" fxLayout="column">
<div fxFlex="100" fxLayout="column">
- {{node.nodeControllerId}}
+<!-- <div style="margin-bottom: 1em">-->
+<!-- <b>Node: </b> {{node.hostname}}-->
+<!-- </div>-->
+
<mat-form-field class="example-chip-list">
<mat-label>Node location tag</mat-label>
@@ -38,6 +41,107 @@
</mat-chip-list>
</mat-form-field>
+
+ <div style="margin-top: 1em" *ngFor="let connectivity of
tmpFieldDeviceResource">
+ <mat-accordion class="example-headers-align">
+ <mat-expansion-panel [expanded]="false">
+ <mat-expansion-panel-header>
+ <mat-panel-title>
+ <b>{{connectivity.deviceName}}</b>
+ </mat-panel-title>
+ <mat-panel-description>
+ Configure connectivity options
+ </mat-panel-description>
+ </mat-expansion-panel-header>
+
+ <div>
+ <div fxFlex="100" fxLayout="row">
+ <div fxFlex="40" fxLayout="row"
fxLayoutAlign="left center">
+ <span>Device name</span>
+ </div>
+ </div>
+ <div fxFlex="60" fxLayout="row"
fxLayoutAlign="start center">
+ <mat-form-field>
+ <mat-label></mat-label>
+ <input matInput
[(ngModel)]="connectivity.deviceName"
+
[value]="connectivity.deviceName">
+ </mat-form-field>
+ </div>
+ </div>
+
+ <div>
+ <div fxFlex="100" fxLayout="row">
+ <div fxFlex="40" fxLayout="row"
fxLayoutAlign="left center">
+ <span>Device type</span>
+ </div>
+ </div>
+ <div fxFlex="60" fxLayout="row"
fxLayoutAlign="start center">
+ <mat-form-field appearance="fill">
+ <mat-label>Device type</mat-label>
+ <mat-select
[(ngModel)]="connectivity.deviceType">
+ <mat-option *ngFor="let deviceType
of deviceTypes"
+
[value]="deviceType.value">
+ {{deviceType.viewValue}}
+ </mat-option>
+ </mat-select>
+ </mat-form-field>
+ </div>
+ </div>
+
+ <div>
+ <div fxFlex="100" fxLayout="row">
+ <div fxFlex="40" fxLayout="row"
fxLayoutAlign="left center">
+ <span>Access type</span>
+ </div>
+ </div>
+ <div fxFlex="60" fxLayout="row"
fxLayoutAlign="start center">
+ <mat-form-field appearance="fill">
+ <mat-label>Acess type</mat-label>
+ <mat-select
[(ngModel)]="connectivity.connectionType">
+ <mat-option *ngFor="let accessType
of accessTypes"
+
[value]="accessType.value">
+ {{accessType.viewValue}}
+ </mat-option>
+ </mat-select>
+ </mat-form-field>
+ </div>
+ </div>
+
+ <div>
+ <div fxFlex="100" fxLayout="row">
+ <div fxFlex="40" fxLayout="row"
fxLayoutAlign="left center">
+ <span>Connection</span>
+ </div>
+ </div>
+ <div fxFlex="60" fxLayout="row"
fxLayoutAlign="start center">
+ <mat-form-field>
+ <mat-label></mat-label>
+ <input matInput
[(ngModel)]="connectivity.connectionString"
+
[value]="connectivity.connectionString">
+ </mat-form-field>
+ </div>
+ </div>
+ <button mat-menu-item
+ matTooltip="Delete connectivity option"
+ matTooltipPosition="above"
+ (click)="delete(connectivity.deviceName)">
+ <mat-icon>
+ delete
+ </mat-icon>
+ <span>delete connectivity option</span>
+ </button>
+ </mat-expansion-panel>
+ </mat-accordion>
+ </div>
+ <div class="sp-dialog-actions" fxLayoutAlign="left center">
+ <button mat-button mat-raised-button
+ color="primary" (click)="addConnectivity()"
style="margin-right:10px;">
+ Add
+ </button>
+ </div>
+
+
+
</div>
</div>
</div>
diff --git
a/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.ts
b/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.ts
index 86560e4..1a67e17 100644
---
a/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.ts
+++
b/ui/src/app/configuration/node-configuration/node-configuration-details/node-configuration-details.component.ts
@@ -17,7 +17,12 @@
*/
import {Component, Input, OnInit} from '@angular/core';
-import {Message, NodeInfoDescription, PipelineOperationStatus} from
"../../../core-model/gen/streampipes-model";
+import {
+ FieldDeviceAccessResource,
+ Message,
+ NodeInfoDescription,
+ PipelineOperationStatus
+} from "../../../core-model/gen/streampipes-model";
import {FormGroup} from "@angular/forms";
import {DialogRef} from "../../../core-ui/dialog/base-dialog/dialog-ref";
import {MatChipInputEvent} from "@angular/material/chips";
@@ -46,6 +51,21 @@ export class NodeConfigurationDetailsComponent implements
OnInit {
addOnBlur = true;
readonly separatorKeysCodes: number[] = [ENTER, COMMA];
tmpTags: string[];
+ tmpFieldDeviceResource : FieldDeviceAccessResource[];
+
+ accessTypes = [
+ {value: 'local', viewValue: 'Local'},
+ {value: 'remote', viewValue: 'Remote'},
+ ];
+
+ deviceTypes = [
+ {value: 'sensor', viewValue: 'Sensor'},
+ {value: 'actuator', viewValue: 'Actuator'},
+ {value: 'camera', viewValue: 'Camera'},
+ {value: 'robot', viewValue: 'Robot'},
+ {value: 'machine', viewValue: 'Machine'},
+ {value: 'iotdevice', viewValue: 'IoT device'},
+ ];
@Input()
node: NodeInfoDescription;
@@ -58,11 +78,13 @@ export class NodeConfigurationDetailsComponent implements
OnInit {
ngOnInit(): void {
this.advancedSettings = false;
this.tmpTags = this.node.staticNodeMetadata.locationTags.map(x => x);
+ this.tmpFieldDeviceResource =
this.node.nodeResources.fieldDeviceAccessResourceList.map(x =>
Object.assign({}, x));
}
updateNodeInfo() {
let updateRequest;
this.node.staticNodeMetadata.locationTags = this.tmpTags;
+ this.node.nodeResources.fieldDeviceAccessResourceList =
this.tmpFieldDeviceResource;
updateRequest = this.nodeService.updateNodeState(this.node);
updateRequest
@@ -109,4 +131,17 @@ export class NodeConfigurationDetailsComponent implements
OnInit {
}
}
+ addConnectivity() {
+ let device = new FieldDeviceAccessResource();
+ device["@class"] =
"org.apache.streampipes.model.node.resources.fielddevice.FieldDeviceAccessResource";
+ console.log(device);
+ this.tmpFieldDeviceResource.push(device);
+ }
+
+ delete(deviceName: string) {
+ this.tmpFieldDeviceResource.forEach( (item, index) => {
+ if(item.deviceName === deviceName)
this.tmpFieldDeviceResource.splice(index, 1);
+ })
+
+ }
}
diff --git
a/ui/src/app/configuration/node-configuration/node-configuration.component.html
b/ui/src/app/configuration/node-configuration/node-configuration.component.html
index ba8bdd4..6cd38e4 100644
---
a/ui/src/app/configuration/node-configuration/node-configuration.component.html
+++
b/ui/src/app/configuration/node-configuration/node-configuration.component.html
@@ -91,19 +91,33 @@
[ngClass]="'node-inactive'">
sell
</mat-icon>
- <span>tags</span>
+ <span>edit</span>
</button>
+<!-- <button mat-menu-item-->
+<!-- [disabled]="!node.active ||
node.condition === 'OFFLINE'"-->
+<!-- matTooltip="Edit node
connectivity"-->
+<!-- matTooltipPosition="above"-->
+<!-- (click)="settings(node)">-->
+<!-- <mat-icon-->
+<!--
[ngClass]="'node-inactive'">-->
+<!-- contactless-->
+<!-- </mat-icon>-->
+<!-- <span>connectivity</span>-->
+<!-- </button>-->
</mat-menu>
</div>
<div mat-card-avatar class="node-header-avatar">
<button mat-icon-button
class="node-mat-icon-button" disabled>
<mat-icon [ngClass]="'node-inactive'">
storage
+<!-- {{node.staticNodeMetadata.type ==
'cloud' ? 'cloud' :-->
+<!-- node.staticNodeMetadata.type ==
'fog' ? 'storage' : 'developer_board'}}-->
</mat-icon>
</button>
</div>
<mat-card-title
- style="font-size:
12pt">{{node.hostname}}</mat-card-title>
+ style="font-size: 12pt">{{node.hostname}} |
+
{{node.staticNodeMetadata.type}}</mat-card-title>
<mat-card-subtitle style="font-size: 10pt">
{{node.nodeControllerId}} |
<b>{{node.nodeResources.softwareResource.os}}</b>
<div style="color: lightgrey;">
diff --git
a/ui/src/app/configuration/node-configuration/node-configuration.component.ts
b/ui/src/app/configuration/node-configuration/node-configuration.component.ts
index 5e2ad46..54a89f4 100644
---
a/ui/src/app/configuration/node-configuration/node-configuration.component.ts
+++
b/ui/src/app/configuration/node-configuration/node-configuration.component.ts
@@ -182,7 +182,7 @@ export class NodeConfigurationComponent implements OnInit{
settings(node: NodeInfoDescription) {
this.DialogService.open(NodeConfigurationDetailsComponent,{
panelType: PanelType.SLIDE_IN_PANEL,
- title: "Edit Node configuration",
+ title: "Edit node configuration: " + node.hostname,
data: {
"node": node
}
diff --git a/ui/src/app/core-model/gen/streampipes-model.ts
b/ui/src/app/core-model/gen/streampipes-model.ts
index 9ae5263..2ff1e90 100644
--- a/ui/src/app/core-model/gen/streampipes-model.ts
+++ b/ui/src/app/core-model/gen/streampipes-model.ts
@@ -19,10 +19,10 @@
/* tslint:disable */
/* eslint-disable */
// @ts-nocheck
-// Generated using typescript-generator version 2.27.744 on 2021-05-10
21:29:45.
+// Generated using typescript-generator version 2.27.744 on 2021-05-28
12:46:35.
export class AbstractStreamPipesEntity {
- "@class": "org.apache.streampipes.model.base.AbstractStreamPipesEntity" |
"org.apache.streampipes.model.base.NamedStreamPipesEntity" |
"org.apache.streampipes.model.connect.adapter.AdapterDescription" |
"org.apache.streampipes.model.connect.adapter.AdapterSetDescription" |
"org.apache.streampipes.model.connect.adapter.GenericAdapterSetDescription" |
"org.apache.streampipes.model.connect.adapter.SpecificAdapterSetDescription" |
"org.apache.streampipes.model.connect.adapter.AdapterStre [...]
+ "@class": "org.apache.streampipes.model.base.AbstractStreamPipesEntity" |
"org.apache.streampipes.model.base.NamedStreamPipesEntity" |
"org.apache.streampipes.model.connect.adapter.AdapterDescription" |
"org.apache.streampipes.model.connect.adapter.AdapterSetDescription" |
"org.apache.streampipes.model.connect.adapter.GenericAdapterSetDescription" |
"org.apache.streampipes.model.connect.adapter.SpecificAdapterSetDescription" |
"org.apache.streampipes.model.connect.adapter.AdapterStre [...]
elementId: string;
static fromData(data: AbstractStreamPipesEntity, target?:
AbstractStreamPipesEntity): AbstractStreamPipesEntity {
@@ -37,7 +37,7 @@ export class AbstractStreamPipesEntity {
}
export class UnnamedStreamPipesEntity extends AbstractStreamPipesEntity {
- "@class": "org.apache.streampipes.model.base.UnnamedStreamPipesEntity" |
"org.apache.streampipes.model.connect.guess.GuessSchema" |
"org.apache.streampipes.model.connect.rules.TransformationRuleDescription" |
"org.apache.streampipes.model.connect.rules.value.ValueTransformationRuleDescription"
|
"org.apache.streampipes.model.connect.rules.value.AddTimestampRuleDescription"
|
"org.apache.streampipes.model.connect.rules.value.AddValueTransformationRuleDescription"
| "org.apache.streamp [...]
+ "@class": "org.apache.streampipes.model.base.UnnamedStreamPipesEntity" |
"org.apache.streampipes.model.connect.guess.GuessSchema" |
"org.apache.streampipes.model.connect.rules.TransformationRuleDescription" |
"org.apache.streampipes.model.connect.rules.value.ValueTransformationRuleDescription"
|
"org.apache.streampipes.model.connect.rules.value.AddTimestampRuleDescription"
|
"org.apache.streampipes.model.connect.rules.value.AddValueTransformationRuleDescription"
| "org.apache.streamp [...]
static fromData(data: UnnamedStreamPipesEntity, target?:
UnnamedStreamPipesEntity): UnnamedStreamPipesEntity {
if (!data) {
@@ -151,8 +151,8 @@ export class NamedStreamPipesEntity extends
AbstractStreamPipesEntity {
instance.applicationLinks =
__getCopyArrayFn(ApplicationLink.fromData)(data.applicationLinks);
instance.internallyManaged = data.internallyManaged;
instance.connectedTo =
__getCopyArrayFn(__identity<string>())(data.connectedTo);
- instance.dom = data.dom;
instance.uri = data.uri;
+ instance.dom = data.dom;
return instance;
}
}
@@ -732,6 +732,40 @@ export class ConsumedMessagesInfo extends MessagesInfo {
}
}
+export class ContainerEnvVar extends UnnamedStreamPipesEntity {
+ "@class": "org.apache.streampipes.model.node.container.ContainerEnvVar";
+ key: string;
+ value: string;
+
+ static fromData(data: ContainerEnvVar, target?: ContainerEnvVar):
ContainerEnvVar {
+ if (!data) {
+ return data;
+ }
+ const instance = target || new ContainerEnvVar();
+ super.fromData(data, instance);
+ instance.key = data.key;
+ instance.value = data.value;
+ return instance;
+ }
+}
+
+export class ContainerLabel extends UnnamedStreamPipesEntity {
+ "@class": "org.apache.streampipes.model.node.container.ContainerLabel";
+ key: string;
+ value: string;
+
+ static fromData(data: ContainerLabel, target?: ContainerLabel):
ContainerLabel {
+ if (!data) {
+ return data;
+ }
+ const instance = target || new ContainerLabel();
+ super.fromData(data, instance);
+ instance.key = data.key;
+ instance.value = data.value;
+ return instance;
+ }
+}
+
export class ContainerRuntime extends UnnamedStreamPipesEntity {
"@class":
"org.apache.streampipes.model.node.resources.software.ContainerRuntime" |
"org.apache.streampipes.model.node.resources.software.DockerContainerRuntime" |
"org.apache.streampipes.model.node.resources.software.NvidiaContainerRuntime";
apiVersion: string;
@@ -1227,9 +1261,9 @@ export class DeploymentContainer extends
UnnamedStreamPipesEntity {
containerName: string;
containerPorts: string[];
dependsOnContainers: string[];
- envVars: string[];
+ envVars: ContainerEnvVar[];
imageTag: string;
- labels: { [index: string]: string };
+ labels: ContainerLabel[];
serviceId: string;
supportedArchitectures: string[];
supportedOperatingSystemTypes: string[];
@@ -1245,8 +1279,8 @@ export class DeploymentContainer extends
UnnamedStreamPipesEntity {
instance.containerName = data.containerName;
instance.serviceId = data.serviceId;
instance.containerPorts =
__getCopyArrayFn(__identity<string>())(data.containerPorts);
- instance.envVars =
__getCopyArrayFn(__identity<string>())(data.envVars);
- instance.labels = __getCopyObjectFn(__identity<string>())(data.labels);
+ instance.envVars =
__getCopyArrayFn(ContainerEnvVar.fromData)(data.envVars);
+ instance.labels =
__getCopyArrayFn(ContainerLabel.fromData)(data.labels);
instance.volumes =
__getCopyArrayFn(__identity<string>())(data.volumes);
instance.supportedArchitectures =
__getCopyArrayFn(__identity<string>())(data.supportedArchitectures);
instance.supportedOperatingSystemTypes =
__getCopyArrayFn(__identity<string>())(data.supportedOperatingSystemTypes);
@@ -1838,8 +1872,8 @@ export class GenericAdapterSetDescription extends
AdapterSetDescription implemen
}
const instance = target || new GenericAdapterSetDescription();
super.fromData(data, instance);
- instance.formatDescription =
FormatDescription.fromData(data.formatDescription);
instance.protocolDescription =
ProtocolDescription.fromData(data.protocolDescription);
+ instance.formatDescription =
FormatDescription.fromData(data.formatDescription);
instance.eventSchema = EventSchema.fromData(data.eventSchema);
return instance;
}
@@ -1857,8 +1891,8 @@ export class GenericAdapterStreamDescription extends
AdapterStreamDescription im
}
const instance = target || new GenericAdapterStreamDescription();
super.fromData(data, instance);
- instance.formatDescription =
FormatDescription.fromData(data.formatDescription);
instance.protocolDescription =
ProtocolDescription.fromData(data.protocolDescription);
+ instance.formatDescription =
FormatDescription.fromData(data.formatDescription);
instance.eventSchema = EventSchema.fromData(data.eventSchema);
return instance;
}
@@ -2305,6 +2339,7 @@ export class NodeInfoDescription extends
UnnamedStreamPipesEntity {
_rev: string;
active: boolean;
condition: NodeCondition;
+ deploymentContainers: DeploymentContainerUnion[];
hostname: string;
lastHeartBeatTime: number;
nodeBroker: NodeBrokerDescription;
@@ -2331,6 +2366,7 @@ export class NodeInfoDescription extends
UnnamedStreamPipesEntity {
instance.staticNodeMetadata =
StaticNodeMetadata.fromData(data.staticNodeMetadata);
instance.nodeResources = NodeResource.fromData(data.nodeResources);
instance.registeredContainers =
__getCopyArrayFn(DeploymentContainer.fromDataUnion)(data.registeredContainers);
+ instance.deploymentContainers =
__getCopyArrayFn(DeploymentContainer.fromDataUnion)(data.deploymentContainers);
instance.supportedElements =
__getCopyArrayFn(__identity<string>())(data.supportedElements);
instance._id = data._id;
instance._rev = data._rev;
diff --git
a/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.html
b/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.html
index d1a7c83..aae39e8 100644
---
a/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.html
+++
b/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.html
@@ -31,21 +31,56 @@
</div>
</form>
<mat-slide-toggle color="primary"
[(ngModel)]="advancedSettings">
- Configure deployment options
+ Configure deployment settings
</mat-slide-toggle>
<mat-divider *ngIf="advancedSettings" style="margin: 1em 0 1em
0;"></mat-divider>
<div *ngIf="advancedSettings">
<div>
- <b>Pipeline Operation Policies</b>
+ <b>Operation Policies</b>
</div>
- <div style="margin-top: 2em">
+ <div style="margin-top: 1em">
+ <div fxFlex="100" fxLayout="row">
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left top">
+ <span>Preemption</span>
+ </div>
+ </div>
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
+ <mat-slide-toggle #preemptionSlideToggle
color="accent"
+ [(ngModel)]="selectedPreemption"
+ [disabled]="true">
+ {{preemptionSlideToggle.checked ? 'enabled' :
'disabled'}}
+ </mat-slide-toggle>
+ </div>
+ </div>
+
+ <div style="margin-top: 1em" *ngIf="selectedPreemption">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left top">
- <span>Relay Policy</span>
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left top">
+ <span>Priority</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row" fxLayoutAlign="start
center">
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
+ <form [formGroup]="priorityForm">
+ <mat-form-field appearance="fill"
color="accent">
+ <mat-label>Select priority
class</mat-label>
+ <mat-select formControlName="priorityForm"
[disabled]="true">
+ <mat-option *ngFor="let priority of
pipelinePriorityClasses" [value]="priority">
+ {{priority.viewValue}}
+ </mat-option>
+ </mat-select>
+ </mat-form-field>
+ </form>
+ </div>
+ </div>
+
+ <div style="margin-top: 1em">
+ <div fxFlex="100" fxLayout="row">
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left top">
+ <span>Event Relay</span>
+ </div>
+ </div>
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
<mat-button-toggle-group
#relayGroup="matButtonToggleGroup" aria-label="Relay strategy"
[disabled]="true"
[value]="selectedRelayStrategyVal"
@@ -58,11 +93,11 @@
<div style="margin-top: 1em">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left top">
- <span>Execution Policy</span>
+ <div fxFlex="66" fxLayout="row"
fxLayoutAlign="left top">
+ <b>Deployment Options</b>
</div>
</div>
- <div fxFlex="50" fxLayout="row" fxLayoutAlign="start
center">
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
<mat-radio-group
aria-labelledby="execution-policy-radio-group-label"
class="execution-policy-radio-group"
@@ -70,7 +105,7 @@
(change)="onExecutionPolicyChange($event.value)">
<mat-radio-button
class="execution-policy-radio-button"
*ngFor="let policy of
pipelineExecutionPolicies"
-
[disabled]="isExecutinoPolicyDisabled()"
+ [disabled]="true"
[value]="policy">{{policy}}
</mat-radio-button>
</mat-radio-group>
@@ -79,9 +114,14 @@
<div style="margin: 1em 0 1em 0;">
<div>
- <b>Node Execution Targets</b>
+ <b>Migration Targets</b>
+ </div>
+ <div style="margin: 1em 0 1em 0;"
*ngIf="selectedPipelineExecutionPolicy == 'custom'">
+ <node-tag-selector [nodes]="edgeNodes"
[selectedTagsAfterUpdate]="selectedNodeTags"
+
(createDynamicallySelectedTags)="nodesFromSelectedTags($event)"
+
(emitSelectedNodeTags)="updateNodeTags($event)"></node-tag-selector>
</div>
- <div style="margin-top: 2em">
+ <div style="margin-top: 1em">
<mat-accordion class="example-headers-align">
<mat-expansion-panel (opened)="panelOpenState
= true"
@@ -89,19 +129,19 @@
[expanded]="panelOpenState">
<mat-expansion-panel-header>
<mat-panel-description>
- Modify deployment targets:
<b>{{selectedPipelineExecutionPolicy}}</b>
+ Select migration targets:
<b>{{selectedPipelineExecutionPolicy}}</b>
<mat-icon>storage</mat-icon>
</mat-panel-description>
</mat-expansion-panel-header>
<div *ngFor="let processors of
tmpPipeline.sepas">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left center">
+ <div fxFlex="66" fxLayout="row"
fxLayoutAlign="left center">
<span>{{processors.name}}</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="start center">
+ <div fxFlex="45" fxLayout="row"
fxLayoutAlign="start center">
<mat-form-field dense>
<mat-select
[(ngModel)]="processors.deploymentTargetNodeId"
[disabled]="disableNodeSelectionForProcessors.value"
@@ -117,12 +157,12 @@
</div>
<div *ngFor="let sinks of
tmpPipeline.actions">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left center">
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left center">
<span>{{sinks.name}}</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="start center">
+ <div fxFlex="45" fxLayout="row"
fxLayoutAlign="start center">
<mat-form-field dense>
<mat-select
[(ngModel)]="sinks.deploymentTargetNodeId"
[disabled]="disableNodeSelectionForSinks.value"
diff --git
a/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.scss
b/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.scss
index fe91666..b915687 100644
---
a/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.scss
+++
b/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.scss
@@ -19,7 +19,7 @@
@import '../../../../scss/sp/sp-dialog.scss';
.sp-dialog-container {
- width: 520px;
+ width: 530px;
}
.customize-section {
diff --git
a/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.ts
b/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.ts
index cbd6276..eb34d39 100644
---
a/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.ts
+++
b/ui/src/app/editor/dialog/migrate-pipeline-processors/migrate-pipeline-processors.component.ts
@@ -23,7 +23,7 @@ import {
Pipeline, PipelineOperationStatus,
StaticNodeMetadata
} from "../../../core-model/gen/streampipes-model";
-import {FormControl, FormGroup, Validators} from "@angular/forms";
+import {FormBuilder, FormControl, FormGroup, Validators} from "@angular/forms";
import {EditorService} from "../../services/editor.service";
import {DialogRef} from "../../../core-ui/dialog/base-dialog/dialog-ref";
import {ObjectProvider} from "../../services/object-provider.service";
@@ -40,7 +40,7 @@ import {map} from "rxjs/operators";
styleUrls: ['./migrate-pipeline-processors.component.scss']
})
export class MigratePipelineProcessorsComponent implements OnInit {
-
+ priorityForm: FormGroup;
submitPipelineForm: FormGroup = new FormGroup({});
saving: boolean = false;
saved: boolean = false;
@@ -56,11 +56,19 @@ export class MigratePipelineProcessorsComponent implements
OnInit {
tmpPipeline: Pipeline;
panelOpenState: boolean;
pipelineExecutionPolicies: string[] = ['default', 'locality-aware',
'custom'];
+ selectedPreemption: boolean;
+ selectedNodeTags: string[];
+
+ pipelinePriorityClasses = [
+ {value: 1, viewValue: 'low'},
+ {value: 5, viewValue: 'medium'},
+ {value: 10, viewValue: 'high'}];
@Input()
pipeline: Pipeline;
constructor(private editorService: EditorService,
+ private formBuilder: FormBuilder,
private dialogRef: DialogRef<MigratePipelineProcessorsComponent>,
private objectProvider: ObjectProvider,
private pipelineService: PipelineService,
@@ -72,9 +80,13 @@ export class MigratePipelineProcessorsComponent implements
OnInit {
}
ngOnInit() {
+ this.loadAndPrepareEdgeNodes();
+
this.tmpPipeline = this.deepCopy(this.pipeline);
- this.loadAndPrepareEdgeNodes();
+ this.priorityForm = this.formBuilder.group({
+ priorityForm: [null, Validators.required]
+ });
this.submitPipelineForm.addControl("pipelineName", new
FormControl(this.tmpPipeline.name,
[Validators.required,
@@ -90,8 +102,12 @@ export class MigratePipelineProcessorsComponent implements
OnInit {
this.tmpPipeline.description = value;
});
- this.selectedRelayStrategyVal = "buffer";
- this.selectedPipelineExecutionPolicy = "custom";
+ this.selectedPipelineExecutionPolicy = 'custom';
+ this.selectedRelayStrategyVal = this.tmpPipeline.eventRelayStrategy;
+ this.selectedPreemption = this.tmpPipeline.preemption;
+
+ const selectedPriorityClass = this.pipelinePriorityClasses.find(c =>
c.value == this.tmpPipeline.priorityScore);
+ this.priorityForm.get('priorityForm').setValue(selectedPriorityClass);
}
@@ -116,7 +132,7 @@ export class MigratePipelineProcessorsComponent implements
OnInit {
}
loadAndPrepareEdgeNodes() {
- this.nodeService.getOnlineNodes().subscribe(response => {
+ this.nodeService.getOnlineNodes().subscribe((response :
NodeInfoDescription[]) => {
this.edgeNodes = response;
this.addAppIds(this.tmpPipeline.sepas, this.edgeNodes);
this.addAppIds(this.tmpPipeline.actions, this.edgeNodes);
@@ -234,7 +250,47 @@ export class MigratePipelineProcessorsComponent implements
OnInit {
this.selectedPipelineExecutionPolicy = value;
}
- isExecutinoPolicyDisabled() {
- return true;
+ isExecutionPolicyDisabled() {
+ return false;
+ }
+
+ private applyTagBasedPolicy(filteredNodes: NodeInfoDescription[]) {
+ if (filteredNodes.length > 0) {
+ this.tmpPipeline.sepas.forEach(processor => {
+ this.deploymentOptions[processor.appId] = [];
+
+ filteredNodes.forEach(filteredNode => {
+
+ if (filteredNode.supportedElements.length != 0 &&
+ filteredNode.supportedElements.some(appId => appId ===
processor.appId)) {
+ this.deploymentOptions[processor.appId].push(filteredNode);
+ }
+ })
+ })
+
+ // this.tmpPipeline.actions.forEach(actions => {
+ // this.deploymentOptions[actions.appId] = [];
+ //
+ // filteredNodes.forEach(filteredNode => {
+ //
+ // if (filteredNode.supportedElements.length != 0 &&
+ // filteredNode.supportedElements.some(appId => appId ===
actions.appId)) {
+ // this.deploymentOptions[actions.appId].push(filteredNode);
+ // }
+ // })
+ // })
+
+ } else {
+ this.addAppIds(this.tmpPipeline.sepas, this.edgeNodes);
+ this.addAppIds(this.tmpPipeline.actions, this.edgeNodes);
+ }
+ }
+
+ nodesFromSelectedTags(filteredNodes: NodeInfoDescription[]) {
+ this.applyTagBasedPolicy(filteredNodes)
+ }
+
+ updateNodeTags($event: any) {
+ this.selectedNodeTags = $event;
}
}
diff --git
a/ui/src/app/editor/dialog/save-pipeline/node-tag-selector/node-tag-selector.component.ts
b/ui/src/app/editor/dialog/save-pipeline/node-tag-selector/node-tag-selector.component.ts
index e5af392..d8a666d 100644
---
a/ui/src/app/editor/dialog/save-pipeline/node-tag-selector/node-tag-selector.component.ts
+++
b/ui/src/app/editor/dialog/save-pipeline/node-tag-selector/node-tag-selector.component.ts
@@ -46,23 +46,26 @@ export class NodeTagSelectorComponent implements OnInit {
dynamicallySelectedTags: NodeTags["name"][] = [];
ngOnInit(): void {
- this.nodes.forEach(node => {
- if (node.supportedElements.length > 0 ||
node.registeredContainers.length > 0) {
- for (let tag of node.staticNodeMetadata.locationTags) {
- if (!this.nodeTags.some(n => n.name === tag)) {
- this.nodeTags.push({'name': tag, selected: false})
+ console.log("here");
+ if (this.nodes != undefined) {
+ this.nodes.forEach(node => {
+ if (node.supportedElements.length > 0 ||
node.registeredContainers.length > 0) {
+ for (let tag of node.staticNodeMetadata.locationTags) {
+ if (!this.nodeTags.some(n => n.name === tag)) {
+ this.nodeTags.push({'name': tag, selected: false})
+ }
}
}
- }
- })
- if (this.selectedTagsAfterUpdate && this.selectedTagsAfterUpdate.length >
0) {
- this.selectedTagsAfterUpdate.forEach(oldTag => {
- this.nodeTags.forEach(entry => {
- if (entry.name === oldTag) {
- entry.selected = true;
- }
- })
})
+ if (this.selectedTagsAfterUpdate && this.selectedTagsAfterUpdate.length
> 0) {
+ this.selectedTagsAfterUpdate.forEach(oldTag => {
+ this.nodeTags.forEach(entry => {
+ if (entry.name === oldTag) {
+ entry.selected = true;
+ }
+ })
+ })
+ }
}
}
diff --git
a/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.html
b/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.html
index f3e27dc..83cb1b6 100644
--- a/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.html
+++ b/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.html
@@ -47,21 +47,21 @@
Start pipeline immediately
</mat-slide-toggle>
<mat-slide-toggle color="primary"
[(ngModel)]="advancedSettings">
- Configure deployment options
+ Advanced deployment settings
</mat-slide-toggle>
<mat-divider *ngIf="advancedSettings" style="margin: 1em 0 1em
0;"></mat-divider>
<div *ngIf="advancedSettings">
<div>
- <b>Pipeline Operation Policies</b>
+ <b>Operation Policies</b>
</div>
<div style="margin-top: 1em">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left top">
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left top">
<span>Preemption</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row" fxLayoutAlign="start
center">
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
<mat-slide-toggle #preemptionSlideToggle
color="accent" [(ngModel)]="selectedPreemption"
(toggleChange)="loadDefaultPreemption()">
{{preemptionSlideToggle.checked ? 'enabled' :
'disabled'}}
@@ -71,11 +71,11 @@
<div style="margin-top: 1em" *ngIf="selectedPreemption">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left top">
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left top">
<span>Priority</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row" fxLayoutAlign="start
center">
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
<form [formGroup]="priorityForm">
<mat-form-field appearance="fill" color="accent">
<mat-label>Select priority class</mat-label>
@@ -93,11 +93,11 @@
<div style="margin-top: 1em">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left top">
- <span>Relay Policy</span>
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left top">
+ <span>Event Relay</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row" fxLayoutAlign="start
center">
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
<mat-button-toggle-group
#relayGroup="matButtonToggleGroup" aria-label="Relay strategy"
[value]="selectedRelayStrategyVal"
(change)="onSelectedRelayStrategyChange(relayGroup.value)">
@@ -109,11 +109,11 @@
<div style="margin-top: 1em">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left top">
- <span>Execution Policy</span>
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left top">
+ <b>Deployment Options</b>
</div>
</div>
- <div fxFlex="50" fxLayout="row" fxLayoutAlign="start
center">
+ <div fxFlex="45" fxLayout="row" fxLayoutAlign="start
center">
<mat-radio-group
aria-labelledby="execution-policy-radio-group-label"
class="execution-policy-radio-group"
@@ -151,11 +151,11 @@
<div *ngFor="let processors of
tmpPipeline.sepas">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left center">
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left center">
<span>{{processors.name}}</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="start center">
+ <div fxFlex="45" fxLayout="row"
fxLayoutAlign="start center">
<mat-form-field dense>
<mat-select
[(ngModel)]="processors.deploymentTargetNodeId"
[disabled]="disableNodeSelection.value"
@@ -171,11 +171,11 @@
</div>
<div *ngFor="let sinks of
tmpPipeline.actions">
<div fxFlex="100" fxLayout="row">
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="left center">
+ <div fxFlex="55" fxLayout="row"
fxLayoutAlign="left center">
<span>{{sinks.name}}</span>
</div>
</div>
- <div fxFlex="50" fxLayout="row"
fxLayoutAlign="start center">
+ <div fxFlex="45" fxLayout="row"
fxLayoutAlign="start center">
<mat-form-field dense>
<mat-select
[(ngModel)]="sinks.deploymentTargetNodeId"
[disabled]="disableNodeSelection.value"
diff --git
a/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.scss
b/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.scss
index 31cc55a..b915687 100644
--- a/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.scss
+++ b/ui/src/app/editor/dialog/save-pipeline/save-pipeline.component.scss
@@ -19,7 +19,7 @@
@import '../../../../scss/sp/sp-dialog.scss';
.sp-dialog-container {
- width: 500px;
+ width: 530px;
}
.customize-section {
diff --git a/ui/src/app/pipelines/services/pipeline-operations.service.ts
b/ui/src/app/pipelines/services/pipeline-operations.service.ts
index df76ba5..2540312 100644
--- a/ui/src/app/pipelines/services/pipeline-operations.service.ts
+++ b/ui/src/app/pipelines/services/pipeline-operations.service.ts
@@ -126,7 +126,7 @@ export class PipelineOperationsService {
this.PipelineService.getPipelineById(pipelineId).subscribe(pipeline => {
this.DialogService.open(MigratePipelineProcessorsComponent,{
panelType: PanelType.SLIDE_IN_PANEL,
- title: "Live-Migrate pipeline processors",
+ title: "Live-Migrate Pipeline Processors",
data: {
"pipeline": pipeline
}