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/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new cae195512 [INLONG-7479][Manager] Forbidden to configure stream when
group configuration fails (#7482)
cae195512 is described below
commit cae195512b9ef023364db36f4f5d3363f67b01c1
Author: fuweng11 <[email protected]>
AuthorDate: Sat Mar 4 19:25:14 2023 +0800
[INLONG-7479][Manager] Forbidden to configure stream when group
configuration fails (#7482)
---
.../org/apache/inlong/manager/common/enums/GroupStatus.java | 3 +--
.../manager/service/listener/queue/QueueResourceListener.java | 2 ++
.../service/listener/queue/StreamQueueResourceListener.java | 5 +++++
.../manager/service/listener/sink/SinkResourceListener.java | 6 +++++-
.../service/listener/sink/StreamSinkResourceListener.java | 11 ++++++++++-
.../manager/service/listener/sort/SortConfigListener.java | 2 ++
.../service/listener/sort/StreamSortConfigListener.java | 5 +++++
inlong-manager/manager-web/bin/restart.sh | 6 +++---
8 files changed, 33 insertions(+), 7 deletions(-)
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupStatus.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupStatus.java
index ac31aed41..565b43576 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupStatus.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/GroupStatus.java
@@ -19,7 +19,6 @@ package org.apache.inlong.manager.common.enums;
import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
-
import org.apache.inlong.manager.common.exceptions.BusinessException;
import java.util.Locale;
@@ -67,7 +66,7 @@ public enum GroupStatus {
GROUP_STATE_AUTOMATON.put(CONFIG_ING, Sets.newHashSet(CONFIG_ING,
CONFIG_FAILED, CONFIG_SUCCESSFUL));
GROUP_STATE_AUTOMATON.put(CONFIG_FAILED,
- Sets.newHashSet(CONFIG_FAILED, CONFIG_SUCCESSFUL,
TO_BE_APPROVAL, DELETING));
+ Sets.newHashSet(CONFIG_FAILED, CONFIG_ING, CONFIG_SUCCESSFUL,
TO_BE_APPROVAL, DELETING));
GROUP_STATE_AUTOMATON.put(CONFIG_SUCCESSFUL,
Sets.newHashSet(CONFIG_SUCCESSFUL, TO_BE_APPROVAL, CONFIG_ING,
SUSPENDING, DELETING, FINISH));
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/QueueResourceListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/QueueResourceListener.java
index 8e50884f0..1b544b2a6 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/QueueResourceListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/QueueResourceListener.java
@@ -21,6 +21,7 @@ import com.google.common.util.concurrent.ThreadFactoryBuilder;
import lombok.extern.slf4j.Slf4j;
import org.apache.inlong.manager.common.consts.InlongConstants;
import org.apache.inlong.manager.common.enums.GroupOperateType;
+import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.enums.TaskEvent;
import org.apache.inlong.manager.common.enums.TaskStatus;
import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
@@ -103,6 +104,7 @@ public class QueueResourceListener implements
QueueOperateListener {
GroupResourceProcessForm groupProcessForm = (GroupResourceProcessForm)
context.getProcessForm();
final String groupId = groupProcessForm.getInlongGroupId();
// ensure the inlong group exists
+ groupService.updateStatus(groupId, GroupStatus.CONFIG_ING.getCode(),
context.getOperator());
InlongGroupInfo groupInfo = groupService.get(groupId);
if (groupInfo == null) {
String msg = "inlong group not found with groupId=" + groupId;
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/StreamQueueResourceListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/StreamQueueResourceListener.java
index a9e778a8f..062eaaf59 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/StreamQueueResourceListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/queue/StreamQueueResourceListener.java
@@ -20,8 +20,10 @@ package org.apache.inlong.manager.service.listener.queue;
import lombok.extern.slf4j.Slf4j;
import org.apache.inlong.manager.common.consts.InlongConstants;
import org.apache.inlong.manager.common.enums.GroupOperateType;
+import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.enums.TaskEvent;
import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
+import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
import org.apache.inlong.manager.pojo.workflow.form.process.ProcessForm;
@@ -83,6 +85,9 @@ public class StreamQueueResourceListener implements
QueueOperateListener {
log.error(msg);
throw new WorkflowListenerException(msg);
}
+ GroupStatus groupStatus = GroupStatus.forCode(groupInfo.getStatus());
+ Preconditions.expectTrue(GroupStatus.CONFIG_FAILED != groupStatus,
+ String.format("group status=%s not support start stream for
groupId=%s", groupStatus, groupId));
final String streamId = streamInfo.getInlongStreamId();
// Read the current information
streamProcessForm.setGroupInfo(groupInfo);
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/SinkResourceListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/SinkResourceListener.java
index 49c07d6a2..f43d89d6a 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/SinkResourceListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/SinkResourceListener.java
@@ -20,11 +20,13 @@ package org.apache.inlong.manager.service.listener.sink;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.inlong.manager.common.consts.InlongConstants;
+import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.enums.TaskEvent;
import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
import org.apache.inlong.manager.pojo.sink.SinkInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
import
org.apache.inlong.manager.pojo.workflow.form.process.GroupResourceProcessForm;
+import org.apache.inlong.manager.service.group.InlongGroupService;
import org.apache.inlong.manager.service.resource.sink.SinkResourceOperator;
import
org.apache.inlong.manager.service.resource.sink.SinkResourceOperatorFactory;
import org.apache.inlong.manager.service.stream.InlongStreamService;
@@ -51,6 +53,8 @@ public class SinkResourceListener implements
SinkOperateListener {
@Autowired
private InlongStreamService streamService;
@Autowired
+ private InlongGroupService groupService;
+ @Autowired
private SinkResourceOperatorFactory sinkOperatorFactory;
@Override
@@ -63,7 +67,7 @@ public class SinkResourceListener implements
SinkOperateListener {
GroupResourceProcessForm form = (GroupResourceProcessForm)
context.getProcessForm();
String groupId = form.getInlongGroupId();
log.info("begin to create sink resources for groupId={}", groupId);
-
+ groupService.updateStatus(groupId, GroupStatus.CONFIG_ING.getCode(),
context.getOperator());
List<String> streamIdList = new ArrayList<>();
List<InlongStreamInfo> streamList = streamService.list(groupId);
if (CollectionUtils.isNotEmpty(streamList)) {
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/StreamSinkResourceListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/StreamSinkResourceListener.java
index 0731f0209..a186a2ba4 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/StreamSinkResourceListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sink/StreamSinkResourceListener.java
@@ -21,12 +21,16 @@ import com.google.common.collect.Lists;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.inlong.manager.common.consts.InlongConstants;
+import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.enums.TaskEvent;
+import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
+import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.sink.SinkInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
import org.apache.inlong.manager.pojo.workflow.form.process.ProcessForm;
import
org.apache.inlong.manager.pojo.workflow.form.process.StreamResourceProcessForm;
+import org.apache.inlong.manager.service.group.InlongGroupService;
import org.apache.inlong.manager.service.resource.sink.SinkResourceOperator;
import
org.apache.inlong.manager.service.resource.sink.SinkResourceOperatorFactory;
import org.apache.inlong.manager.workflow.WorkflowContext;
@@ -49,6 +53,8 @@ public class StreamSinkResourceListener implements
SinkOperateListener {
@Autowired
private StreamSinkEntityMapper sinkEntityMapper;
@Autowired
+ private InlongGroupService groupService;
+ @Autowired
private SinkResourceOperatorFactory resourceOperatorFactory;
@Override
@@ -69,7 +75,10 @@ public class StreamSinkResourceListener implements
SinkOperateListener {
final String groupId = streamInfo.getInlongGroupId();
final String streamId = streamInfo.getInlongStreamId();
log.info("begin to create sink resource for groupId={}, streamId={}",
groupId, streamId);
-
+ InlongGroupInfo groupInfo = groupService.get(groupId);
+ GroupStatus groupStatus = GroupStatus.forCode(groupInfo.getStatus());
+ Preconditions.expectTrue(GroupStatus.CONFIG_FAILED != groupStatus,
+ String.format("group status=%s not support start stream for
groupId=%s", groupStatus, groupId));
List<SinkInfo> sinkInfos = sinkEntityMapper.selectAllConfig(groupId,
Lists.newArrayList(streamId));
List<SinkInfo> needCreateResources = sinkInfos.stream()
.filter(sinkInfo ->
InlongConstants.ENABLE_CREATE_RESOURCE.equals(sinkInfo.getEnableCreateResource()))
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/SortConfigListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/SortConfigListener.java
index ee1b24e86..e15c8a803 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/SortConfigListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/SortConfigListener.java
@@ -19,6 +19,7 @@ package org.apache.inlong.manager.service.listener.sort;
import org.apache.commons.collections.CollectionUtils;
import org.apache.inlong.manager.common.enums.GroupOperateType;
+import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.enums.TaskEvent;
import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
@@ -87,6 +88,7 @@ public class SortConfigListener implements
SortOperateListener {
return ListenerResult.success();
}
// ensure the inlong group exists
+ groupService.updateStatus(groupId, GroupStatus.CONFIG_ING.getCode(),
context.getOperator());
InlongGroupInfo groupInfo = groupService.get(groupId);
if (groupInfo == null) {
String msg = "inlong group not found with groupId=" + groupId;
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/StreamSortConfigListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/StreamSortConfigListener.java
index c2e943d03..293c35655 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/StreamSortConfigListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/listener/sort/StreamSortConfigListener.java
@@ -19,8 +19,10 @@ package org.apache.inlong.manager.service.listener.sort;
import org.apache.commons.collections.CollectionUtils;
import org.apache.inlong.manager.common.enums.GroupOperateType;
+import org.apache.inlong.manager.common.enums.GroupStatus;
import org.apache.inlong.manager.common.enums.TaskEvent;
import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
+import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.sink.StreamSink;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
@@ -92,6 +94,9 @@ public class StreamSortConfigListener implements
SortOperateListener {
}
InlongGroupInfo groupInfo = groupService.get(groupId);
+ GroupStatus groupStatus = GroupStatus.forCode(groupInfo.getStatus());
+ Preconditions.expectTrue(GroupStatus.CONFIG_FAILED != groupStatus,
+ String.format("group status=%s not support start stream for
groupId=%s", groupStatus, groupId));
List<StreamSink> streamSinks = streamInfo.getSinkList();
if (CollectionUtils.isEmpty(streamSinks)) {
LOGGER.warn("not build sort config for groupId={}, streamId={}, as
not found any sinks", groupId, streamId);
diff --git a/inlong-manager/manager-web/bin/restart.sh
b/inlong-manager/manager-web/bin/restart.sh
index aba17039a..696b2451f 100755
--- a/inlong-manager/manager-web/bin/restart.sh
+++ b/inlong-manager/manager-web/bin/restart.sh
@@ -36,10 +36,10 @@ BIN_PATH=$(
echo "restart in" ${BIN_PATH}
# Stop service
-sh "$BIN_PATH"/shutdown.sh
+bash +x "$BIN_PATH"/shutdown.sh
sleep 1s
-echo ""
+echo "begin to execute the startup command..."
# Start service
-sh "$BIN_PATH"/startup.sh
+bash +x "$BIN_PATH"/startup.sh