This is an automated email from the ASF dual-hosted git repository.
healchow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 41e5f768a [INLONG-4369][Manager] Support for modification of
information after approval (#4370)
41e5f768a is described below
commit 41e5f768a93b06eddc378aae5c28cc1920b40877
Author: healzhou <[email protected]>
AuthorDate: Wed May 25 17:43:32 2022 +0800
[INLONG-4369][Manager] Support for modification of information after
approval (#4370)
---
.../manager/common/auth/DefaultAuthentication.java | 8 ++++----
.../common/pojo/group/InlongGroupApproveRequest.java | 3 +++
.../common/pojo/source/pulsar/PulsarSourceDTO.java | 1 +
.../manager/service/core/plugin/PluginClassLoader.java | 12 +++++++++---
.../manager/service/group/InlongGroupServiceImpl.java | 14 +++++++++++---
inlong-manager/manager-web/pom.xml | 4 ++++
.../manager/workflow/event/LogableEventListener.java | 4 ++--
.../manager/workflow/util/WorkflowBeanUtils.java | 18 +++++++++++++++---
8 files changed, 49 insertions(+), 15 deletions(-)
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/auth/DefaultAuthentication.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/auth/DefaultAuthentication.java
index b2a146dbe..9661dcb03 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/auth/DefaultAuthentication.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/auth/DefaultAuthentication.java
@@ -30,7 +30,7 @@ import java.util.Map;
@NoArgsConstructor
public class DefaultAuthentication implements Authentication {
- public static final String USER_NAME = "user_name";
+ public static final String USERNAME = "username";
public static final String PASSWORD = "password";
@@ -52,15 +52,15 @@ public class DefaultAuthentication implements
Authentication {
@Override
public void configure(Map<String, String> properties) {
- AssertUtils.notEmpty(properties, "Properties should not be empty when
init DefaultAuthentification");
- this.userName = properties.get(USER_NAME);
+ AssertUtils.notEmpty(properties, "Properties should not be empty when
init DefaultAuthentication");
+ this.userName = properties.get(USERNAME);
this.password = properties.get(PASSWORD);
}
@Override
public String toString() {
ObjectNode objectNode = OBJECT_MAPPER.createObjectNode();
- objectNode.put(USER_NAME, this.getUserName());
+ objectNode.put(USERNAME, this.getUserName());
objectNode.put(PASSWORD, this.getPassword());
return objectNode.toString();
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/group/InlongGroupApproveRequest.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/group/InlongGroupApproveRequest.java
index 3f6845127..3d66691e0 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/group/InlongGroupApproveRequest.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/group/InlongGroupApproveRequest.java
@@ -40,6 +40,9 @@ public class InlongGroupApproveRequest {
@ApiModelProperty(value = "MQ resource, for Tube, it is Topic, for Pulsar,
it is Namespace")
private String mqResource;
+ @ApiModelProperty(value = "Inlong cluster tag, inlong group will be
associated with the cluster")
+ private String inlongClusterTag;
+
@ApiModelProperty(value = "The partition num of Pulsar topic, between
1-20")
private Integer topicPartitionNum;
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/source/pulsar/PulsarSourceDTO.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/source/pulsar/PulsarSourceDTO.java
index 28f6bf6f8..9f943441b 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/source/pulsar/PulsarSourceDTO.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/source/pulsar/PulsarSourceDTO.java
@@ -60,6 +60,7 @@ public class PulsarSourceDTO {
@ApiModelProperty("Configure the Source's startup mode. "
+ "Available options are earliest, latest, external-subscription,
and specific-offsets.")
+ @Builder.Default
private String scanStartupMode = "earliest";
/**
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/plugin/PluginClassLoader.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/plugin/PluginClassLoader.java
index e64940a54..29a56b538 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/plugin/PluginClassLoader.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/plugin/PluginClassLoader.java
@@ -122,8 +122,14 @@ public class PluginClassLoader extends URLClassLoader {
* load pluginDefinition in **.jar/META-INF/plugin.yaml
*/
private void loadPluginDefinition() throws IOException {
- List<PluginDefinition> definitions = new ArrayList();
- for (File jarFile : pluginDirectory.listFiles()) {
+ File[] files = pluginDirectory.listFiles();
+ if (files == null) {
+ log.warn("plugin directory {} has no files", pluginDirectory);
+ return;
+ }
+
+ List<PluginDefinition> definitions = new ArrayList<>();
+ for (File jarFile : files) {
if (!jarFile.getName().endsWith(".jar")) {
log.warn("{} is not valid plugin jar, skip to load", jarFile);
continue;
@@ -138,7 +144,7 @@ public class PluginClassLoader extends URLClassLoader {
definitions.add(definition);
}
pluginDefinitionMap = definitions.stream()
- .collect(Collectors.toMap(definition -> definition.getName(),
definition -> definition));
+ .collect(Collectors.toMap(PluginDefinition::getName,
definition -> definition));
}
private void checkPluginValid(File jarFile, PluginDefinition
pluginDefinition) {
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
index 29e01a414..d40eb0a7d 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
@@ -23,6 +23,7 @@ import com.github.pagehelper.PageInfo;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import org.apache.commons.collections.CollectionUtils;
+import org.apache.commons.lang3.StringUtils;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.enums.SourceType;
@@ -342,9 +343,8 @@ public class InlongGroupServiceImpl implements
InlongGroupService {
return topicInfo;
}
- // TODO
@Override
- @Transactional(rollbackFor = Throwable.class)
+ @Transactional(rollbackFor = Throwable.class, propagation =
Propagation.REQUIRES_NEW)
public boolean updateAfterApprove(InlongGroupApproveRequest approveInfo,
String operator) {
LOGGER.debug("begin to update inlong group after approve={}",
approveInfo);
@@ -356,9 +356,17 @@ public class InlongGroupServiceImpl implements
InlongGroupService {
Preconditions.checkNotNull(mqType, "MQ type cannot by empty");
// Update status to [GROUP_APPROVE_PASSED]
- // If you need to change inlong group info after approve, just do in
here
this.updateStatus(groupId, GroupStatus.APPROVE_PASSED.getCode(),
operator);
+ // update other info for inlong group after approve
+ if (StringUtils.isNotBlank(approveInfo.getInlongClusterTag())) {
+ InlongGroupEntity entity = new InlongGroupEntity();
+ entity.setInlongGroupId(approveInfo.getInlongGroupId());
+ entity.setInlongClusterTag(approveInfo.getInlongClusterTag());
+ entity.setModifier(operator);
+ groupMapper.updateByIdentifierSelective(entity);
+ }
+
LOGGER.info("success to update inlong group status after approve for
groupId={}", groupId);
return true;
}
diff --git a/inlong-manager/manager-web/pom.xml
b/inlong-manager/manager-web/pom.xml
index 4fdf77fec..c6968d787 100644
--- a/inlong-manager/manager-web/pom.xml
+++ b/inlong-manager/manager-web/pom.xml
@@ -105,6 +105,10 @@
<groupId>javax.validation</groupId>
<artifactId>validation-api</artifactId>
</exclusion>
+ <exclusion>
+ <artifactId>javassist</artifactId>
+ <groupId>org.javassist</groupId>
+ </exclusion>
</exclusions>
</dependency>
<dependency>
diff --git
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/event/LogableEventListener.java
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/event/LogableEventListener.java
index 105543049..6e1affa2b 100644
---
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/event/LogableEventListener.java
+++
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/event/LogableEventListener.java
@@ -65,7 +65,7 @@ public abstract class LogableEventListener<EventType extends
WorkflowEvent> impl
log.debug("listener execute result: {} - {}",
workflowEventLogEntity, result);
return result;
} catch (Exception e) {
- log.error("execute listener {} error: {}", workflowEventLogEntity,
e);
+ log.error("execute listener " + workflowEventLogEntity + " error:
", e);
if (!async()) {
throw new WorkflowListenerException(e);
}
@@ -86,7 +86,7 @@ public abstract class LogableEventListener<EventType extends
WorkflowEvent> impl
} catch (Exception e) {
logEntity.setStatus(EventStatus.FAILED.getStatus());
logEntity.setException(e.getMessage());
- log.error("execute listener {} error: {}", logEntity, e);
+ log.error("execute listener " + logEntity + " error: ", e);
if (!async()) {
throw new WorkflowListenerException(e.getMessage());
}
diff --git
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/util/WorkflowBeanUtils.java
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/util/WorkflowBeanUtils.java
index c89475860..e4a885f4f 100644
---
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/util/WorkflowBeanUtils.java
+++
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/util/WorkflowBeanUtils.java
@@ -17,7 +17,6 @@
package org.apache.inlong.manager.workflow.util;
-import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.commons.lang3.StringUtils;
@@ -72,7 +71,8 @@ public class WorkflowBeanUtils {
if (taskEntity == null) {
return null;
}
- return TaskResponse.builder()
+
+ TaskResponse taskResponse = TaskResponse.builder()
.id(taskEntity.getId())
.type(taskEntity.getType())
.processId(taskEntity.getProcessId())
@@ -88,6 +88,18 @@ public class WorkflowBeanUtils {
.startTime(taskEntity.getStartTime())
.endTime(taskEntity.getEndTime())
.build();
+
+ try {
+ JsonNode formData = null;
+ if (StringUtils.isNotBlank(taskEntity.getFormData())) {
+ formData = OBJECT_MAPPER.readTree(taskEntity.getFormData());
+ }
+ taskResponse.setFormData(formData);
+ } catch (Exception e) {
+ LOGGER.error("parse form data error: ", e);
+ }
+
+ return taskResponse;
}
/**
@@ -122,7 +134,7 @@ public class WorkflowBeanUtils {
extParams = OBJECT_MAPPER.readTree(entity.getExtParams());
}
processResponse.setExtParams(extParams);
- } catch (JsonProcessingException e) {
+ } catch (Exception e) {
LOGGER.error("parse form data error: ", e);
throw new JsonException("parse form data or ext params error");
}